"""任务处理服务""" import sys import os sys.path.append(os.path.abspath(os.path.join(os.path.dirname(__file__), os.path.pardir))) from datetime import datetime from common.helpers import get_f_block_positions import time # from busisness.blls import PDRecordBll class PD_TaskService: def __init__(self): # self.api_client = APIClient() self.half_volume = [0, 0] self.task_before = {"block_number":None, "beton_volume":None, "artifact_id":None} self.artifact_timestamps = {} # self.tcp_server = tcp_server # self.pd_record_bll=PDRecordBll() # from config.settings import TCP_CLIENT_HOST, TCP_CLIENT_PORT # self.data_client = OPCClient(url=f'opc.tcp://{TCP_CLIENT_HOST}:{TCP_CLIENT_PORT}') # self.data_client.start() def process_not_pour_info(self,artifact_list): """处理未浇筑信息""" #artifact_list[{'ArtifactActionID','BetonTaskID','BetonVolume','BetonTaskID'},{},{}] # artifact_list = self.api_client.get_not_pour_info() #mould_code 当前正在浇筑的模具编号 #artifact_list = self.pd_record_bll.get_last_pds(mould_code,3) # 如果API调用失败,返回空列表而不是抛出异常 if artifact_list is None: return [], [], [], self.half_volume if not artifact_list: return [], [], [], self.half_volume artifact_list = [person.__dict__ for person in artifact_list] # 处理F块信息 f_blocks_info = self._process_f_blocks(artifact_list) f_blocks = f_blocks_info["f_blocks"] f_block_count = f_blocks_info["f_block_count"] total_f_volume = f_blocks_info["total_f_volume"] f_positions = f_blocks_info["f_positions"] # 处理当前任务 #task{BetonGrade,MixID} current_task = self._process_current_task(artifact_list[0]) # 如果获取任务信息失败,则返回空结果 if current_task is None or not current_task.get("artifact_id"): return [], [], [], self.half_volume # 根据F块情况处理任务 task_result = self._handle_tasks_by_f_blocks( f_block_count, f_positions, current_task, f_blocks, total_f_volume, artifact_list ) # 更新上一个任务信息 self.task_before = { "beton_task_id": current_task["beton_task_id"], "beton_volume": round(current_task["beton_volume"],1), "artifact_id": current_task["artifact_id"], "block_number": current_task["block_number"] } return task_result def _process_f_blocks(self, artifact_list): """处理F块相关信息""" f_blocks = [artifact for artifact in artifact_list if artifact.get("BlockNumber") == "F"] f_block_count = len(f_blocks) total_f_volume = sum(artifact["BetonVolume"] for artifact in f_blocks) f_positions = get_f_block_positions(artifact_list) return { "f_blocks": f_blocks, "f_block_count": f_block_count, "total_f_volume": total_f_volume, "f_positions": f_positions } def _process_current_task(self, latest_artifact): """处理当前任务信息""" #task_data = self.api_client.get_task_info(latest_artifact["BetonTaskID"]) # 如果API调用失败,返回None #if task_data is None: # return None #"beton_task_id": latest_artifact["BetonTaskID"], return { "beton_task_id": latest_artifact["TaskID"], "beton_volume": latest_artifact["BetonVolume"], "artifact_id": latest_artifact["ArtifactActionID"], "block_number": latest_artifact.get("BlockNumber", ""), "produce_mix_id": latest_artifact.get("ProduceMixID", ""), "planned_volume": latest_artifact.get("PlannedVolume", 0), "id": latest_artifact.get("ID", 0), "task_data": latest_artifact } def _handle_tasks_by_f_blocks(self, f_block_count, f_positions, current_task, f_blocks, total_f_volume, artifact_list): """根据F块数量和位置处理任务""" # 多个F块情况 if f_block_count > 2: return self._handle_multiple_f_blocks(current_task, total_f_volume, artifact_list) # 两个F块情况 elif f_block_count == 2: return self._handle_two_f_blocks(f_positions, current_task, total_f_volume, artifact_list) # 一个F块情况 elif f_block_count == 1: return self._handle_single_f_block(f_positions, current_task, f_blocks, total_f_volume, artifact_list) # 无F块情况 elif f_block_count == 0: return self._handle_no_f_blocks(current_task, artifact_list) else: print("报警") return [], [], [], self.half_volume def _handle_multiple_f_blocks(self, current_task, total_f_volume, artifact_list): """处理多个F块的情况""" if self.task_before.get("block_number") == "F": print("报警:,超出正常补块逻辑") else: adjusted_volume = total_f_volume - self.half_volume[0] tasks = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": adjusted_volume, "artifact_id": current_task["artifact_id"], "block_number": current_task["block_number"], "produce_mix_id": current_task.get("produce_mix_id", ""), "planned_volume": current_task.get("planned_volume", 0), "id": current_task.get("id", 0) }] send_list = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": adjusted_volume, "artifact_id": current_task["artifact_id"], "block_number": current_task["block_number"], "beton_grade": current_task["task_data"]["BetonGrade"], "mix_id": current_task["task_data"]["ProduceMixID"], "time": self.artifact_timestamps.get(current_task["artifact_id"], datetime.now()) }] # 更新时间戳 #self._update_artifact_timestamps(send_list) return tasks, artifact_list, send_list, self.half_volume def _handle_two_f_blocks(self, f_positions, current_task, total_f_volume, artifact_list): """处理两个F块的情况""" if f_positions == [0, 1] and self.task_before.get("block_number") == "F": adjusted_volume = 0 block_number = "补方" elif f_positions == [1, 2]: adjusted_volume = artifact_list[0]["BetonVolume"] block_number = current_task["block_number"] elif f_positions == [0, 1] and self.task_before.get("block_number") != "F": adjusted_volume = total_f_volume - self.half_volume[0] block_number = current_task["block_number"] else: print("报警") tasks = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": adjusted_volume, "artifact_id": current_task["artifact_id"], "block_number": block_number, "produce_mix_id": current_task.get("produce_mix_id", ""), "planned_volume": current_task.get("planned_volume", 0), "id": current_task.get("id", 0) }] send_list = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": adjusted_volume, "artifact_id": current_task["artifact_id"], "block_number": block_number, "beton_grade": current_task["task_data"]["BetonGrade"], "mix_id": current_task["task_data"]["ProduceMixID"], "time": self.artifact_timestamps.get(current_task["artifact_id"], datetime.now()) }] # 更新时间戳 #self._update_artifact_timestamps(send_list) return tasks, artifact_list, send_list, self.half_volume def _handle_single_f_block(self, f_positions, current_task, f_blocks, total_f_volume, artifact_list): """处理单个F块的情况""" if f_positions == [2]: f_volume = f_blocks[0].get("BetonVolume") if f_blocks else 0 self.half_volume[0] = 0.1#round(total_f_volume / 2, 2)#暂时修改为0.1 self.half_volume[1] = 0.1#f_volume - self.half_volume[0]#暂时修改 adjusted_volume = round(current_task["beton_volume"],1) + self.half_volume[0] elif f_positions == [1]: adjusted_volume = round(current_task["beton_volume"],1) + self.half_volume[1] elif f_positions == [0]: adjusted_volume = 0 else: adjusted_volume = round(current_task["beton_volume"],1) block_number = "补方" if f_positions == [0] else current_task["block_number"] tasks = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": adjusted_volume, "artifact_id": current_task["artifact_id"], "block_number": block_number, "produce_mix_id": current_task.get("produce_mix_id", ""), "planned_volume": current_task.get("planned_volume", 0), "id": current_task.get("id", 0) }] send_list = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": adjusted_volume, "artifact_id": current_task["artifact_id"], "block_number": block_number, "beton_grade": current_task["task_data"]["BetonGrade"], "mix_id": current_task["task_data"]["ProduceMixID"], "time": self.artifact_timestamps.get(current_task["artifact_id"], datetime.now()) }] # 更新时间戳 #self._update_artifact_timestamps(send_list) return tasks, artifact_list, send_list, self.half_volume def _handle_no_f_blocks(self, current_task, artifact_list): """处理无F块的情况""" tasks = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": round(current_task["beton_volume"],1), "artifact_id": current_task["artifact_id"], "block_number": current_task["block_number"], "produce_mix_id": current_task.get("produce_mix_id", ""), "planned_volume": current_task.get("planned_volume", 0), "id": current_task.get("id", 0), }] send_list = [{ "beton_task_id": current_task["beton_task_id"], "beton_volume": round(current_task["beton_volume"],1), "artifact_id": current_task["artifact_id"], "block_number": current_task["block_number"], "beton_grade": current_task["task_data"]["BetonGrade"], "mix_id": current_task["task_data"]["ProduceMixID"], "time": self.artifact_timestamps.get(current_task["artifact_id"], datetime.now()) }] # 更新时间戳 #self._update_artifact_timestamps(send_list) return tasks, artifact_list, send_list, self.half_volume def _update_artifact_timestamps(self, send_list): """更新artifact时间戳""" current_artifact_ids = {item["artifact_id"] for item in send_list} for artifact_id in current_artifact_ids: if artifact_id not in self.artifact_timestamps: self.artifact_timestamps[artifact_id] = datetime.now() def insert_into_produce_table(self, connection, task_info, beton_volume, erp_id, artifact_id, half_volume, status): """插入数据到Produce表""" sql_db = SQLServerDB() if status == 1: # 准备插入数据 insert_data = { "ErpID": erp_id, "Code": task_info["TaskID"], "DatTim": datetime.now(), "Recipe": task_info["ProduceMixID"], "MorRec": "", "ProdMete": beton_volume, "MorMete": 0.0, # 砂浆方量,根据实际需求填写 "TotVehs": 0, # 累计车次,根据实际需求填写 "TotMete": task_info["PlannedVolume"], # 累计方量 "Qualitor": "", # 质检员,根据实际需求填写 "Acceptor": "", # 现场验收,根据实际需求填写 "Attamper": "", # 调度员,根据实际需求填写 "Flag": "1", # 标识,根据实际需求填写 "Note": "" # 备注,根据实际需求填写 } sql_db.insert_produce_data(insert_data) print(f"数据已成功插入到Produce表中,ERP ID: {erp_id}") try: # 参数包括: MISID(即erp_id), Flag, TaskID, ProduceMixID, ProjectName, BetonGrade, 调整后的方量 self.save_to_custom_table( misid=erp_id, flag="已插入", # 初始Flag值 task_id=task_info["TaskID"], produce_mix_id=task_info["ProduceMixID"], project_name=task_info["ProjectName"], beton_grade=task_info["BetonGrade"], adjusted_volume=round(beton_volume,1), artifact_id=artifact_id, half_volume=half_volume, # 已经调整后的方量 ) print(f"任务 {erp_id} 的数据已保存到自定义数据表") except Exception as e: print(f"调用保存函数时出错: {e}") if self.tcp_server: try: time.sleep(2) task_data = { "erp_id": erp_id, # 车号,相当于序号 "task_id": task_info["TaskID"], # 任务单号 "artifact_id": artifact_id, "produce_mix_id": task_info["ProduceMixID"], # 配比号 "project_name": task_info["ProjectName"], # 任务名 "beton_grade": task_info["BetonGrade"], # 砼强度 "adjusted_volume": round(beton_volume,1), # 方量 "flag": "已插入", # 状态 "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S") # 时间 } self.tcp_server.send_data(task_data) print(f"任务 {artifact_id} 的数据已发送给TCP客户端") except Exception as e: print(f"发送数据给TCP客户端时出错: {e}") return erp_id else: try: # 参数包括: MISID(即erp_id), Flag, TaskID, ProduceMixID, ProjectName, BetonGrade, 调整后的方量 self.save_to_custom_table( misid=erp_id, flag="已插入", # 初始Flag值 task_id=task_info["TaskID"], produce_mix_id=task_info["ProduceMixID"], project_name=task_info["ProjectName"], beton_grade=task_info["BetonGrade"], adjusted_volume=round(beton_volume,1), artifact_id=artifact_id, half_volume=half_volume, ) print(f"任务 {erp_id} 的数据已保存到自定义数据表") except Exception as e: print(f"调用保存函数时出错: {e}") return erp_id def save_to_custom_table(self, misid, flag, task_id, produce_mix_id, project_name, beton_grade, adjusted_volume, artifact_id,half_volume): try: task_data = { "erp_id": misid, "task_id": task_id, "artifact_id": artifact_id, "produce_mix_id": produce_mix_id, "project_name": project_name, "beton_grade": beton_grade, "adjusted_volume": adjusted_volume, "flag": flag, "half_volume":half_volume, "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S") } # self.data_client.send_data(task_data) print(f"任务 {artifact_id} 的数据已发送到另一台电脑") except Exception as e: print(f"发送数据到另一台电脑时出错: {e}") print(f"原计划保存到自定义数据表: MISID={misid}, Flag={flag}, TaskID={task_id}, 调整后方量={adjusted_volume}") def update_custom_table_status(self, erp_id, status): try: status_data = { "cmd": "update_status", "erp_id": erp_id, "status": status, "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S") } # self.data_client.send_data(status_data) print(f"任务状态更新已发送到另一台电脑: ERP ID={erp_id}, 状态={status}") except Exception as e: print(f"发送状态更新到另一台电脑时出错: {e}") print(f"原计划更新自定义数据表状态: ERP ID={erp_id}, 状态={status}") if __name__ == "__main__": pd_task_service = PD_TaskService() ret=pd_task_service.process_not_pour_info('SHR2B3-5') i=100