614 lines
22 KiB
Python
614 lines
22 KiB
Python
"""
|
||
阿龙电竞 - 信号处理器(修复Redis导入问题)
|
||
"""
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
# 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}")
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
# 🔥 关键修复:在最开始强制正确导入redis模块
|
||
'''import sys
|
||
|
||
# 1. 如果redis模块已加载但有问题,先移除
|
||
if 'redis' in sys.modules:
|
||
print(f"🔧 移除已加载的redis模块: {sys.modules['redis']}")
|
||
del sys.modules['redis']
|
||
# 同时移除相关模块
|
||
for key in list(sys.modules.keys()):
|
||
if key.startswith('redis.') or key == 'redis':
|
||
del sys.modules[key]
|
||
|
||
# 2. 重新导入redis模块
|
||
try:
|
||
import importlib
|
||
redis = importlib.import_module('redis')
|
||
print(f"✅ 重新导入redis成功: {redis}")
|
||
print(f" Redis类: {redis.Redis}")
|
||
print(f" Redis类是否为None: {redis.Redis is None}")
|
||
|
||
# 验证Redis类是否可用
|
||
if redis.Redis is None:
|
||
print("❌ Redis类为None,尝试从源码级别修复")
|
||
# 强制重新加载
|
||
import importlib.util
|
||
import os
|
||
|
||
# 找到redis包路径
|
||
import pkgutil
|
||
redis_spec = pkgutil.find_loader('redis')
|
||
if redis_spec:
|
||
print(f" Redis包路径: {redis_spec.origin}")
|
||
# 强制重新加载
|
||
importlib.invalidate_caches()
|
||
redis = importlib.import_module('redis')
|
||
importlib.reload(redis)
|
||
print(f" Redis类(重载后): {redis.Redis}")
|
||
|
||
except Exception as e:
|
||
print(f"❌ 导入redis失败: {e}")
|
||
import traceback
|
||
traceback.print_exc()
|
||
# 尝试直接安装
|
||
print("尝试直接安装redis包...")
|
||
import subprocess
|
||
subprocess.call([sys.executable, "-m", "pip", "install", "redis==4.5.4", "-q"])
|
||
|
||
# 3. 确保kombu.transport.redis也能正确导入
|
||
try:
|
||
import kombu.transport.redis
|
||
print(f"✅ kombu.transport.redis导入成功")
|
||
except Exception as e:
|
||
print(f"❌ kombu.transport.redis导入失败: {e}")
|
||
# 修复:先导入redis,再导入kombu.transport.redis
|
||
if 'redis' not in sys.modules:
|
||
import redis
|
||
# 重新导入
|
||
import importlib
|
||
if 'kombu.transport.redis' in sys.modules:
|
||
del sys.modules['kombu.transport.redis']
|
||
kombu_transport_redis = importlib.import_module('kombu.transport.redis')
|
||
print(f"✅ 重新导入kombu.transport.redis成功")
|
||
|
||
# 4. 继续原来的信号处理器代码
|
||
import logging
|
||
from django.db.models.signals import pre_save
|
||
from django.dispatch import receiver
|
||
from django.utils import timezone
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
@receiver(pre_save, sender='dingdan.Dingdan')
|
||
def order_status_8_handler_fixed(sender, instance, **kwargs):
|
||
"""
|
||
修复Redis导入问题后的信号处理器
|
||
"""
|
||
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 = 60 # 测试用60秒
|
||
|
||
# 🔥 再次验证redis模块
|
||
try:
|
||
import redis
|
||
logger.info(f"🔍 信号处理器内redis模块状态: {redis.Redis}")
|
||
if redis.Redis is None:
|
||
logger.error("❌ 信号处理器内redis.Redis为None!")
|
||
# 使用备用方案
|
||
instance.auto_task_id = f"redis_none_{timezone.now().timestamp()}"
|
||
return
|
||
except Exception as e:
|
||
logger.error(f"❌ 信号处理器内验证redis失败: {e}")
|
||
instance.auto_task_id = f"redis_check_failed_{timezone.now().timestamp()}"
|
||
return
|
||
|
||
# 提交Celery任务
|
||
try:
|
||
from dingdan.tasks import process_expired_order
|
||
|
||
logger.info(f"🔧 开始提交任务,订单: {instance.dingdan_id}")
|
||
|
||
# 使用更安全的方式提交任务
|
||
from celery import current_app
|
||
|
||
# 验证current_app的backend
|
||
if hasattr(current_app, 'backend'):
|
||
logger.info(f"✅ current_app.backend存在: {current_app.backend}")
|
||
else:
|
||
logger.error("❌ current_app没有backend属性")
|
||
|
||
task_result = process_expired_order.apply_async(
|
||
args=[instance.dingdan_id],
|
||
countdown=expire_seconds,
|
||
queue='order_tasks',
|
||
priority=9
|
||
)
|
||
|
||
logger.info(f"✅ Celery任务提交成功,订单: {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"❌ Celery任务提交失败: {type(e).__name__}: {e}")
|
||
|
||
# 特别检查Redis相关错误
|
||
if "'NoneType' object has no attribute 'Redis'" in str(e):
|
||
logger.error("💥 检测到Redis类为None错误")
|
||
logger.error("💥 尝试从当前环境重新导入redis模块")
|
||
|
||
# 尝试修复
|
||
try:
|
||
import importlib
|
||
import redis as redis_module
|
||
importlib.reload(redis_module)
|
||
logger.info(f"💥 重新加载redis模块后: {redis_module.Redis}")
|
||
except Exception as reload_e:
|
||
logger.error(f"💥 重新加载失败: {reload_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"订单状态变更处理失败,订单ID: {getattr(instance, 'dingdan_id', '未知')}, 错误: {e}")'''
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
|
||
'''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
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
# 在你的 dingdan/tasks.py 文件末尾添加以下代码
|
||
|
||
from django.db.models.signals import pre_save
|
||
from django.dispatch import receiver
|
||
from django.utils import timezone
|
||
import logging
|
||
|
||
from dingdan.models import Dingdan
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
@receiver(pre_save, sender=Dingdan)
|
||
def handle_order_status_change(sender, instance, **kwargs):
|
||
"""
|
||
订单状态变化监听器(纯净标记版)
|
||
严格遵循原始逻辑,仅更新数据库标记。
|
||
"""
|
||
# 1. 跳过新建的订单
|
||
if instance.pk is None:
|
||
return
|
||
|
||
# 2. 获取数据库中此订单的旧状态
|
||
try:
|
||
old_status = Dingdan.objects.filter(pk=instance.pk).values_list('zhuangtai', flat=True).first()
|
||
except Dingdan.DoesNotExist:
|
||
return # 旧订单不存在,不处理
|
||
except Exception as e:
|
||
logger.error(f"获取订单 {instance.pk} 旧状态失败: {e},跳过信号处理")
|
||
return # 出错时静默退出,绝不阻塞订单保存
|
||
|
||
old_status = old_status or 0 # 防None
|
||
|
||
# 3. 【核心逻辑】状态从 非8 变为 8
|
||
if old_status != 8 and instance.zhuangtai == 8:
|
||
instance.pending_dispatch = True # 标记为待处理
|
||
instance.status_8_time = timezone.now() # 记录变8时间
|
||
# 清空可能存在的旧任务信息(安全起见)
|
||
instance.auto_task_id = ''
|
||
instance.auto_expire_at = None
|
||
logger.info(f"📝 订单 {instance.dingdan_id} 状态变8,已标记为待调度。")
|
||
|
||
# 4. 【核心逻辑】状态从 8 变为 非8
|
||
elif old_status == 8 and instance.zhuangtai != 8:
|
||
# 清除所有调度标记和任务信息
|
||
instance.pending_dispatch = False
|
||
instance.status_8_time = None
|
||
instance.auto_task_id = ''
|
||
instance.auto_expire_at = None
|
||
logger.info(f"📝 订单 {instance.dingdan_id} 状态离开8,已清除所有调度标记。")'''
|
||
|
||
|
||
from django.db.models.signals import pre_save
|
||
from django.dispatch import receiver
|
||
from django.utils import timezone
|
||
from datetime import timedelta
|
||
from celery.result import AsyncResult
|
||
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):
|
||
"""
|
||
监听订单状态变化(稳定生产版)- 修复配置问题
|
||
"""
|
||
# 1. 跳过新建的订单(无旧状态可比)
|
||
if instance.pk is None:
|
||
return
|
||
|
||
try:
|
||
old = sender.objects.get(pk=instance.pk)
|
||
except sender.DoesNotExist:
|
||
return
|
||
# dingdan/signals.py - 修改状态变8的逻辑部分
|
||
|
||
|
||
# 2. 状态变成 8:提交任务
|
||
if old.zhuangtai != 8 and instance.zhuangtai == 8:
|
||
from django.conf import settings
|
||
from celery import current_app
|
||
|
||
# 🔥【核心修复】确保信号上下文中的Celery应用有正确配置
|
||
# 直接从你的Django settings中获取Redis连接字符串,并设置给current_app
|
||
if hasattr(settings, 'CELERY_BROKER_URL'):
|
||
current_app.conf.update(broker_url=settings.CELERY_BROKER_URL)
|
||
else:
|
||
# 如果settings里没有,就按你的配置拼接(和celery.py里一致)
|
||
redis_password = getattr(settings, 'REDIS_PASSWORD', 'Dujieduze5.3')
|
||
redis_host = getattr(settings, 'REDIS_HOST', '172.19.0.3')
|
||
redis_port = getattr(settings, 'REDIS_PORT', 6379)
|
||
redis_db_celery = getattr(settings, 'REDIS_DB_CELERY', 0)
|
||
broker_url = f'redis://:{redis_password}@{redis_host}:{redis_port}/{redis_db_celery}'
|
||
current_app.conf.update(broker_url=broker_url)
|
||
|
||
countdown_seconds = getattr(settings, 'ORDER_EXPIRE_SECONDS', 172800) # 默认48小时
|
||
|
||
try:
|
||
# 现在再提交任务
|
||
job = process_expired_order.apply_async(
|
||
args=[instance.dingdan_id],
|
||
countdown=countdown_seconds,
|
||
queue='order_tasks',
|
||
priority=9
|
||
)
|
||
# ✅ 核心:将任务ID和计划时间记录到当前实例
|
||
instance.auto_task_id = job.id
|
||
instance.auto_expire_at = timezone.now() + timedelta(seconds=countdown_seconds)
|
||
logger.info(f"✅ 订单 {instance.dingdan_id} 状态变8,已提交延时任务,ID: {job.id}")
|
||
except Exception as e:
|
||
logger.error(f"❌ 提交订单 {instance.dingdan_id} 的延时任务失败: {str(e)}")
|
||
return # 处理完毕,直接返回
|
||
|
||
# 3. 状态从 8 被改成别的 -> 尝试撤销任务 (这部分代码不变)
|
||
if old.zhuangtai == 8 and instance.zhuangtai != 8:
|
||
if old.auto_task_id:
|
||
try:
|
||
AsyncResult(old.auto_task_id).revoke()
|
||
logger.info(f"订单 {instance.dingdan_id} 状态离开8,已撤销原任务: {old.auto_task_id}")
|
||
except Exception as e:
|
||
logger.warning(f"撤销订单 {instance.dingdan_id} 的任务失败: {str(e)}")
|
||
instance.auto_task_id = ''
|
||
instance.auto_expire_at = None'''
|
||
|
||
|
||
|
||
'''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):
|
||
print(f"🚨【信号触发】订单ID: {getattr(instance, 'dingdan_id', 'N/A')}, 新状态: {instance.zhuangtai}")
|
||
"""
|
||
轻量级订单状态变化监听器
|
||
设计原则:
|
||
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)}")
|
||
# 不抛出异常,避免影响订单保存的主流程
|
||
# 信号处理失败不影响订单保存,有周期性检查作为备用方案''' |