import requests import time import logging import requests from ewcc2 import Sinopec_ewcc # 配置日志 handlers = [logging.FileHandler('log\query_processor.log'), logging.StreamHandler()] logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s', handlers=handlers ) logger = logging.getLogger(__name__) # API基础URL BASE_URL = 'http://localhost:99' # 演示数据处理函数 def process_query_data(query_data, parameters): """ 处理查询数据的演示函数 这里使用演示代码代替实际的数据处理逻辑 """ # 获取用户ID,如果存在 user_id = query_data.get('user_id', '未知') logger.info(f"处理查询数据: ID={query_data['id']}, 车牌号={query_data['plate_number']}, 用户ID={user_id}") # 这里是演示处理逻辑 # 在实际应用中,这里应该是真实的数据处理代码 logger.info(f"使用参数进行处理: {parameters}") # 模拟处理过程 time.sleep(2) # 模拟处理耗时 # 返回模拟的处理结果(简单地将results_count设置为随机数1-5) # 实际应用中,应该根据真实处理结果返回相应的值 results_count = 3 # 演示用的固定结果数 logger.info(f"处理完成,结果数量: {results_count}") return results_count # 读取参数配置 def get_parameters(): """ 从后端API读取参数配置 """ try: url = f'{BASE_URL}/api/parameters' response = requests.get(url) response.raise_for_status() data = response.json() if data['success']: logger.info(f"成功获取参数配置: {data['data']}") return data['data'] else: logger.error(f"获取参数配置失败: {data['error']}") return None except Exception as e: logger.error(f"获取参数配置时发生错误: {str(e)}") return None # 获取待处理查询 def get_pending_queries(): """ 从后端API获取待处理查询 """ try: url = f'{BASE_URL}/api/pending-queries' response = requests.get(url) response.raise_for_status() data = response.json() if data['success']: logger.info(f"成功获取待处理查询,共 {data['count']} 条") return data['data'] else: logger.error(f"获取待处理查询失败: {data['error']}") return [] except Exception as e: logger.error(f"获取待处理查询时发生错误: {str(e)}") return [] # 更新查询结果 def update_query_result(query_id, results_count, cars_info=None): """ 更新查询结果到后端API,可选地包含车辆信息 参数: query_id: 查询ID results_count: 结果数量 cars_info: 可选,车辆信息列表,将被保存到car_owner表 返回: bool: 更新是否成功 """ try: url = f'{BASE_URL}/api/update-query-result' data = { 'id': query_id, 'results_count': results_count } # 如果提供了车辆信息,添加到请求数据中 if cars_info: data['cars_info'] = cars_info response = requests.post(url, json=data) response.raise_for_status() result = response.json() if result['success']: logger.info(f"成功更新查询结果: ID={query_id}, 结果数量={results_count}") return True else: logger.error(f"更新查询结果失败: {result['error']}") return False except Exception as e: logger.error(f"更新查询结果时发生错误: {str(e)}") return False # 保存param4值到后端API def save_param4(param4_value): """ 保存param4的值到后端API 如果param4_value是复杂对象,会尝试将其转换为JSON字符串 """ try: url = f'{BASE_URL}/api/save-param4' # 尝试将cookies对象转换为JSON字符串,以便能够序列化 import json try: # 首先尝试直接序列化,如果失败则进行转换 json.dumps(param4_value) # 如果能够直接序列化,就使用原始值 processed_param4 = param4_value except (TypeError, OverflowError): # 如果不能直接序列化,尝试将其转换为字典或字符串 if hasattr(param4_value, '__dict__'): # 如果是有__dict__属性的对象,转换为字典 processed_param4 = param4_value.__dict__ else: # 否则转换为字符串 processed_param4 = str(param4_value) data = { 'param4': processed_param4 } response = requests.post(url, json=data) response.raise_for_status() result = response.json() if result['success']: logger.info(f"成功保存param4值") return True else: logger.error(f"保存param4值失败: {result['error']}") return False except Exception as e: logger.error(f"保存param4值时发生错误: {str(e)}") return False # 主处理函数 def main(): logger.info("查询处理代码启动") # 主循环 while True: ewcc = Sinopec_ewcc() try: # 获取参数配置 parameters = get_parameters() if not parameters: logger.warning("无法获取参数配置,停止30秒后继续") time.sleep(30) continue # 登录浏览器 if not ewcc.Autologin(parameters['param1'],parameters['param2'],parameters['param3'],parameters['param4']): time.sleep(30) continue # 只留浏览器cookies save_param4( ewcc.tab.cookies()) # 获取待处理查询 pending_queries = get_pending_queries() # 处理每条查询 ewcc.query_car_info("湘AB5W21") for query in pending_queries: # 处理查询数据 rtn_cars_info=ewcc.query_car_info(query['plate_number'],query['phone']) #results_count = process_query_data(query, parameters) # 更新查询结果 update_query_result(query['id'], len(rtn_cars_info), rtn_cars_info) # 等待一段时间后再次检查 wait_time = 60 # 默认等待60秒 # 如果有parameters参数,可以根据参数调整等待时间 if parameters and 'polling_interval' in parameters: try: wait_time = int(parameters['polling_interval']) except (ValueError, TypeError): logger.warning(f"无效的轮询间隔参数: {parameters['polling_interval']}") logger.info(f"本次处理完成,等待 {wait_time} 秒后再次检查") time.sleep(wait_time) except KeyboardInterrupt: logger.info("查询处理器被用户中断") break except Exception as e: logger.error(f"处理过程中发生错误: {str(e)}") # 发生错误时,等待更短的时间后重试 time.sleep(30) logger.info("查询处理器停止") if __name__ == '__main__': main()