data_entry_by_meta.py 8.6 KB
# coding=utf-8
#author:        4N
#createtime:    2021/1/27
#email:         nheweijun@sina.com

from osgeo.ogr import *
import uuid

import time
from ..models import *

import json
import re
from app.util.component.ApiTemplate import ApiTemplate
from app.util.component.PGUtil import PGUtil
from app.util.component.StructurePrint import StructurePrint
from sqlalchemy.orm import Session
import configure
import datetime
import multiprocessing
from .util.EntryDataVacuate import EntryDataVacuate
from app.util.component.TaskController import TaskController
from app.util.component.TaskWriter import TaskWriter
class Api(ApiTemplate):

    api_name = "通过meta入库"

    def process(self):

        #设置任务信息
        self.para["task_guid"] = uuid.uuid1().__str__()
        self.para["task_time"] = time.time()

        #返回结果
        res={}
    
        try:
            #检测目录
            if Catalog.query.filter_by(pguid=self.para.get("guid")).all():
                raise Exception("目录非子目录,不可入库!")

            # 图层重名检查
            meta_list:list = json.loads(self.para.get("meta").__str__())
            check_meta_only = int(self.para.get("check_meta_only",0))

            res["data"] = {}
            if check_meta_only:

                database = Database.query.filter_by(guid=self.para.get("database_guid")).one_or_none()
                if not database:
                    raise Exception("数据库不存在!")
                pg_ds: DataSource = PGUtil.open_pg_data_source(1, DES.decode(database.sqlalchemy_uri))

                res["result"] = True

                for meta in meta_list:
                    layers:dict = meta.get("layer")

                    for layer_name_origin in layers.keys():

                        layer_name = layers.get(layer_name_origin)
                        if pg_ds.GetLayerByName(layer_name) or InsertingLayerName.query.filter_by(name=layer_name).one_or_none():
                            res["data"][layer_name]=0
                            res["result"] = False
                        # 判断特殊字符
                        elif re.search(r"\W",layer_name):
                            res["data"][layer_name]=-1
                            res["result"] = False
                        else :
                            res["data"][layer_name] = 1

                if pg_ds:
                    try:
                        pg_ds.Destroy()
                    except:
                       print("关闭数据库失败!")
                return res




            # 录入数据后台进程,录入主函数为entry
            # 初始化task
            task = Task(guid=self.para.get("task_guid"),
                        name="入库 | {}".format(self.para.get("task_name")),
                        create_time=datetime.datetime.now(),
                        state=0,
                        task_type=1,
                        creator=self.para.get("creator"),
                        file_name=meta_list[0].get("filename"),
                        database_guid=self.para.get("database_guid"),
                        catalog_guid=self.para.get("catalog_guid"),
                        process="等待入库",
                        parameter=json.dumps(self.para),
                        # task_pid=entry_thread.pid
                        )
            db.session.add(task)
            db.session.commit()


            entry_thread = multiprocessing.Process(target=self.entry,args=(self.para.get("task_guid"),))
            entry_thread.start()

            Task.query.filter_by(guid=task_guid).update({"task_pid": entry_thread.pid})
            db.session.commit()

            res["result"] = True
            res["msg"] = "数据录入提交成功!"
            res["data"] = self.para["task_guid"]
        except Exception as e:
            raise e
        return res

    def entry(self,task_guid):

        task_writer = None
        this_task_layer = []
        try:

            #任务控制,等待执行
            TaskController.wait(task_guid)

            task_writer = TaskWriter(task_guid)

            task:Task = task_writer.session.query(Task).filter_by(guid=task_guid).one_or_none()
            parameter = json.loads(task.parameter)
            task_writer.update_task({"state": 2, "process": "入库中"})

            #处理修改入库信息
            metas: list = json.loads(parameter.get("meta").__str__())
            parameter["meta"] = metas

            database = task_writer.session.query(Database).filter_by(guid=task.database_guid).one_or_none()
            pg_ds: DataSource = PGUtil.open_pg_data_source(1, DES.decode(database.sqlalchemy_uri))


            for meta in metas:
                overwrite = parameter.get("overwrite", "no")

                for layer_name_origin, layer_name in meta.get("layer").items():
                    origin_name = layer_name
                    no = 1

                    while (overwrite.__eq__("no") and pg_ds.GetLayerByName(layer_name)) or task_writer.session.query(
                            InsertingLayerName).filter_by(name=layer_name).one_or_none():
                        layer_name = origin_name + "_{}".format(no)
                        no += 1

                    # 添加到正在入库的列表中
                    iln = InsertingLayerName(guid=uuid.uuid1().__str__(),
                                             task_guid=task.guid,
                                             name=layer_name)

                    task_writer.session.add(iln)
                    task_writer.session.commit()
                    this_task_layer.append(layer_name)
                    # 修改表名
                    meta["layer"][layer_name_origin] = layer_name
            pg_ds.Destroy()

            #入库
            EntryDataVacuate().entry(parameter)

            #完成后
            for ln in this_task_layer:
                iln = task_writer.session.query(InsertingLayerName).filter_by(name=ln).one_or_none()
                task_writer.session.delete(iln)

        except Exception as e:
            StructurePrint().print(e.__str__(), "error")
            task_writer.update_task({"state": -1, "process": "入库失败"})
            for ln in this_task_layer:
                iln = task_writer.session.query(InsertingLayerName).filter_by(name=ln).one_or_none()
                task_writer.session.delete(iln)
            task_writer.update_process(e.__str__())
            task_writer.update_process("任务中止!")
            StructurePrint().print(e.__str__(), "error")
        finally:
            task_writer.session.commit()
            task_writer.close()


    api_doc={
    "tags":["IO接口"],
    "parameters":[
        {"name": "meta",
         "in": "formData",
         "type": "string",
         "description": "数据meta"},
        {"name": "encoding",
         "in": "formData",
         "type": "string",
         "description": "原shp文件编码,非必要,优先使用cpg文件中编码,没有则默认GBK","enum":["UTF-8","GBK"]},
        {"name": "overwrite",
         "in": "formData",
         "type": "string",
         "description": "是否覆盖",
         "enum":["yes","no"]},
        {"name": "fid",
         "in": "formData",
         "type": "string",
         "description": "fid列名"},
        {"name": "geom_name",
         "in": "formData",
         "type": "string",
         "description": "空间属性列名"},
        {"name": "task_name",
         "in": "formData",
         "type": "string",
         "description": "任务名",
         "required":"true"},
        {"name": "creator",
         "in": "formData",
         "type": "string",
         "description": "创建人"},
        {"name": "database_guid",
         "in": "formData",
         "type": "string",
         "description": "数据库guid",
         "required": "true"},
        {"name": "catalog_guid",
         "in": "formData",
         "type": "string",
         "description": "目录guid"},

        {"name": "vacuate",
         "in": "formData",
         "type": "string",
         "description": "是否抽稀",
         "enum":[1,0]},

        {"name": "check_meta_only",
         "in": "formData",
         "type": "int",
         "description": "是否只检查meta","enum":[0,1]}

    ],
    "responses":{
        200:{
            "schema":{
                "properties":{
                }
            }
            }
        }
    }