""" Celery任务定义 """ from celery import current_task from sqlalchemy import create_engine, text import json from app.celery_app import celery_app from app.config import settings from app.logger import get_logger from app.database import CallbackFailureLog, ExternalApiLog, CallbackFailureData from app.callback_service import ( save_callback_data_items, get_uncompleted_callback_log, get_callback_log_data, check_phone_number_threshold, call_external_api_with_retry, mark_callback_log_completed, _log_dtc_push_call ) 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(self): """ 推送数据给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: # 获取一条未完成的回调请求(按创建时间取最小值) 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': len(valid_phone_numbers), 'total': len(phone_numbers), 'status': f'检查手机号: {phone_number}'} ) # 解析请求头和请求体 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}失败,任务停止执行") # 提取所有手机号并去重 phone_numbers_set = set() data_list = request_body.get('data', []) for item in data_list: number_data = item.get('number_data', {}) phone_number = number_data.get('number') if phone_number: phone_numbers_set.add(phone_number) phone_numbers = list(phone_numbers_set) logger.info(f"📱 提取到的手机号列表(去重后): {phone_numbers}") # 过滤出需要处理的手机号(不超过阈值的) valid_phone_numbers = [] for phone_number in phone_numbers: # 更新任务状态 current_task.update_state( state='PROGRESS', meta={'current': len(valid_phone_numbers), 'total': len(phone_numbers), 'status': f'检查手机号: {phone_number}'} ) # 检查手机号是否超过阈值 exceeds_threshold, query_success = check_phone_number_threshold(conn, phone_number) if not query_success: logger.error(f"❌ 查询手机号 {phone_number} 失败,跳过处理") continue if exceeds_threshold: logger.info(f"✅ 手机号 {phone_number} 出现次数超过阈值,跳过处理") continue valid_phone_numbers.append(phone_number) logger.info(f"📱 需要处理的手机号数量: {len(valid_phone_numbers)}") # 推送数据给DTC(在循环外执行一次) processed_count = 0 skipped_count = len(phone_numbers) - len(valid_phone_numbers) if valid_phone_numbers: # 更新任务状态 current_task.update_state( state='PROGRESS', meta={'current': 0, 'total': len(valid_phone_numbers), 'status': '开始推送数据给DTC'} ) success, retry_count = call_external_api_with_retry( conn=conn, request_body=request_body, max_retries=settings.external_api_retry_max, callback_failure_log_id=callback_log_id ) if success: logger.info(f"✅ 推送数据给DTC成功,处理的手机号数量: {len(valid_phone_numbers)}, 重试次数: {retry_count}") processed_count = len(valid_phone_numbers) else: logger.error(f"❌ 推送数据给DTC失败,重试次数: {retry_count}") else: logger.info("📋 没有有效的手机号需要处理") logger.info(f"🎉 推送数据给DTC处理完成,成功: {processed_count}, 跳过: {skipped_count}") # 标记CallbackFailureLog为已完成 mark_callback_log_completed(conn, callback_log_id) return { "status": "completed", "processed": processed_count, "skipped": skipped_count, "total": len(phone_numbers) } 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)}"}