Files
ai-talk-callback/app/celery_tasks.py
mark.tian d00bb5e475 测试
2025-12-09 10:01:00 +08:00

207 lines
9.3 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.
"""
Celery任务定义
"""
from sqlalchemy import create_engine
import json
from app.celery_app import celery_app
from app.config import settings
from app.logger import get_celery_tasks_logger
from app.callback_service import (
save_callback_data_items,
get_uncompleted_callback_log,
call_external_api_with_retry,
mark_callback_log_completed,
get_related_records_by_unique_data_list
)
from app.redis_lock import redis_manager
logger = get_celery_tasks_logger()
# 创建同步数据库连接用于Celery任务
engine = create_engine(
settings.database_url,
echo=settings.debug,
future=True
)
@celery_app.task(bind=True, name='push_data_to_dtc')
def push_data_to_dtc_task(self):
"""
推送数据给DTC的Celery任务
自动获取一条未完成的回调请求进行处理
"""
task_name = 'push_data_to_dtc'
logger.info(f"🌿 开始推送数据给DTC任务")
try:
# 获取分布式锁,使用任务名称作为锁标识
connection_success = redis_manager.connect()
if not connection_success:
logger.error(f"❌ Redis管理器连接失败,任务停止执行")
raise Exception("Redis连接失败")
lock = redis_manager.create_lock(f"celery_task:{task_name}", timeout=300) # 5分钟超时
# 尝试获取锁
if not lock.acquire(blocking=False):
logger.warning(f"⚠️ 任务 {task_name} 正在执行中,跳过本次执行")
return {"status": "skipped", "message": f"任务 {task_name} 正在执行中,跳过本次执行"}
logger.info(f"🔒 成功获取任务 {task_name} 的分布式锁")
with engine.connect() as conn:
logger.info("1")
# 更新任务状态
self.update_state(
state='PROGRESS',
meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'}
)
logger.info("2")
# 获取一条未完成的回调请求(按创建时间取最小值)
query_success, callback_log_data = get_uncompleted_callback_log(conn)
if not query_success:
logger.error("❌ 查询未完成的回调请求失败,尝试再查一次")
query_success, callback_log_data = get_uncompleted_callback_log(conn)
if not query_success:
logger.error("❌ 再次查询未完成的回调请求失败,停止任务执行")
raise Exception("查询未完成的回调请求失败,任务停止执行")
logger.info("3")
if not callback_log_data:
logger.info("📋 没有找到未完成的回调请求")
return {"status": "skipped", "message": "没有找到未完成的回调请求"}
callback_log_id, site_id, request_headers_json, request_body_json = callback_log_data
logger.info(f"📋 获取到未完成的回调请求: ID={callback_log_id}, site_id={site_id}")
# 更新任务状态
self.update_state(
state='PROGRESS',
meta={'current': 25, 'total': 100, 'status': f'获取回调请求成功(ID={callback_log_id}, site_id={site_id}),开始分析处理...'}
)
# 解析请求头和请求体
request_headers = json.loads(request_headers_json)
request_body = json.loads(request_body_json)
# 保存callback_data.data中的数据
data_list = request_body.get('data', [])
if data_list and len(data_list) > 0:
save_success = save_callback_data_items(conn, data_list, callback_log_id)
if not save_success:
logger.error(f"❌ 保存callback_data:{callback_log_id}失败,停止任务执行")
raise Exception(f"保存callback_data:{callback_log_id}失败,任务停止执行")
# 更新任务状态
self.update_state(
state='PROGRESS',
meta={'current': 50, 'total': 100, 'status': f'请求体中通话明细处理成功(ID={callback_log_id}, site_id={site_id}),开始过滤需要转发的通话记录...'}
)
# 根据手机号、task_id、user_id给data_list去重
data_list = request_body.get('data', [])
unique_data_list = []
seen_records = set()
original_count = len(data_list)
for item in data_list:
# 提取手机号
number_data = item.get('number_data', {})
phone_number = number_data.get('number', '')
# 提取task_id
task = item.get('task', {})
task_id = task.get('id', '')
# 提取user_id
user_id = item.get('user_id', '')
# 创建唯一标识
unique_key = (phone_number, task_id, user_id)
# 如果这个组合没见过,则添加到去重列表中
if unique_key not in seen_records:
seen_records.add(unique_key)
unique_data_list.append(item)
logger.info(f"🔄 数据去重完成: 原始数据 {original_count} 条,去重后 {len(unique_data_list)} 条 - callback_log_id: {callback_log_id}, site_id: {site_id}")
if len(unique_data_list) < original_count:
logger.info(f"🗑️ 移除了 {original_count - len(unique_data_list)} 条重复数据 - callback_log_id: {callback_log_id}, site_id: {site_id}")
# 查询相关数据:根据手机号、task_id、user_id作为条件,按创建时间排序获取前三条记录
query_success, related_records = get_related_records_by_unique_data_list(
conn, unique_data_list, callback_log_id, settings.count_threshold
)
if not query_success:
logger.error("❌ 查询相关记录失败,停止任务执行")
raise Exception("查询相关记录失败,任务停止执行")
# 判断是否有相关记录需要处理
if not related_records:
logger.info(f"📋 没有有效的通过记录需要处理,直接返回 - callback_log_id: {callback_log_id}, site_id={site_id}")
return {"status": "completed", "message": "没有有效的通过记录需要处理"}
# 更新任务状态
self.update_state(
state='PROGRESS',
meta={'current': 70, 'total': 100, 'status': f'需要推送的通过记录已获取成功(ID={callback_log_id}, site_id={site_id}),准备转发...'}
)
if related_records:
# 创建只包含有效数据项的请求体
filtered_request_body = request_body.copy()
filtered_request_body['data'] = related_records
# 更新任务状态
self.update_state(
state='PROGRESS',
meta={'current': 85, 'total': 100, 'status': f'开始推送数据给DTC(ID={callback_log_id}, site_id={site_id})...'}
)
success, retry_count = call_external_api_with_retry(
conn=conn,
request_body=filtered_request_body,
request_headers=request_headers,
max_retries=settings.external_api_retry_max,
callback_failure_log_id=callback_log_id
)
if success:
logger.info(f"✅ 推送数据给DTC成功,处理的数据项数量: {len(related_records)}, 重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}")
else:
logger.error(f"❌ 推送数据给DTC失败,重试次数: {retry_count} - callback_log_id: {callback_log_id}, site_id: {site_id}")
else:
logger.info(f"📋 没有有效的通过记录需要处理 - callback_log_id: {callback_log_id}, site_id: {site_id}")
logger.info(f"🎉 推送数据给DTC处理完成 - callback_log_id: {callback_log_id}, site_id: {site_id}")
# 标记CallbackFailureLog为已完成
mark_callback_log_completed(conn, callback_log_id)
# 更新任务状态
self.update_state(
state='PROGRESS',
meta={'current': 100, 'total': 100, 'status': f'标记回调请求日志为已完成(ID={callback_log_id}, site_id={site_id})'}
)
return {
"status": "completed",
"message": "任务完成"
}
except Exception as e:
logger.error(f"❌ 任务执行失败: {e}", exc_info=True)
return {"status": "error", "message": str(e)}
finally:
# 释放分布式锁
try:
if 'lock' in locals():
lock.release()
logger.info(f"🔓 释放任务 {task_name} 的分布式锁")
except Exception as e:
logger.error(f"❌ 释放任务 {task_name} 的分布式锁失败: {e}")