diff --git a/a_long_dianjing/celery.py b/a_long_dianjing/celery.py index d25e27d..b591226 100644 --- a/a_long_dianjing/celery.py +++ b/a_long_dianjing/celery.py @@ -1,157 +1,40 @@ """ -阿龙电竞 - Celery定时任务主配置 -生产环境就绪版本,支持精确延时任务和周期性任务 +阿龙电竞 - Celery 配置(仅新订单服务号广播) + +严禁 autodiscover / beat 加载结算、清零、排行榜等任务。 +生产只启动 broadcast 队列 worker,不要启动 celery beat。 """ - - - - import os from celery import Celery -from celery.schedules import crontab - -# 在 Django Shell 中执行 - - -# 设置Django默认设置模块 os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'a_long_dianjing.settings') -# 创建Celery应用实例 app = Celery('a_long_dianjing') - -# 从Django settings中加载Celery配置(CELERY_前缀) app.config_from_object('django.conf:settings', namespace='CELERY') -# 自动发现所有已注册app中的tasks.py文件 -#app.autodiscover_tasks() -app.autodiscover_tasks(['yonghu.tasks', 'yonghu.ranking_tasks', 'dingdan.tasks', 'peizhi.tasks']) -# 从ranking_tasks导入所有任务 +# 只注册新订单服务号广播任务(dingdan.tongzhi_tasks.dingdan_guangbo) +app.autodiscover_tasks(['dingdan.tongzhi_tasks']) +# 禁止 Beat 周期性任务(清零、排行榜、自动结算补偿等一律不调度) +app.conf.beat_schedule = {} - - -# 配置周期性任务(Celery Beat Schedule) -app.conf.beat_schedule = { - - # 在 app.conf.beat_schedule 中添加以下配置 - - # 在 app.conf.beat_schedule 中添加以下配置 - # 在 app.conf.beat_schedule 中修改以下任务 - - # 月榜数据转移(每月最后一天23:50执行) - 'yuebang_zhuanyi': { - 'task': 'yonghu.ranking_tasks.zhuanyi_yuebang', - 'schedule': crontab(hour=23, minute=50, day_of_month='28-31'), # ✅ 修改这里 - 'options': {'queue': 'periodic_tasks', 'priority': 7}, - }, - - - - # 日榜数据转移(每天23:55执行,在清零前) - 'ribang_zhuanyi': { - 'task': 'yonghu.ranking_tasks.zhuanyi_ribang', - 'schedule': crontab(hour=23, minute=55), # 每天23:55 - 'options': {'queue': 'periodic_tasks', 'priority': 7}, # 优先级高于清零任务 - }, - - - - # 清理旧历史数据(每月1号凌晨1点执行) - 'qingli_jiulishuju': { - 'task': 'yonghu.ranking_tasks.qingli_jiulishuju', - 'schedule': crontab(hour=1, minute=0, day_of_month=1), # 每月1号凌晨1点 - 'options': {'queue': 'periodic_tasks', 'priority': 4}, - }, - - - - - - - # 🔥 核心任务:订单调度器(每2分钟运行) - - - # 🔥 补偿检查任务(每5分钟运行)— 已禁用:见 settings.ORDER_AUTO_SETTLE_ENABLED - # 'check_order_expire_task': { - # 'task': 'dingdan.tasks.check_order_expire_task', - # 'schedule': crontab(minute='*/5'), - # 'options': {'queue': 'order_tasks', 'priority': 3}, - # }, - - - # 1. 每日凌晨0点执行 - 清零任务 - 'daily_reset_task': { - 'task': 'yonghu.tasks.daily_reset_task', - 'schedule': crontab(hour=0, minute=0), # 每天0点 - 'options': { - 'queue': 'periodic_tasks', - 'priority': 5 - }, - 'args': (), - 'kwargs': {} - }, - - # 2. 每月1日凌晨0点执行 - 月度清零任务 - 'monthly_reset_task': { - 'task': 'yonghu.tasks.monthly_reset_task', - 'schedule': crontab(hour=0, minute=0, day_of_month=1), # 每月1日0点 - 'options': { - 'queue': 'periodic_tasks', - 'priority': 5 - }, - 'args': (), - 'kwargs': {} - }, - - - - - - # 4. 每天凌晨0点5分清理收支记录 - 'daily_sz_reset_task': { - 'task': 'peizhi.tasks.daily_sz_reset_task', - 'schedule': crontab(hour=0, minute=5), # 每天0点5分 - 'options': { - 'queue': 'periodic_tasks', - 'priority': 4 - }, - 'args': (), - 'kwargs': {} - } -} - - - - -# 时区设置 app.conf.timezone = 'Asia/Shanghai' app.conf.enable_utc = True -# 队列路由配置 +# 仅广播任务路由到 broadcast 队列 app.conf.task_routes = { - 'dingdan.tasks.*': {'queue': 'order_tasks'}, - 'yonghu.tasks.*': {'queue': 'periodic_tasks'}, - 'peizhi.tasks.*': {'queue': 'periodic_tasks'}, + 'dingdan.tongzhi_tasks.dingdan_guangbo': {'queue': 'broadcast'}, } -# 任务序列化 app.conf.accept_content = ['json'] app.conf.task_serializer = 'json' app.conf.result_serializer = 'json' -# 任务超时设置 -app.conf.task_time_limit = 300 # 任务最大执行时间300秒 -app.conf.task_soft_time_limit = 240 # 软超时240秒 - -# Worker并发设置 -app.conf.worker_concurrency = 4 +app.conf.task_time_limit = 300 +app.conf.task_soft_time_limit = 240 +app.conf.worker_concurrency = 2 app.conf.worker_prefetch_multiplier = 1 - -# 任务确认设置 app.conf.task_acks_late = True app.conf.worker_disable_rate_limits = True - -# 结果过期时间 -app.conf.result_expires = 3600 # 任务结果保留1小时 \ No newline at end of file +app.conf.result_expires = 3600 diff --git a/a_long_dianjing/celery_app.py b/a_long_dianjing/celery_app.py index 10629c2..f238d0f 100644 --- a/a_long_dianjing/celery_app.py +++ b/a_long_dianjing/celery_app.py @@ -1,10 +1,4 @@ -# celery_app.py -import os -from celery import Celery +# 与 a_long_dianjing.celery 共用同一 Celery 实例,避免双 app 注册任务失败 +from a_long_dianjing.celery import app -os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'a_long_dianjing.settings') - -app = Celery('a_long_dianjing') -app.config_from_object('django.conf:settings', namespace='CELERY') -# 手动导入任务模块,避免与已有 tasks.py 冲突 -app.autodiscover_tasks(['dingdan.tongzhi_tasks']) \ No newline at end of file +__all__ = ('app',) diff --git a/a_long_dianjing/settings.py b/a_long_dianjing/settings.py index 82a48db..37eac54 100644 --- a/a_long_dianjing/settings.py +++ b/a_long_dianjing/settings.py @@ -428,6 +428,9 @@ CELERY_ENABLE_UTC = True # settings.py 末尾添加(在其他 Celery 配置之后) CELERY_IMPORTS = ('dingdan.tongzhi_tasks',) +# 生产硬性开关:仅允许新订单服务号广播 Celery 任务,其余一律禁止执行 +CELERY_ONLY_ORDER_BROADCAST = True + # 任务超时设置 CELERY_TASK_TIME_LIMIT = 300 # 任务最大执行时间300秒 CELERY_TASK_SOFT_TIME_LIMIT = 240 # 软超时240秒 @@ -463,7 +466,6 @@ CELERY_TASK_QUEUES = { CELERY_TASK_ROUTES = { 'dingdan.tongzhi_tasks.dingdan_guangbo': {'queue': 'broadcast'}, - 'dingdan.tasks.push_order_status_chat_task': {'queue': 'default'}, } # ==================== 定时任务时间配置 ==================== diff --git a/dingdan/signals.py b/dingdan/signals.py index c71e2a7..cfd00da 100644 --- a/dingdan/signals.py +++ b/dingdan/signals.py @@ -61,8 +61,22 @@ ensure_redis_loaded() @receiver(pre_save, sender='dingdan.Dingdan') def handle_order_status_8(sender, instance, **kwargs): """ - 订单状态变为8时,提交定时任务 + 订单状态变为8时,提交定时任务(已禁用:仅保留新订单广播) """ + if getattr(settings, 'CELERY_ONLY_ORDER_BROADCAST', True): + 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() + if old == 8 and instance.zhuangtai != 8: + instance.status_8_time = None + instance.pending_dispatch = False + instance.auto_task_id = '' + except Exception as e: + logger.error(f"信号处理失败: {e}") + return + # 跳过新增的订单 if instance.pk is None: return diff --git a/dingdan/tasks.py b/dingdan/tasks.py index 525dd72..21f589a 100644 --- a/dingdan/tasks.py +++ b/dingdan/tasks.py @@ -17,6 +17,7 @@ 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__) @@ -33,6 +34,11 @@ 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}' @@ -166,6 +172,10 @@ def process_expired_order(self, dingdan_id): @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) @@ -186,6 +196,10 @@ def retry_establish_order_chat_task(self, dingdan_id): @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) @@ -209,6 +223,11 @@ 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} diff --git a/dingdan/tongzhi_tasks.py b/dingdan/tongzhi_tasks.py index 490f955..a884fa1 100644 --- a/dingdan/tongzhi_tasks.py +++ b/dingdan/tongzhi_tasks.py @@ -3,7 +3,7 @@ import logging from django.conf import settings -from a_long_dianjing.celery_app import app +from a_long_dianjing.celery import app from utils.weixin_broadcast import WeixinBroadcastSender logger = logging.getLogger('weixin_broadcast') @@ -12,11 +12,13 @@ logger = logging.getLogger('weixin_broadcast') @app.task( bind=True, name='dingdan.tongzhi_tasks.dingdan_guangbo', - max_retries=2, + max_retries=3, default_retry_delay=30, queue='broadcast', soft_time_limit=240, time_limit=300, + acks_late=True, + reject_on_worker_lost=True, ) def dingdan_guangbo(self, order_info): """ diff --git a/dingdan/views.py b/dingdan/views.py index 6899a5a..4c5234e 100644 --- a/dingdan/views.py +++ b/dingdan/views.py @@ -4211,12 +4211,13 @@ class QiangdanView(APIView): chat_success = establish_order_chat_with_retry(dingdan_id, max_rounds=5) if not chat_success: - try: - from dingdan.tasks import retry_establish_order_chat_task - retry_establish_order_chat_task.apply_async(args=[dingdan_id], countdown=8) - retry_establish_order_chat_task.apply_async(args=[dingdan_id], countdown=30) - except Exception as chat_task_err: - logger.error(f"提交建聊补偿任务失败: {chat_task_err}") + if not getattr(settings, 'CELERY_ONLY_ORDER_BROADCAST', True): + try: + from dingdan.tasks import retry_establish_order_chat_task + retry_establish_order_chat_task.apply_async(args=[dingdan_id], countdown=8) + retry_establish_order_chat_task.apply_async(args=[dingdan_id], countdown=30) + except Exception as chat_task_err: + logger.error(f"提交建聊补偿任务失败: {chat_task_err}") try: update_dashou_daily_by_action( diff --git a/peizhi/tasks.py b/peizhi/tasks.py index f49f474..504abbc 100644 --- a/peizhi/tasks.py +++ b/peizhi/tasks.py @@ -14,6 +14,7 @@ import logging # 导入模型 from peizhi.models import Szjilu from utils.celery_utils import log_task_execution, rollback_on_failure +from utils.celery_guard import celery_non_broadcast_blocked logger = logging.getLogger(__name__) @@ -24,6 +25,11 @@ def daily_sz_reset_task(): 收支记录每日清零 - 每天凌晨0点5分执行 清零收支记录表的今日相关字段,固定使用ID=1的记录 """ + blocked = celery_non_broadcast_blocked() + if blocked: + logger.info('跳过 daily_sz_reset_task: %s', blocked) + return {"success": True, "message": blocked} + try: with transaction.atomic(): # 获取今天的日期 diff --git a/scripts/run_celery_broadcast_only.sh b/scripts/run_celery_broadcast_only.sh new file mode 100644 index 0000000..bbafbb7 --- /dev/null +++ b/scripts/run_celery_broadcast_only.sh @@ -0,0 +1,18 @@ +#!/bin/bash +# 仅新订单服务号广播 — 唯一允许的 Celery Worker 启动方式 +# 严禁同时启动 celery beat 或其它队列 worker(default/order_tasks/periodic_tasks) + +set -euo pipefail + +cd "$(dirname "$0")/.." + +export DJANGO_SETTINGS_MODULE=a_long_dianjing.settings + +exec celery -A a_long_dianjing worker \ + -l info \ + -Q broadcast \ + -c 2 \ + --hostname=broadcast@%h \ + --without-gossip \ + --without-mingle \ + -Ofair diff --git a/utils/celery_guard.py b/utils/celery_guard.py new file mode 100644 index 0000000..eed5546 --- /dev/null +++ b/utils/celery_guard.py @@ -0,0 +1,13 @@ +""" +Celery 任务安全闸:生产环境仅允许「新订单服务号广播」任务执行。 +""" +from django.conf import settings + +_DISABLED_MSG = 'CELERY_ONLY_ORDER_BROADCAST:非新订单广播任务已禁止执行' + + +def celery_non_broadcast_blocked(): + """若当前禁止非广播任务,返回说明字符串;否则返回 None。""" + if getattr(settings, 'CELERY_ONLY_ORDER_BROADCAST', True): + return _DISABLED_MSG + return None diff --git a/utils/chat_utils.py b/utils/chat_utils.py index 960cb63..e7bcf37 100644 --- a/utils/chat_utils.py +++ b/utils/chat_utils.py @@ -850,6 +850,8 @@ STATUS_NOTIFY_TEXT = { def schedule_order_status_chat_push(dingdan_id): """订单状态变更后,事务提交后异步推送群聊状态更新""" + if getattr(settings, 'CELERY_ONLY_ORDER_BROADCAST', True): + return if not dingdan_id: return diff --git a/yonghu/ranking_tasks.py b/yonghu/ranking_tasks.py index 4614f9f..d103efc 100644 --- a/yonghu/ranking_tasks.py +++ b/yonghu/ranking_tasks.py @@ -15,6 +15,7 @@ from yonghu.models import ( UserMain, UserDashou, UserShangjia, UserGuanshi, UserBoss, RankingRecord ) +from utils.celery_guard import celery_non_broadcast_blocked logger = logging.getLogger(__name__) @@ -44,6 +45,11 @@ def zhuanyi_ribang(): 转移日榜数据 - 每天23:55执行 在清零前将今日数据转移到排行榜表 """ + blocked = celery_non_broadcast_blocked() + if blocked: + logger.info('跳过 zhuanyi_ribang: %s', blocked) + return {"success": True, "message": blocked} + try: # 统计日期(今天) today = timezone.now().date() @@ -141,6 +147,11 @@ def zhuanyi_yuebang(): 转移月榜数据 - 每月最后一天23:50执行 在清零前将本月数据转移到排行榜表 """ + blocked = celery_non_broadcast_blocked() + if blocked: + logger.info('跳过 zhuanyi_yuebang: %s', blocked) + return {"success": True, "message": blocked} + try: # 🔥 获取当前时间 now = timezone.now() @@ -258,6 +269,11 @@ def qingli_jiulishuju(): 清理旧历史数据 - 每月1号凌晨1点执行 清理3个月前的排行榜数据,保留最近3个月的数据 """ + blocked = celery_non_broadcast_blocked() + if blocked: + logger.info('跳过 qingli_jiulishuju: %s', blocked) + return {"success": True, "message": blocked} + try: # 计算3个月前的日期 three_months_ago = timezone.now().date() - timedelta(days=90) diff --git a/yonghu/tasks.py b/yonghu/tasks.py index ce07131..3e61883 100644 --- a/yonghu/tasks.py +++ b/yonghu/tasks.py @@ -14,6 +14,7 @@ import logging # 导入模型 from yonghu.models import UserDashou, UserShangjia, UserGuanshi from utils.celery_utils import log_task_execution, rollback_on_failure +from utils.celery_guard import celery_non_broadcast_blocked logger = logging.getLogger(__name__) @@ -24,6 +25,11 @@ def daily_reset_task(): 每日清零任务 - 每天凌晨0点执行 清零打手、商家、管事的今日相关字段 """ + blocked = celery_non_broadcast_blocked() + if blocked: + logger.info('跳过 daily_reset_task: %s', blocked) + return {"success": True, "message": blocked} + try: with transaction.atomic(): start_time = timezone.now() @@ -103,6 +109,11 @@ def monthly_reset_task(): 每月清零任务 - 每月1日凌晨0点执行 清零打手、商家、管事的本月相关字段 """ + blocked = celery_non_broadcast_blocked() + if blocked: + logger.info('跳过 monthly_reset_task: %s', blocked) + return {"success": True, "message": blocked} + try: with transaction.atomic(): start_time = timezone.now()