第一次提交:小程序前端最新版本
This commit is contained in:
272
dingdan/signals.py
Normal file
272
dingdan/signals.py
Normal file
@@ -0,0 +1,272 @@
|
||||
"""
|
||||
阿龙电竞 - 订单信号处理(优化版)
|
||||
核心原则:绝对不影响订单提交的性能和稳定性
|
||||
设计理念:订单状态变化只记录,不执行,由接口或定时任务处理具体逻辑
|
||||
"""
|
||||
|
||||
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)}")
|
||||
# 不抛出异常,避免影响订单保存的主流程
|
||||
# 信号处理失败不影响订单保存,有周期性检查作为备用方案'''
|
||||
Reference in New Issue
Block a user