正在显示
5 个修改的文件
包含
26 行增加
和
32 行删除
... | ... | @@ -44,7 +44,6 @@ class Api(ApiTemplate): |
44 | 44 | file_name=None, |
45 | 45 | process="数据下载中", |
46 | 46 | database_guid=self.para.get("database_guid"), |
47 | - # task_pid=download_process.pid | |
48 | 47 | ) |
49 | 48 | |
50 | 49 | db.session.add(task) |
... | ... | @@ -53,9 +52,6 @@ class Api(ApiTemplate): |
53 | 52 | download_process = multiprocessing.Process(target=self.download, args=(task_guid,self.para)) |
54 | 53 | download_process.start() |
55 | 54 | |
56 | - Task.query.filter_by(guid=task_guid).update({"task_pid": download_process.pid}) | |
57 | - db.session.commit() | |
58 | - | |
59 | 55 | res["data"] = "下载任务已提交!" |
60 | 56 | |
61 | 57 | # 提示信息 |
... | ... | @@ -84,10 +80,12 @@ class Api(ApiTemplate): |
84 | 80 | |
85 | 81 | try: |
86 | 82 | |
83 | + task_writer = TaskWriter(task_guid) | |
84 | + pid = multiprocessing.current_process().pid | |
85 | + task_writer.update_task({"task_pid":pid}) | |
87 | 86 | #任务控制,等待执行 |
88 | 87 | |
89 | 88 | TaskController.wait(task_guid) |
90 | - task_writer = TaskWriter(task_guid) | |
91 | 89 | |
92 | 90 | task_writer.update_task({"state":2,"update_time":datetime.datetime.now(),"process" : "下载中"}) |
93 | 91 | task_writer.update_process("开始下载...") | ... | ... |
... | ... | @@ -101,9 +101,6 @@ class Api(ApiTemplate): |
101 | 101 | entry_thread = multiprocessing.Process(target=self.entry,args=(self.para.get("task_guid"),)) |
102 | 102 | entry_thread.start() |
103 | 103 | |
104 | - Task.query.filter_by(guid=task_guid).update({"task_pid": entry_thread.pid}) | |
105 | - db.session.commit() | |
106 | - | |
107 | 104 | res["result"] = True |
108 | 105 | res["msg"] = "数据录入提交成功!" |
109 | 106 | res["data"] = self.para["task_guid"] |
... | ... | @@ -116,11 +113,13 @@ class Api(ApiTemplate): |
116 | 113 | task_writer = None |
117 | 114 | this_task_layer = [] |
118 | 115 | try: |
116 | + task_writer = TaskWriter(task_guid) | |
117 | + pid = multiprocessing.current_process().pid | |
118 | + task_writer.update_task({"task_pid": pid}) | |
119 | 119 | |
120 | 120 | #任务控制,等待执行 |
121 | - TaskController.wait(task_guid) | |
122 | 121 | |
123 | - task_writer = TaskWriter(task_guid) | |
122 | + TaskController.wait(task_guid) | |
124 | 123 | |
125 | 124 | task:Task = task_writer.session.query(Task).filter_by(guid=task_guid).one_or_none() |
126 | 125 | parameter = json.loads(task.parameter) | ... | ... |
... | ... | @@ -43,8 +43,7 @@ class Api(ApiTemplate): |
43 | 43 | creator=self.para.get("creator"), |
44 | 44 | file_name=None, |
45 | 45 | database_guid=database.guid, |
46 | - process="数据库更新中", | |
47 | - # task_pid=refresh_process.pid | |
46 | + process="数据库更新中" | |
48 | 47 | ) |
49 | 48 | |
50 | 49 | db.session.add(task) |
... | ... | @@ -53,13 +52,12 @@ class Api(ApiTemplate): |
53 | 52 | refresh_process = multiprocessing.Process(target=self.table_refresh,args=(database,task_guid,self.para.get("creator"))) |
54 | 53 | refresh_process.start() |
55 | 54 | |
56 | - Task.query.filter_by(guid=task_guid).update({"task_pid":refresh_process.pid}) | |
57 | - db.session.commit() | |
58 | 55 | |
59 | 56 | res["msg"] = "数据库更新已提交!" |
60 | 57 | res["data"] = task_guid |
61 | 58 | res["result"] = True |
62 | 59 | |
60 | + | |
63 | 61 | except Exception as e: |
64 | 62 | print(traceback.format_exc()) |
65 | 63 | raise e |
... | ... | @@ -69,24 +67,24 @@ class Api(ApiTemplate): |
69 | 67 | def table_refresh(self,database,task_guid,creator): |
70 | 68 | |
71 | 69 | pg_ds =None |
72 | - sys_ds =None | |
70 | + # sys_ds =None | |
73 | 71 | data_session=None |
74 | 72 | result = {} |
75 | 73 | sys_session = None |
76 | 74 | task_writer = None |
77 | 75 | |
78 | - | |
79 | 76 | db_tuple = PGUtil.get_info_from_sqlachemy_uri(DES.decode(database.sqlalchemy_uri)) |
80 | 77 | |
81 | 78 | try: |
82 | - | |
83 | - #任务控制,等待执行 | |
84 | - TaskController.wait(task_guid) | |
85 | 79 | task_writer = TaskWriter(task_guid) |
80 | + pid = multiprocessing.current_process().pid | |
81 | + task_writer.update_task({"task_pid":pid}) | |
86 | 82 | |
83 | + # 任务控制,等待执行 | |
84 | + TaskController.wait(task_guid) | |
87 | 85 | |
88 | 86 | sys_session = PGUtil.get_db_session(configure.SQLALCHEMY_DATABASE_URI) |
89 | - sys_ds = PGUtil.open_pg_data_source(0,configure.SQLALCHEMY_DATABASE_URI) | |
87 | + # sys_ds = PGUtil.open_pg_data_source(0,configure.SQLALCHEMY_DATABASE_URI) | |
90 | 88 | # 实体库连接 |
91 | 89 | data_session: Session = PGUtil.get_db_session(DES.decode(database.sqlalchemy_uri)) |
92 | 90 | |
... | ... | @@ -113,7 +111,6 @@ class Api(ApiTemplate): |
113 | 111 | # 空间表处理完毕 |
114 | 112 | sys_session.commit() |
115 | 113 | |
116 | - | |
117 | 114 | # 注册普通表 |
118 | 115 | |
119 | 116 | |
... | ... | @@ -173,8 +170,8 @@ class Api(ApiTemplate): |
173 | 170 | data_session.close() |
174 | 171 | if sys_session: |
175 | 172 | sys_session.close() |
176 | - if sys_ds: | |
177 | - sys_ds.Destroy() | |
173 | + # if sys_ds: | |
174 | + # sys_ds.Destroy() | |
178 | 175 | return result |
179 | 176 | |
180 | 177 | def add_spatail_table(self,database,pg_ds,data_session,sys_session,spatial_tables_names,this_time,db_tuple,creator): | ... | ... |
... | ... | @@ -92,9 +92,6 @@ class Api(ApiTemplate): |
92 | 92 | vacuate_process = multiprocessing.Process(target=self.task,args=(table,task_guid)) |
93 | 93 | vacuate_process.start() |
94 | 94 | |
95 | - Task.query.filter_by(guid=task_guid).update({"task_pid": vacuate_process.pid}) | |
96 | - db.session.commit() | |
97 | - | |
98 | 95 | res["msg"] = "矢量金字塔构建已提交!" |
99 | 96 | res["data"] = task_guid |
100 | 97 | res["result"] = True |
... | ... | @@ -114,11 +111,13 @@ class Api(ApiTemplate): |
114 | 111 | vacuate_process = None |
115 | 112 | try: |
116 | 113 | |
114 | + task_writer = TaskWriter(task_guid) | |
115 | + pid = multiprocessing.current_process().pid | |
116 | + task_writer.update_task({"task_pid":pid}) | |
117 | + | |
117 | 118 | #任务控制,等待执行 |
118 | 119 | TaskController.wait(task_guid) |
119 | 120 | |
120 | - task_writer = TaskWriter(task_guid) | |
121 | - | |
122 | 121 | task_writer.update_table(table.guid,{"is_vacuate": 2, "update_time": datetime.datetime.now()}) |
123 | 122 | task_writer.update_task({"state":2,"process":"构建中"}) |
124 | 123 | task_writer.update_process("开始构建...") | ... | ... |
... | ... | @@ -98,9 +98,6 @@ class Api(ApiTemplate): |
98 | 98 | vacuate_process = multiprocessing.Process(target=self.task,args=(table,task_guid,grids)) |
99 | 99 | vacuate_process.start() |
100 | 100 | |
101 | - Task.query.filter_by(guid=task_guid).update({"task_pid": vacuate_process.pid}) | |
102 | - db.session.commit() | |
103 | - | |
104 | 101 | res["msg"] = "矢量金字塔构建已提交!" |
105 | 102 | res["data"] = task_guid |
106 | 103 | res["result"] = True |
... | ... | @@ -122,9 +119,13 @@ class Api(ApiTemplate): |
122 | 119 | |
123 | 120 | try: |
124 | 121 | |
122 | + | |
123 | + task_writer = TaskWriter(task_guid) | |
124 | + pid = multiprocessing.current_process().pid | |
125 | + task_writer.update_task({"task_pid":pid}) | |
126 | + | |
125 | 127 | #任务控制,等待执行 |
126 | 128 | TaskController.wait(task_guid) |
127 | - task_writer = TaskWriter(task_guid) | |
128 | 129 | |
129 | 130 | task_writer.update_table(table.guid,{"is_vacuate": 2, "update_time": datetime.datetime.now()}) |
130 | 131 | task_writer.update_task({"state":2,"process":"构建中"}) | ... | ... |
请
注册
或
登录
后发表评论