单文件 852 行按职责拆分,保持 BzzoiroSource 与 get_source("bzzoiro") 行为不变:
- bzzoiro_common HTTP 抓取(多 key 轮换) + 字段转换原语
- bzzoiro_events fetch_bzzoiro_events + BzzoiroSource.ingest + Bronze 补写
- bzzoiro_standings standings 管线
- bzzoiro_stats stats 回填
- pipeline_write RawEvent/IngestFailure/DataLineage 写入助手
子模块运行期经聚合门面 src.data.bzzoiro 解析可替换协作者,
单文件时代的 bz.* monkeypatch 语义完全保留。
路由 import 已指向新模块(ingest.py / schedules.py)。
186 lines
6.5 KiB
Python
186 lines
6.5 KiB
Python
"""定时任务管理路由。"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from datetime import datetime, timezone
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException
|
|
from sqlalchemy import select, delete
|
|
|
|
from src.api.deps import require_admin
|
|
from src.api.schemas import ScheduleIn, ScheduleUpdate, ScheduleOut
|
|
from src.core.scheduler import scheduler
|
|
from src.data.bzzoiro_standings import ingest_bzzoiro_standings
|
|
from src.data.bzzoiro_stats import ingest_bzzoiro_event_stats
|
|
from src.data.sources import get_source
|
|
from src.data.config import BZZOIRO_LEAGUE_IDS
|
|
from src.db.base import AsyncSession, get_db_read
|
|
from src.db.models import Schedule
|
|
from src.db.unit_of_work import get_uow
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
router = APIRouter(prefix="/api/v1/admin", tags=["schedule"], dependencies=[Depends(require_admin)])
|
|
|
|
|
|
async def _run_scheduled_task(schedule_id: str) -> None:
|
|
"""执行定时任务的回调函数。"""
|
|
async with get_uow() as session:
|
|
stmt = select(Schedule).where(Schedule.id == schedule_id)
|
|
sched = (await session.execute(stmt)).scalar_one_or_none()
|
|
if sched is None or not sched.enabled:
|
|
return
|
|
|
|
leagues = sched.leagues.split(",") if sched.leagues else list(BZZOIRO_LEAGUE_IDS.keys())
|
|
task = sched.task
|
|
|
|
try:
|
|
if task in ("events", "all"):
|
|
statuses = ["finished", "scheduled"]
|
|
source = get_source("bzzoiro")
|
|
for st in statuses:
|
|
await source.ingest(session, leagues=leagues, status=st)
|
|
|
|
if task in ("standings", "all"):
|
|
await ingest_bzzoiro_standings(session, leagues=leagues)
|
|
|
|
if task in ("stats", "all"):
|
|
await ingest_bzzoiro_event_stats(session, leagues=leagues, limit=500, only_missing=True)
|
|
|
|
sched.last_status = "success"
|
|
except Exception:
|
|
logger.exception("定时任务执行失败: %s", schedule_id)
|
|
sched.last_status = "failed"
|
|
finally:
|
|
sched.last_run_at = datetime.now(timezone.utc)
|
|
|
|
|
|
def _sync_scheduler() -> None:
|
|
"""同步数据库中的调度配置到调度器。"""
|
|
# 这是一个简化版本:实际应该在 lifespan 中异步同步
|
|
pass
|
|
|
|
|
|
@router.get("/schedules")
|
|
async def list_schedules(db: AsyncSession = Depends(get_db_read)):
|
|
"""列出所有定时任务。"""
|
|
rows = (await db.execute(select(Schedule).order_by(Schedule.created_at))).scalars().all()
|
|
return [
|
|
ScheduleOut(
|
|
id=s.id,
|
|
task=s.task,
|
|
cron=s.cron,
|
|
leagues=s.leagues,
|
|
enabled=s.enabled,
|
|
last_run_at=s.last_run_at.isoformat() if s.last_run_at else None,
|
|
last_status=s.last_status,
|
|
)
|
|
for s in rows
|
|
]
|
|
|
|
|
|
@router.post("/schedules")
|
|
async def create_schedule(req: ScheduleIn, db: AsyncSession = Depends(get_db_read)):
|
|
"""创建定时任务。"""
|
|
sched = Schedule(
|
|
id=req.id,
|
|
task=req.task,
|
|
cron=req.cron,
|
|
leagues=",".join(req.leagues) if req.leagues else None,
|
|
enabled=req.enabled,
|
|
)
|
|
db.add(sched)
|
|
await db.commit()
|
|
|
|
# 注册到调度器(注意: 第 2 参是 cron 表达式,不是 task 类型名)
|
|
scheduler.register(req.id, req.cron, lambda: _run_scheduled_task(req.id), enabled=req.enabled)
|
|
|
|
return {"ok": True, "id": req.id}
|
|
|
|
|
|
@router.put("/schedules/{schedule_id}")
|
|
async def update_schedule(schedule_id: str, req: ScheduleUpdate, db: AsyncSession = Depends(get_db_read)):
|
|
"""更新定时任务(部分更新)。"""
|
|
stmt = select(Schedule).where(Schedule.id == schedule_id)
|
|
sched = (await db.execute(stmt)).scalar_one_or_none()
|
|
if sched is None:
|
|
raise HTTPException(404, "定时任务不存在")
|
|
|
|
if req.task is not None:
|
|
sched.task = req.task
|
|
if req.cron is not None:
|
|
sched.cron = req.cron
|
|
if req.leagues is not None:
|
|
sched.leagues = ",".join(req.leagues) if req.leagues else None
|
|
if req.enabled is not None:
|
|
sched.enabled = req.enabled
|
|
await db.commit()
|
|
|
|
# 更新调度器(第 2 参传 cron 表达式;使用最终值)
|
|
scheduler.register(schedule_id, sched.cron, lambda: _run_scheduled_task(schedule_id), enabled=sched.enabled)
|
|
|
|
return {"ok": True}
|
|
|
|
|
|
@router.delete("/schedules/{schedule_id}")
|
|
async def delete_schedule(schedule_id: str, db: AsyncSession = Depends(get_db_read)):
|
|
"""删除定时任务。"""
|
|
await db.execute(delete(Schedule).where(Schedule.id == schedule_id))
|
|
await db.commit()
|
|
scheduler.remove(schedule_id)
|
|
return {"ok": True}
|
|
|
|
|
|
@router.post("/schedules/{schedule_id}/run")
|
|
async def run_schedule_now(schedule_id: str):
|
|
"""手动触发定时任务。"""
|
|
import asyncio
|
|
asyncio.create_task(_run_scheduled_task(schedule_id))
|
|
return {"ok": True, "message": "任务已启动"}
|
|
|
|
|
|
# ── 采集失败重试 ────────────────────────────────────────────────
|
|
|
|
|
|
@router.get("/ingest-failures")
|
|
async def list_ingest_failures(db: AsyncSession = Depends(get_db_read)):
|
|
"""列出采集失败记录。"""
|
|
from src.db.models import IngestFailure
|
|
rows = (
|
|
await db.execute(
|
|
select(IngestFailure).order_by(IngestFailure.created_at.desc()).limit(50)
|
|
)
|
|
).scalars().all()
|
|
return [
|
|
{
|
|
"id": f.id,
|
|
"source": f.source_system,
|
|
"entity_type": f.entity_type,
|
|
"source_record_id": f.source_record_id,
|
|
"error_type": f.error_type,
|
|
"error_detail": f.error_detail,
|
|
"retry_count": f.retry_count,
|
|
"status": f.status,
|
|
"next_retry_at": f.next_retry_at.isoformat() if f.next_retry_at else None,
|
|
"created_at": f.created_at.isoformat() if f.created_at else None,
|
|
}
|
|
for f in rows
|
|
]
|
|
|
|
|
|
@router.post("/ingest-failures/{failure_id}/retry")
|
|
async def retry_ingest_failure(failure_id: int, db: AsyncSession = Depends(get_db_read)):
|
|
"""重试一次采集失败。"""
|
|
from src.db.models import IngestFailure
|
|
stmt = select(IngestFailure).where(IngestFailure.id == failure_id)
|
|
failure = (await db.execute(stmt)).scalar_one_or_none()
|
|
if failure is None:
|
|
raise HTTPException(404, "失败记录不存在")
|
|
|
|
failure.status = "retrying"
|
|
failure.retry_count += 1
|
|
await db.commit()
|
|
|
|
# 触发重试(简化版:仅标记状态,实际重试逻辑由调度器处理)
|
|
return {"ok": True, "message": f"已标记重试 (第 {failure.retry_count} 次)"}
|