65 lines
2.1 KiB
Python
65 lines
2.1 KiB
Python
from __future__ import annotations
|
||
|
||
from celery import Celery
|
||
from celery.schedules import crontab
|
||
|
||
from app.core.config import get_settings
|
||
|
||
|
||
def _env_prefix(app_env: str) -> str:
|
||
"""
|
||
根据环境生成前缀:
|
||
- dev -> dev
|
||
- prod -> pro
|
||
|
||
说明:Redis ACL 限制使用 `dev:*` / `pro:*`。
|
||
"""
|
||
|
||
return "dev" if app_env == "dev" else "pro"
|
||
|
||
|
||
settings = get_settings()
|
||
prefix = _env_prefix(settings.app_env)
|
||
|
||
# 关键:使用 Redis transport 的全局 key 前缀,确保所有 broker key 都在 ACL 允许范围内
|
||
broker_transport_options = {"global_keyprefix": f"{prefix}:"}
|
||
|
||
celery_app = Celery(
|
||
"mindfulness",
|
||
broker=settings.celery_broker_url,
|
||
backend=settings.celery_result_backend,
|
||
broker_transport_options=broker_transport_options,
|
||
)
|
||
|
||
# 默认不存结果(降低 Redis 占用)。如需结果存储,可在业务中显式开启并设置 TTL。
|
||
celery_app.conf.update(
|
||
task_ignore_result=True if not settings.celery_result_backend else False,
|
||
task_default_queue=f"{prefix}:celery",
|
||
task_default_exchange=f"{prefix}:celery",
|
||
task_default_routing_key=f"{prefix}:celery",
|
||
)
|
||
|
||
# 定时任务(Celery Beat)
|
||
# 说明:
|
||
# - 每天运行一次“生成推送计划”,为所有开启每日提醒的用户生成当天/明天的随机抖动时间点,并投递 ETA 发送任务
|
||
# - 这里按 UTC 00:10 触发一次;具体时间可按运维习惯调整
|
||
celery_app.conf.timezone = "UTC"
|
||
celery_app.conf.beat_schedule = {
|
||
# Beat 心跳:用于 API 健康检查判断 beat 是否在跑
|
||
"ops-beat-heartbeat": {
|
||
"task": "tasks.ops.beat_heartbeat",
|
||
"schedule": crontab(minute="*/1"),
|
||
"options": {"queue": f"{prefix}:celery"},
|
||
},
|
||
"push-generate-daily-schedule": {
|
||
"task": "tasks.push.generate_daily_schedule",
|
||
"schedule": crontab(minute=10, hour=0),
|
||
"kwargs": {"max_users": 5000},
|
||
"options": {"queue": f"{prefix}:celery"},
|
||
}
|
||
}
|
||
|
||
# 自动发现任务(约定:导入 app.tasks 触发其内部对子模块的显式导入)
|
||
celery_app.autodiscover_tasks(["app"])
|
||
|