Files
QueryCarInfo2/query_processor.py
T
2025-10-15 17:37:40 +08:00

250 lines
8.5 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import requests
import time
from applog import logger
import requests
import datetime
from ewcc2 import Sinopec_ewcc
# API基础URL
BASE_URL = 'http://10.190.20.156: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"成功获取参数配置: ")
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
# 显示cars_info信息
logger.info(f"车辆信息: {cars_info}")
else:
logger.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("查询处理代码启动")
ewcc = Sinopec_ewcc()
i =0
# 主循环
while True:
try:
# 获取参数配置
parameters = get_parameters()
if not parameters:
logger.warning("无法获取配置参数,停止5秒后继续")
time.sleep(5)
continue
if not ewcc.browser.states.is_alive:
logger.info("浏览器连接已断开,尝试重新连接")
ewcc.quit( )
ewcc=Sinopec_ewcc()
continue
# 自动登录
if not ewcc.Autologin(parameters['param1'],parameters['param2'],parameters['param3'],parameters['param4']):
time.sleep(30)
continue
# 保留浏览器cookies
save_param4( ewcc.tab.cookies())
# 查询员工拉新权益会员明细
#if not ewcc.is_query_member and datetime.datetime.now().hour >= 3:
#logger.info("开始查询员工拉新权益会员明细")
#ewcc.is_query_member=ewcc.query_memeber()
# 查询员工拉新权益会员明细
if not ewcc.is_query_oil_fw_tj and datetime.datetime.now().hour >= 3:
logger.info("开始查询油站服务号应用统计")
ewcc.query_oil_fw_tj()
#ewcc.is_query_oil_fw_tj=ewcc.query_oil_fw_tj()
#continue
# 查询油站汇总本月累计
if not ewcc.is_query_oil_current_month_lj and datetime.datetime.now().hour >= 3:
logger.info("开始查询 按油站汇总本月累计")
ewcc.query_oil_current_month_lj()
#ewcc.is_query_oil_current_month_lj=ewcc.query_oil_current_month_lj()
continue
# 获取待处理查询
pending_queries = get_pending_queries()
# 处理每条查询
for query in pending_queries:
# 处理查询数据
rtn_cars_info=[]
if ewcc.query_car_info(query['plate_number'],query['phone'],rtn_cars_info) :
# 更新查询结果
logger.info(f"查询结果数量: {len(rtn_cars_info)}")
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 ewcc:
ewcc.quit()
if __name__ == '__main__':
main()