272 lines
9.9 KiB
Python
272 lines
9.9 KiB
Python
"""
|
||
阿龙电竞 - 订单信号处理(优化版)
|
||
核心原则:绝对不影响订单提交的性能和稳定性
|
||
设计理念:订单状态变化只记录,不执行,由接口或定时任务处理具体逻辑
|
||
"""
|
||
|
||
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)}")
|
||
# 不抛出异常,避免影响订单保存的主流程
|
||
# 信号处理失败不影响订单保存,有周期性检查作为备用方案''' |