"""采集路由(bzzoiro 单一数据源)。 任务类型: events — 比赛日程/比分(/events/) standings — 联赛积分榜(/leagues/{id}/standings/) stats — 已完赛比赛详细统计回填(/events/{id}/stats/) all — 依次执行以上三项 """ from __future__ import annotations import asyncio import logging from fastapi import APIRouter, Depends, HTTPException from src.api.deps import require_admin from src.api.schemas import IngestBzzoiroRequest from src.data.config import BZZOIRO_LEAGUE_IDS from src.data.bzzoiro import ingest_bzzoiro_event_stats, ingest_bzzoiro_standings from src.data.sources import get_source from src.db.unit_of_work import get_uow logger = logging.getLogger(__name__) router = APIRouter(prefix="/api/v1", tags=["ingest"]) # 后台采集任务注册表:持强引用防止被 GC _background_tasks: set[asyncio.Task] = set() VALID_TASKS = {"events", "standings", "stats", "all"} def _spawn(coro) -> None: """启动后台采集任务;异常已在任务内记录到系统日志。""" task = asyncio.create_task(coro) _background_tasks.add(task) task.add_done_callback(_background_tasks.discard) @router.post("/ingest/bzzoiro", dependencies=[Depends(require_admin)]) async def ingest_bzzoiro_route(req: IngestBzzoiroRequest): """触发 bzzoiro 采集(events / standings / stats / all)。""" if req.task not in VALID_TASKS: raise HTTPException(status_code=422, detail=f"未知任务类型: {req.task}(可选: {', '.join(sorted(VALID_TASKS))})") leagues = req.leagues or list(BZZOIRO_LEAGUE_IDS.keys()) task_label = {"events": "比赛数据", "standings": "积分榜", "stats": "统计回填", "all": "全量(比赛+积分榜+统计)"}[req.task] _spawn(_run_bzzoiro(req.task, leagues, req)) return { "ok": True, "message": f"采集任务已启动(后台执行,任务: {task_label}),请在「系统日志」查看进度与结果", } async def _run_bzzoiro(task: str, leagues: list[str], req: IngestBzzoiroRequest) -> None: """后台执行 bzzoiro 采集:上游限速时单次可能耗时数分钟,必须脱离请求生命周期。""" try: if task in ("events", "all"): statuses = [req.status] if req.status else ["finished", "scheduled"] source = get_source("bzzoiro") merged: dict = {"leagues": {}, "total_inserted": 0, "total_updated": 0, "errors": []} # D2 修复: 按联赛分批提交,避免超长事务 for code in leagues: for st in statuses: async with get_uow() as session: r = await source.ingest( session, leagues=[code], date_from=req.date_from, date_to=req.date_to, status=st, ) merged["total_inserted"] += r.get("total_inserted", 0) merged["total_updated"] += r.get("total_updated", 0) merged["errors"].extend(r.get("errors", [])) acc = merged["leagues"].setdefault(code, {"inserted": 0, "updated": 0, "errors": []}) acc["inserted"] += r.get("inserted", 0) acc["updated"] += r.get("updated", 0) acc["errors"].extend(r.get("errors", [])) logger.info( "bzzoiro 比赛采集完成: 新增 %d, 更新 %d, 联赛 %d 个, 状态 %s", merged["total_inserted"], merged["total_updated"], len(merged["leagues"]), statuses, ) if merged["errors"]: logger.warning("bzzoiro 比赛采集错误 %d 条: %s", len(merged["errors"]), merged["errors"][:3]) if task in ("standings", "all"): async with get_uow() as session: r = await ingest_bzzoiro_standings(session, leagues=leagues, season=req.season) if r["errors"]: logger.warning("bzzoiro 积分榜采集部分失败: %s", r["errors"][:3]) else: logger.info("bzzoiro 积分榜采集完成: upsert %d 条", r["total_upserted"]) if task in ("stats", "all"): async with get_uow() as session: r = await ingest_bzzoiro_event_stats( session, leagues=leagues, limit=req.limit, only_missing=True ) if r["errors"]: logger.warning("bzzoiro 统计回填错误 %d 条: %s", len(r["errors"]), r["errors"][:3]) except Exception: logger.exception("bzzoiro 采集任务失败(task=%s)", task)