from __future__ import annotations from datetime import datetime, timezone import redis from celery import shared_task from app.core.config import get_settings def _env_prefix(app_env: str) -> str: """ 根据环境生成前缀: - dev -> dev - prod -> pro """ return "dev" if str(app_env) == "dev" else "pro" @shared_task(name="tasks.ops.beat_heartbeat") def beat_heartbeat() -> dict[str, str]: """ Beat 心跳任务(用于健康检查)。 作用: - 由 Celery Beat 每分钟触发一次 - 写入 Redis 心跳 key,并设置 TTL - API 侧读取该 key,可判断 beat 是否在运行 """ settings = get_settings() prefix = _env_prefix(settings.app_env) key = f"{prefix}:beat:heartbeat" now = datetime.now(timezone.utc).isoformat() r = redis.Redis.from_url(settings.celery_broker_url, decode_responses=True) # TTL 设短一些:一旦 beat 挂了,很快就能从“过期/缺失”判断出来 r.set(key, now, ex=180) return {"status": "ok", "key": key, "at": now}