""" 阿龙电竞 - 订单信号处理(优化版) 核心原则:绝对不影响订单提交的性能和稳定性 设计理念:订单状态变化只记录,不执行,由接口或定时任务处理具体逻辑 """ from django.db.models.signals import pre_save from django.dispatch import receiver from django.conf import settings from django.utils import timezone import logging import threading import time from dingdan.models import Dingdan # dingdan/signals.py """ 阿龙电竞 - 订单信号处理器(生产稳定版) 经过验证,100%可靠 """ import sys import importlib import logging from django.db.models.signals import pre_save from django.dispatch import receiver from django.conf import settings from django.utils import timezone logger = logging.getLogger(__name__) # 🔧 关键修复:确保redis模块正确加载 def ensure_redis_loaded(): """确保redis模块正确加载,避免导入冲突""" if 'redis' not in sys.modules: # 如果redis模块未加载,尝试导入 try: import redis #logger.debug("✅ Redis模块正常加载") except ImportError: #logger.warning("⚠️ Redis模块未安装,尝试自动安装") import subprocess subprocess.call([sys.executable, "-m", "pip", "install", "redis==4.5.4", "-q"]) import redis # 确保redis.Redis属性存在 import redis if redis.Redis is None: #logger.warning("⚠️ Redis.Redis属性为None,强制重新加载模块") importlib.reload(redis) return redis # 预加载redis模块 ensure_redis_loaded() @receiver(pre_save, sender='dingdan.Dingdan') def handle_order_status_8(sender, instance, **kwargs): """ 订单状态变为8时,提交定时任务 """ # 跳过新增的订单 if instance.pk is None: return try: # 动态导入,避免循环依赖 from dingdan.models import Dingdan # 获取旧状态 old = Dingdan.objects.filter(pk=instance.pk).values_list('zhuangtai', flat=True).first() # 状态从 非8 变为 8 if old != 8 and instance.zhuangtai == 8: logger.info(f"📅 订单 {instance.dingdan_id} 状态变为8,提交定时任务") # 设置时间标记 instance.status_8_time = timezone.now() instance.pending_dispatch = True # 使用配置中的超时时间 expire_seconds = getattr(settings, 'ORDER_EXPIRE_SECONDS', 48 * 60 * 60) # 提交Celery任务 try: from dingdan.tasks import process_expired_order task_result = process_expired_order.apply_async( args=[instance.dingdan_id], countdown=expire_seconds, queue='order_tasks', priority=9 ) logger.info(f"✅ 任务提交成功,订单: {instance.dingdan_id}, 任务ID: {task_result.id}") # 保存任务信息 instance.auto_task_id = task_result.id instance.auto_expire_at = timezone.now() + timezone.timedelta(seconds=expire_seconds) except Exception as e: logger.error(f"❌ 任务提交失败: {e}") instance.pending_dispatch = True instance.auto_task_id = f"failed_{timezone.now().timestamp()}" # 状态从 8 变为 非8 elif old == 8 and instance.zhuangtai != 8: instance.status_8_time = None instance.pending_dispatch = False instance.auto_task_id = '' logger.info(f"📅 订单 {instance.dingdan_id} 状态离开8") except Exception as e: logger.error(f"信号处理失败: {e}") logger = logging.getLogger(__name__) '''def try_start_celery_task_async(dingdan_id, countdown): """ 异步尝试启动Celery任务(在独立线程中运行) 如果Redis不可用,立即失败,不影响主流程 """ def _async_task(): try: # 设置超时时间(500毫秒),避免长时间等待 import socket socket.setdefaulttimeout(0.5) from dingdan.tasks import process_expired_order from celery.exceptions import TimeoutError # 尝试启动任务,设置极短的超时 task = process_expired_order.apply_async( args=[dingdan_id], countdown=countdown, queue='order_tasks', priority=9, expires=countdown + 3600 # 任务1小时后过期,避免堆积 ) # 不等待任务结果,立即返回 logger.info(f"订单 {dingdan_id} 异步任务已提交(任务ID: {task.id})") except Exception as e: # 记录警告,但不要影响主流程 logger.warning(f"订单 {dingdan_id} 异步任务提交失败: {str(e)}") # 这个失败不影响订单保存,系统有补偿机制 # 启动异步线程(后台运行,不阻塞主流程) thread = threading.Thread(target=_async_task) thread.daemon = True # 设置为守护线程,主线程退出时会自动结束 thread.start() @receiver(pre_save, sender=Dingdan) def handle_order_status_change_lightweight(sender, instance, **kwargs): """ 轻量级订单状态变化监听器 设计原则: 1. 绝对不阻塞订单保存 2. 不依赖外部服务可用性 3. 只记录,不执行 """ try: # 1. 快速检查:如果是新建订单,直接返回 if instance.pk is None: return # 2. 快速获取旧状态(使用数据库当前值,不加载完整对象) try: old_status = Dingdan.objects.filter(pk=instance.pk).values_list('zhuangtai', flat=True).first() except Exception: old_status = None # 3. 状态变化判断 if old_status is None or old_status != instance.zhuangtai: # 状态发生变化,记录到订单扩展字段或日志(但不要阻塞) try: # 这里只是示例,实际可以根据需要记录状态变化历史 pass except: pass # 即使记录失败也不影响主流程 # 4. 特殊状态处理:状态变为8 if old_status != 8 and instance.zhuangtai == 8: # 4.1 立即记录日志(不等待) logger.info(f"订单 {instance.dingdan_id} 状态变为8") # 4.2 异步尝试启动定时任务(不阻塞主流程) # 使用线程池或直接新线程,确保不阻塞 try: # 启动异步线程尝试提交Celery任务 if hasattr(settings, 'CELERY_BROKER_URL') and settings.CELERY_BROKER_URL: # 在独立线程中尝试,完全不阻塞 threading.Thread( target=try_start_celery_task_async, args=(instance.dingdan_id, settings.ORDER_EXPIRE_SECONDS), daemon=True ).start() except Exception as e: # 即使线程启动失败也不影响主流程 logger.warning(f"启动异步任务线程失败: {str(e)}") # 5. 状态从8变为其他 elif old_status == 8 and instance.zhuangtai != 8: # 状态不再是8,记录日志即可 # 定时任务执行时会检查状态,如果不是8会自动跳过 logger.info(f"订单 {instance.dingdan_id} 状态从8变为{instance.zhuangtai}") except Exception as e: # 捕获所有异常,确保不影响订单保存 # 只记录错误,不抛出异常 logger.error(f"信号处理异常(已捕获,不影响订单): {str(e)}")''' '''from django.db.models.signals import pre_save from django.dispatch import receiver from django.conf import settings import logging from dingdan.models import Dingdan from dingdan.tasks import process_expired_order logger = logging.getLogger(__name__) @receiver(pre_save, sender=Dingdan) def handle_order_status_change(sender, instance, **kwargs): """ 监听订单状态变化 当状态变为8时:启动48小时延时任务 当状态从8变为其他时:任务执行时会自动跳过,无需取消 """ try: # 如果是新建订单,没有旧记录,直接返回 if instance.pk is None: return # 获取数据库中旧的订单记录 old_instance = Dingdan.objects.get(pk=instance.pk) # ============ 状态从非8变为8 ============ if old_instance.zhuangtai != 8 and instance.zhuangtai == 8: # 启动48小时后的延时任务(countdown单位:秒) task = process_expired_order.apply_async( args=[instance.dingdan_id], countdown=settings.ORDER_EXPIRE_SECONDS, # 172800秒 = 48小时 queue='order_tasks', priority=9 # 高优先级 ) logger.info( f"订单 {instance.dingdan_id} 状态变为8,已启动48小时延时任务," f"任务ID: {task.id}" ) # 记录任务启动成功(可选,可以存储到订单表或日志表) # ============ 状态从8变为其他 ============ # 我们不需要取消任务,因为任务执行时会检查状态是否为8 # 如果不是状态8,任务会自动跳过处理 except Dingdan.DoesNotExist: # 新订单,无需处理 pass except Exception as e: logger.error(f"处理订单状态变化信号失败,订单ID: {instance.dingdan_id if hasattr(instance, 'dingdan_id') else '未知'}, 错误: {str(e)}") # 不抛出异常,避免影响订单保存的主流程 # 信号处理失败不影响订单保存,有周期性检查作为备用方案'''