""" 阿龙电竞 - 订单定时任务 处理订单超时自动确认等逻辑 完整版本,包含所有必要的异常处理 """ from celery import shared_task from django.db import transaction from django.db.models import F, Q from django.conf import settings from django.utils import timezone from datetime import timedelta from decimal import Decimal import logging # 导入模型 - 确保路径正确 from dingdan.models import Dingdan, DingdanShangjia from yonghu.models import UserMain, UserDashou, UserShangjia from utils.celery_utils import safe_decimal_operation, log_task_execution, rollback_on_failure from utils.celery_guard import celery_non_broadcast_blocked logger = logging.getLogger(__name__) @shared_task(bind=True, max_retries=3, default_retry_delay=60) def process_expired_order(self, dingdan_id): """ 处理单个超时订单(订单状态=8,超过48小时未处理) """ blocked = celery_non_broadcast_blocked() if blocked: logger.info('跳过 process_expired_order(%s): %s', dingdan_id, blocked) return blocked if not getattr(settings, 'ORDER_AUTO_SETTLE_ENABLED', False): logger.info('订单自动结算已禁用,跳过 process_expired_order(%s)', dingdan_id) return f'自动结算已禁用,跳过订单{dingdan_id}' try: guanshi_settle_ctx = None # 使用select_for_update锁定记录,防止并发修改 with transaction.atomic(): # 获取订单,同时锁定记录 order = Dingdan.objects.select_for_update().get( dingdan_id=dingdan_id ) # 其余代码不变... # ... 保持原有代码 # 再次检查订单状态是否为8 if order.zhuangtai != 8: logger.info(f"订单{dingdan_id}状态已不是8(当前:{order.zhuangtai}),跳过处理") return f"订单{dingdan_id}状态已变更,跳过处理" # 获取发单平台类型 fadan_pingtai = order.fadan_pingtai # 获取打手ID和分成金额 dashou_id = order.jiedan_dashou_id dashou_fencheng = safe_decimal_operation(order.dashou_fencheng) # 1. 更新订单状态为3(已完成) order.zhuangtai = 3 order.save(update_fields=['zhuangtai']) try: from utils.chat_utils import schedule_order_status_chat_push schedule_order_status_chat_push(dingdan_id) except Exception: pass # 2. 更新打手信息(如果存在) updated_dashou = False if dashou_id: try: # 查询打手用户主表 dashou_user = UserMain.objects.filter(yonghuid=dashou_id).first() if dashou_user: # 获取打手扩展表 dashou_profile = getattr(dashou_user, 'dashou_profile', None) if dashou_profile: # 使用F表达式原子更新 UserDashou.objects.filter(pk=dashou_profile.pk).update( chengjiaozongliang=F('chengjiaozongliang') + 1, yue=F('yue') + dashou_fencheng, zonge=F('zonge') + dashou_fencheng, jinrishouyi=F('jinrishouyi') + dashou_fencheng, jinyueshouyi=F('jinyueshouyi') + dashou_fencheng ) updated_dashou = True logger.info(f"更新打手{dashou_id}信息成功") except Exception as e: logger.error(f"更新打手信息失败,打手ID: {dashou_id}, 错误: {str(e)}") # 继续执行,不中断流程 # 3. 如果是商家发单,更新商家信息 updated_shangjia = False if fadan_pingtai == 2: # 商家发单 try: # 获取商家扩展信息 shangjia_kuozhan = getattr(order, 'shangjia_kuozhan', None) if shangjia_kuozhan: shangjia_id = shangjia_kuozhan.shangjia_id if shangjia_id: # 查询商家用户 shangjia_user = UserMain.objects.filter(yonghuid=shangjia_id).first() if shangjia_user: shangjia_profile = getattr(shangjia_user, 'shop_profile', None) if shangjia_profile: # 更新商家成交订单数量 UserShangjia.objects.filter(pk=shangjia_profile.pk).update( chengjiao=F('chengjiao') + 1 ) updated_shangjia = True logger.info(f"更新商家{shangjia_id}信息成功") except Exception as e: logger.error(f"更新商家信息失败,订单ID: {dingdan_id}, 错误: {str(e)}") # 继续执行,不中断流程 # 记录任务执行结果 result_msg = f"订单{dingdan_id}自动结算完成" if updated_dashou: result_msg += ",打手信息已更新" if updated_shangjia: result_msg += ",商家信息已更新" log_task_execution(f"自动结算订单 {dingdan_id}", True, result_msg) if fadan_pingtai == 2 and dashou_id: guanshi_settle_ctx = (dingdan_id, dashou_id) if guanshi_settle_ctx: try: from dingdan.utils import settle_shangjia_order_guanshi_fenhong settle_order = Dingdan.objects.get(dingdan_id=guanshi_settle_ctx[0]) settle_shangjia_order_guanshi_fenhong(settle_order, guanshi_settle_ctx[1]) except Exception as e: logger.error(f"自动结算管事分红失败 订单{dingdan_id}: {e}", exc_info=True) return result_msg except Dingdan.DoesNotExist: error_msg = f"订单{dingdan_id}不存在" logger.warning(error_msg) return error_msg except Exception as e: error_msg = f"处理订单{dingdan_id}时发生错误: {str(e)}" logger.error(error_msg) # 重试机制 if self.request.retries < self.max_retries: logger.info(f"订单{dingdan_id}处理失败,准备第{self.request.retries + 1}次重试") raise self.retry(exc=e) log_task_execution(f"自动结算订单 {dingdan_id}", False, error_msg) return error_msg @shared_task(bind=True, max_retries=3, default_retry_delay=15) def retry_establish_order_chat_task(self, dingdan_id): """抢单后建聊失败时的延迟补偿任务""" blocked = celery_non_broadcast_blocked() if blocked: return blocked try: from utils.chat_utils import establish_order_chat_with_retry ok = establish_order_chat_with_retry(dingdan_id, max_rounds=5) if ok: logger.info(f"建聊补偿成功: 订单{dingdan_id}") return f'chat ok {dingdan_id}' logger.warning(f"建聊补偿仍失败: 订单{dingdan_id}") if self.request.retries < self.max_retries: raise self.retry() return f'chat failed {dingdan_id}' except Exception as e: logger.error(f"建聊补偿任务异常 订单{dingdan_id}: {e}", exc_info=True) if self.request.retries < self.max_retries: raise self.retry(exc=e) return str(e) @shared_task(bind=True, max_retries=2, default_retry_delay=10) def push_order_status_chat_task(self, dingdan_id): """订单状态变更后推送群聊状态卡片""" blocked = celery_non_broadcast_blocked() if blocked: return blocked try: from utils.chat_utils import push_order_status_chat_update ok = push_order_status_chat_update(dingdan_id) if ok: return f'status chat ok {dingdan_id}' if self.request.retries < self.max_retries: raise self.retry() return f'status chat failed {dingdan_id}' except Exception as e: logger.error(f"订单状态聊天推送异常 {dingdan_id}: {e}", exc_info=True) if self.request.retries < self.max_retries: raise self.retry(exc=e) return str(e) @shared_task @rollback_on_failure("批量检查超时订单") def check_order_expire_task(): """ 批量检查超时订单(补偿机制) 每5分钟执行一次,作为信号处理的备用方案 查找状态为8且创建时间超过72小时的订单(宽松时间,防止误处理) """ blocked = celery_non_broadcast_blocked() if blocked: logger.info('跳过 check_order_expire_task: %s', blocked) return {"success": True, "message": blocked, "count": 0} if not getattr(settings, 'ORDER_AUTO_SETTLE_ENABLED', False): logger.info('订单自动结算已禁用,跳过 check_order_expire_task') return {"success": True, "message": "自动结算已禁用", "count": 0} try: # 计算更宽松的时间点:当前时间减去72小时(48+24小时) # 这样避免误处理刚变为状态8的订单 check_time = timezone.now() - timedelta(seconds=settings.ORDER_EXPIRE_SECONDS + 86400) # 48+24=72小时 # 查询状态为8且创建时间早于检查时间的订单 expired_orders = Dingdan.objects.filter( zhuangtai=8, create_time__lt=check_time ).values_list('dingdan_id', flat=True) expired_count = expired_orders.count() if expired_count == 0: logger.info("补偿检查:无超时订单") return {"success": True, "message": "无超时订单", "count": 0} # 分批处理超时订单(避免一次处理太多) batch_size = 10 processed_count = 0 for i in range(0, expired_count, batch_size): batch = expired_orders[i:i+batch_size] # 异步处理每个超时订单 for dingdan_id in batch: process_expired_order.apply_async( args=[dingdan_id], queue='order_tasks', priority=8 # 补偿任务的优先级低一些 ) processed_count += 1 result_msg = f"补偿检查:发现{expired_count}个可能超时的订单,已提交处理{processed_count}个" log_task_execution("批量检查超时订单", True, result_msg) return { "success": True, "message": result_msg, "total_count": expired_count, "processed_count": processed_count } except Exception as e: error_msg = f"批量检查超时订单失败: {str(e)}" logger.error(error_msg) log_task_execution("批量检查超时订单", False, error_msg) return { "success": False, "message": error_msg, "error": str(e) }