增加定时服务的健康检测
This commit is contained in:
@@ -4,6 +4,7 @@ from datetime import datetime, timezone
|
||||
from typing import Any, Literal, Optional
|
||||
|
||||
import httpx
|
||||
import redis
|
||||
from fastapi import APIRouter, Depends, Header, HTTPException, Query
|
||||
from pydantic import BaseModel, Field
|
||||
from sqlalchemy import select
|
||||
@@ -13,8 +14,10 @@ from app.api.limits import rate_limit_push_by_ip
|
||||
from app.core.config import get_settings
|
||||
from app.db.models.push_preference import PushPreference
|
||||
from app.db.models.push_token import PushToken
|
||||
from app.db.models.push_send_log import PushSendLog
|
||||
from app.db.session import get_db
|
||||
from app.features.user_profile_scoring.types import UserProfileV1_2
|
||||
from app.worker import celery_app
|
||||
|
||||
|
||||
router = APIRouter(
|
||||
@@ -260,3 +263,84 @@ async def test_push(
|
||||
_ = accept_language
|
||||
return {"status": "ok", "expo": expo_res}
|
||||
|
||||
|
||||
def _env_prefix(app_env: str) -> str:
|
||||
"""
|
||||
根据环境生成前缀:
|
||||
- dev -> dev
|
||||
- prod -> pro
|
||||
"""
|
||||
|
||||
return "dev" if str(app_env) == "dev" else "pro"
|
||||
|
||||
|
||||
@router.get("/scheduler/health")
|
||||
async def scheduler_health(db: AsyncSession = Depends(get_db)) -> dict[str, Any]:
|
||||
"""
|
||||
推送“定时服务”健康检查(用于容器内验证)。
|
||||
|
||||
返回内容(尽量不暴露敏感信息):
|
||||
- Redis:是否可连通
|
||||
- Worker:是否至少有一个 worker 在线(inspect ping)
|
||||
- Beat:是否在跑(beat 心跳 key 是否在持续刷新)
|
||||
- DB:是否可查询到 push_send_log 的最新时间(辅助定位排程是否生成)
|
||||
"""
|
||||
|
||||
settings = get_settings()
|
||||
prefix = _env_prefix(settings.app_env)
|
||||
beat_key = f"{prefix}:beat:heartbeat"
|
||||
|
||||
out: dict[str, Any] = {
|
||||
"env": settings.app_env,
|
||||
"redis": {"ok": False},
|
||||
"worker": {"ok": False, "worker_count": 0},
|
||||
"beat": {"ok": False, "last_heartbeat_at": None, "age_seconds": None},
|
||||
"db": {"ok": False, "push_send_log_latest_created_at": None},
|
||||
"now_utc": datetime.now(timezone.utc).isoformat(),
|
||||
}
|
||||
|
||||
# 1) Redis 连通性 + 读取 beat 心跳
|
||||
try:
|
||||
r = redis.Redis.from_url(settings.celery_broker_url, decode_responses=True)
|
||||
r.ping()
|
||||
out["redis"]["ok"] = True
|
||||
|
||||
hb = r.get(beat_key)
|
||||
if hb:
|
||||
out["beat"]["last_heartbeat_at"] = hb
|
||||
try:
|
||||
# Python 3.11+ 支持解析 ISO8601(含 +00:00)
|
||||
hb_dt = datetime.fromisoformat(hb.replace("Z", "+00:00"))
|
||||
now = datetime.now(timezone.utc)
|
||||
age = int((now - hb_dt.astimezone(timezone.utc)).total_seconds())
|
||||
out["beat"]["age_seconds"] = age
|
||||
# 2 分钟内认为健康(beat 每分钟刷新一次)
|
||||
out["beat"]["ok"] = age <= 120
|
||||
except Exception:
|
||||
# 解析失败:至少说明 key 存在,但时间格式异常
|
||||
out["beat"]["ok"] = False
|
||||
except Exception as e:
|
||||
out["redis"]["error"] = f"{type(e).__name__}: {e}"
|
||||
|
||||
# 2) Worker 在线性(inspect ping)
|
||||
try:
|
||||
insp = celery_app.control.inspect(timeout=1.0)
|
||||
pings = insp.ping() or {}
|
||||
if isinstance(pings, dict):
|
||||
out["worker"]["worker_count"] = len(pings)
|
||||
out["worker"]["ok"] = len(pings) > 0
|
||||
except Exception as e:
|
||||
out["worker"]["error"] = f"{type(e).__name__}: {e}"
|
||||
|
||||
# 3) DB:查询 push_send_log 最新创建时间(用于判断排程是否有生成)
|
||||
try:
|
||||
q = select(PushSendLog.created_at).order_by(PushSendLog.created_at.desc()).limit(1)
|
||||
row = await db.execute(q)
|
||||
latest = row.scalar_one_or_none()
|
||||
out["db"]["ok"] = True
|
||||
out["db"]["push_send_log_latest_created_at"] = latest.isoformat() if latest else None
|
||||
except Exception as e:
|
||||
out["db"]["error"] = f"{type(e).__name__}: {e}"
|
||||
|
||||
return out
|
||||
|
||||
|
||||
@@ -10,4 +10,5 @@ Celery 任务集合。
|
||||
from app.tasks import ping as _ping # noqa: F401
|
||||
from app.tasks import reco as _reco # noqa: F401
|
||||
from app.tasks import push as _push # noqa: F401
|
||||
from app.tasks import ops as _ops # noqa: F401
|
||||
|
||||
|
||||
41
server/app/tasks/ops.py
Normal file
41
server/app/tasks/ops.py
Normal file
@@ -0,0 +1,41 @@
|
||||
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}
|
||||
|
||||
@@ -45,6 +45,12 @@ celery_app.conf.update(
|
||||
# - 这里按 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),
|
||||
|
||||
BIN
server/celerybeat-schedule
Normal file
BIN
server/celerybeat-schedule
Normal file
Binary file not shown.
Reference in New Issue
Block a user