From 851e92ad3a9dd4cd686b9e31e26fc54b51857ae1 Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Wed, 10 Dec 2025 12:18:28 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96=E5=9B=9E=E8=B0=83=E8=AF=B7?= =?UTF-8?q?=E6=B1=82=E5=A4=84=E7=90=86=E4=B8=BA=E5=B7=B2=E5=AE=8C=E6=88=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/callback_service.py | 22 +++++++++++-- app/celery_tasks.py | 68 +++++++++++++++++++++-------------------- 2 files changed, 54 insertions(+), 36 deletions(-) diff --git a/app/callback_service.py b/app/callback_service.py index c0b2474..7bb7ce9 100644 --- a/app/callback_service.py +++ b/app/callback_service.py @@ -327,11 +327,27 @@ async def mark_callback_log_completed(db: AsyncSession, callback_log_id: int) -> ) success_count = success_result.fetchone()[0] - if success_count == 0: - logger.warning(f"⚠️ 回调日志 {callback_log_id} 没有成功的推送记录,不标记为完成") + # 检查失败记录数量 + failure_result = await db.execute( + text(""" + SELECT COUNT(*) as failure_count + FROM external_api_logs + WHERE callback_failure_log_id = :log_id + AND response_status != 200 + """), + {"log_id": callback_log_id} + ) + failure_count = failure_result.fetchone()[0] + + # 如果有成功记录,或者失败记录达到最大重试次数,则标记为完成 + if success_count == 0 and failure_count < settings.external_api_retry_max: + logger.warning(f"⚠️ 回调日志 {callback_log_id} 没有成功的推送记录,且失败记录({failure_count}条)未达到最大重试次数({settings.external_api_retry_max}条),不标记为完成") return False - logger.info(f"✅ 回调日志 {callback_log_id} 找到 {success_count} 条成功推送记录,开始标记为完成") + if success_count > 0: + logger.info(f"✅ 回调日志 {callback_log_id} 找到 {success_count} 条成功推送记录,开始标记为完成") + else: + logger.info(f"✅ 回调日志 {callback_log_id} 失败记录({failure_count}条)已达到最大重试次数({settings.external_api_retry_max}条),开始标记为完成") # 更新指定日志记录为已完成 result = await db.execute( diff --git a/app/celery_tasks.py b/app/celery_tasks.py index 6bf026e..d9a2e1e 100644 --- a/app/celery_tasks.py +++ b/app/celery_tasks.py @@ -165,7 +165,7 @@ def push_data_to_dtc_task(self): 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": "没有有效的通过记录需要处理"} @@ -176,46 +176,48 @@ def push_data_to_dtc_task(self): 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 = await call_external_api_with_retry( - db=db, - 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为已完成 - await mark_callback_log_completed(db, callback_log_id) + # 由于已经确认有记录,直接进入推送逻辑 + # 创建只包含有效数据项的请求体 + filtered_request_body = request_body.copy() + filtered_request_body['data'] = related_records # 更新任务状态 self.update_state( state='PROGRESS', - meta={'current': 100, 'total': 100, 'status': f'标记回调请求日志为已完成(ID={callback_log_id}, site_id={site_id})'} + meta={'current': 85, 'total': 100, 'status': f'开始推送数据给DTC(ID={callback_log_id}, site_id={site_id})...'} ) + success, retry_count = await call_external_api_with_retry( + db=db, + 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}") + logger.info(f"🎉 推送数据给DTC处理完成 - 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}") + + # 标记CallbackFailureLog为已完成 + mark_success = await mark_callback_log_completed(db, callback_log_id) + + if not mark_success: + logger.error(f"❌ 标记回调请求日志状态失败,任务终止执行") + raise Exception("标记回调请求日志状态失败,任务终止执行") + + # 更新任务状态 + self.update_state( + state='PROGRESS', + meta={'current': 100, 'total': 100, 'status': f'标记回调请求日志状态(ID={callback_log_id}, site_id={site_id})'} + ) + + logger.info(f"✅ 回调日志 {callback_log_id} 已成功标记为完成") return { "status": "completed", - "message": "任务完成" + "message": "任务完成,回调日志已标记为完成" } except Exception as e: