""" Celery任务定义 """ from celery import current_task from sqlalchemy import create_engine import json from app.celery_app import celery_app from app.config import settings from app.logger import get_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_logger("celery_tasks") # 创建同步数据库连接用于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(): """ 推送数据给DTC的Celery任务 自动获取一条未完成的回调请求进行处理 """ task_name = 'push_data_to_dtc' logger.info(f"🌿 开始推送数据给DTC任务") # 获取分布式锁,使用任务名称作为锁标识 lock = redis_manager.create_lock(f"celery_task:{task_name}", timeout=300) # 5分钟超时 try: # 尝试获取锁 if not lock.acquire(blocking=False): logger.warning(f"⚠️ 任务 {task_name} 正在执行中,跳过本次执行") return {"status": "skipped", "message": f"任务 {task_name} 正在执行中,跳过本次执行"} logger.info(f"🔒 成功获取任务 {task_name} 的分布式锁") try: with engine.connect() as conn: # 更新任务状态 current_task.update_state( state='PROGRESS', meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'} ) # 获取一条未完成的回调请求(按创建时间取最小值) 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("查询未完成的回调请求失败,任务停止执行") 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}") # 更新任务状态 current_task.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}失败,任务停止执行") # 更新任务状态 current_task.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": "没有有效的通过记录需要处理"} # 更新任务状态 current_task.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 # 更新任务状态 current_task.update_state( state='PROGRESS', meta={'current': 85, 'total': 100, 'status': '开始推送数据给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) # 更新任务状态 current_task.update_state( state='PROGRESS', meta={'current': 100, 'total': 100, 'status': '标记回调请求日志为已完成(ID={callback_log_id}, site_id={site_id})'} ) return { "status": "completed", "message": "任务完成" } except Exception as e: logger.error(f"❌ 推送数据给DTC时发生错误: {e}", exc_info=True) return {"status": "error", "message": str(e)} finally: # 释放分布式锁 try: lock.release() logger.info(f"🔓 释放任务 {task_name} 的分布式锁") except Exception as e: logger.error(f"❌ 释放任务 {task_name} 的分布式锁失败: {e}") except Exception as e: logger.error(f"❌ 获取分布式锁失败: {e}", exc_info=True) return {"status": "error", "message": f"获取分布式锁失败: {str(e)}"}