"""dramatiq 启动 setup 模块。 worker 启动命令引用本模块以完成 broker 初始化、actor 注册与 crontab 周期任务接管: dramatiq config.dramatiq --processes 1 --threads 4 周期任务(dramatiq-crontab,装饰器用法 @cron("分 时 日 月 周")): - 每日 08:00 全租户业务预警扫描(含 AI 应收风险) - 每日 02:00 电商渠道自动拉单(有凭证店铺) - 每日 03:00 套餐试用/订阅到期检查(降级 free + 通知) - 每日 04:00 审计日志保留清理(默认保留 180 天) """ # 0. 初始化 Django(worker 进程必须先加载 Django 再注册模型 actor) import os import django os.environ.setdefault("DJANGO_SETTINGS_MODULE", "config.settings.prod") django.setup() # 1. 初始化 broker(必须在导入 actor 之前) from .broker import broker as _broker # noqa: F401 # 2. crontab 调度 + 业务 actors import dramatiq from dramatiq_crontab import cron from apps.notify.tasks import run_alert_checks_for_tenant from apps.channel.tasks import sync_orders_for_account from apps.billing.tasks import expire_trials_and_notify def _enqueue_all_alert_checks(): """拆分到单租户 actor,避免单条消息超时。""" from apps.core.models import Tenant for t in Tenant.objects.filter(is_active=True): run_alert_checks_for_tenant.send(t.id) def _enqueue_all_channel_syncs(): from apps.channel.models import ChannelAccount for acc in ChannelAccount.objects.filter(is_active=True).exclude(app_key=""): sync_orders_for_account.send(acc.id) def _run_billing_expiry(): """套餐到期检查(租户数少,直接同步跑,无需拆消息)。""" from apps.billing import tasks as billing_tasks return billing_tasks.expire_trials() @dramatiq.actor(max_retries=1) def scheduler_tick_alerts(): _enqueue_all_alert_checks() @dramatiq.actor(max_retries=1) def scheduler_tick_channel_syncs(): _enqueue_all_channel_syncs() @dramatiq.actor(max_retries=1) def scheduler_tick_billing_expiry(): return _run_billing_expiry() @cron("0 8 * * *") @dramatiq.actor(max_retries=1) def cron_daily_alert_checks(): _enqueue_all_alert_checks() @cron("0 2 * * *") @dramatiq.actor(max_retries=1) def cron_daily_channel_syncs(): _enqueue_all_channel_syncs() @cron("0 3 * * *") @dramatiq.actor(max_retries=1) def cron_daily_billing_expiry(): """每日 03:00 扫描试用/订阅到期 → 降级 free 并发通知。""" return _run_billing_expiry() @cron("0 4 * * *") @dramatiq.actor(max_retries=1) def cron_daily_audit_purge(): """每日 04:00 清理过期审计日志(P1-3,默认保留 180 天)。""" from apps.core import tasks as core_tasks return core_tasks.purge_expired_audit_logs()