"""Bzzoiro 数据源:抓取 + 入库。 迁移自旧项目 app/data/sources/bzzoiro/,改成 async + 简化入库。 使用 Repository 模式进行数据访问,不直接控制事务。 """ from __future__ import annotations import asyncio import json as _json import logging import random from collections.abc import Iterable from datetime import datetime, timezone from sqlalchemy import select from src.core.config import settings from src.core.http_client import get_client from src.data.config import BZZOIRO_LEAGUE_IDS, LEAGUE_COUNTRIES, LEAGUE_NAMES, REQUEST_INTERVAL from src.data.normalize import normalize_bzzoiro from src.data.sources import register from src.db.models import League, Match, MatchStats, Team logger = logging.getLogger(__name__) def _to_date(value): """把 datetime / date / str 统一成 `date`。""" if value is None: return None if hasattr(value, "date") and callable(value.date): return value.date() return value def _match_key(home_team_id: int, away_team_id: int, match_date) -> tuple[int, int, str]: """比赛去重键:(主队, 客队, 天级日期 ISO 字符串)。 统一在这里构造,避免"预加载时用 str(date)、写入时用 isoformat()"这类 隐式格式依赖 —— 两者当前恰好相等,但一旦有人改动其一就会静默失配, 导致所有比赛被判为不存在而重复插入。 """ d = _to_date(match_date) return (home_team_id, away_team_id, d.isoformat() if d is not None else "") async def _fetch_json_async(path: str, params: dict | None = None, max_retries: int = 3) -> dict | list: """异步 HTTP(bzzoiro 使用 httpx,不再阻塞事件循环线程池)。""" base = settings.BZZOIRO_BASE.rstrip("/") url = f"{base}/{path.lstrip('/')}" key = settings.BZZOIRO_KEY if not key: raise RuntimeError("BZZOIRO_KEY 未设置") headers = { "Authorization": f"Token {key}", "Accept": "application/json", } last_exc: Exception | None = None for attempt in range(max_retries): try: client = get_client() resp = await client.get(url, headers=headers, params=params, timeout=30) resp.raise_for_status() return resp.json() except Exception as e: last_exc = e status = getattr(getattr(e, "response", None), "status_code", None) if status == 429: delay = min(2 ** attempt, 16) + random.uniform(0, 1) logger.warning("bzzoiro 429, retry %d in %.1fs", attempt + 1, delay) await asyncio.sleep(delay) continue if 500 <= (status or 0) < 600: delay = min(2 ** attempt, 16) + random.uniform(0, 1) logger.warning("bzzoiro %d, retry %d in %.1fs", status, attempt + 1, delay) await asyncio.sleep(delay) continue # 网络错误(连接失败/超时)也退避重试 if isinstance(e, (TimeoutError, ConnectionError, OSError)): delay = min(2 ** attempt, 16) + random.uniform(0, 1) logger.warning("bzzoiro network error, retry %d in %.1fs: %s", attempt + 1, delay, e) await asyncio.sleep(delay) continue raise raise RuntimeError(f"bzzoiro request failed after {max_retries} attempts: {last_exc}") async def fetch_bzzoiro_events( league_code: str, *, status: str = "finished", date_from: str | None = None, date_to: str | None = None, limit: int = 200, ) -> list[dict]: """抓取 bzzoiro 原始事件(纯异步,无需 run_in_executor)。""" league_id = BZZOIRO_LEAGUE_IDS.get(league_code) if league_id is None: raise ValueError(f"未知联赛代码: {league_code}") rows: list[dict] = [] offset = 0 while True: params: dict = { "league_id": league_id, "status": status, "limit": limit, "offset": offset, } if date_from: params["date_from"] = str(date_from)[:10] if date_to: params["date_to"] = str(date_to)[:10] payload = await _fetch_json_async("/events/", params) batch = payload.get("results") or [] if not batch: break rows.extend(batch) total = payload.get("total") offset += limit if total is not None and offset >= total: break if len(batch) < limit: break await asyncio.sleep(REQUEST_INTERVAL) return rows @register class BzzoiroSource: """bzzoiro 数据源(实现 DataSource 协议)。""" name = "bzzoiro" async def ingest( self, db, *, leagues: Iterable[str], date_from: str | None = None, date_to: str | None = None, status: str = "finished", ) -> dict: """采集 bzzoiro → 入库。返回统计。 注意: 本方法不控制事务(commit/rollback),由调用方通过 UnitOfWork 控制。 """ result: dict = {"leagues": {}, "total_inserted": 0, "total_updated": 0, "errors": []} for code in leagues: league_r: dict = {"inserted": 0, "updated": 0, "errors": []} try: raw_events = await fetch_bzzoiro_events(code, status=status, date_from=date_from, date_to=date_to) except Exception as e: logger.exception("bzzoiro fetch failed for %s", code) league_r["errors"].append(f"fetch failed: {e}") result["leagues"][code] = league_r continue # 获取或创建联赛 stmt = select(League).where(League.code == code) league = (await db.execute(stmt)).scalar_one_or_none() if league is None: league = League(code=code, name=LEAGUE_NAMES.get(code, code), country=LEAGUE_COUNTRIES.get(code)) db.add(league) await db.flush() # === 批量优化: 预加载球队和已有比赛到内存 === team_name_to_id: dict[str, int] = {} existing_matches: dict[tuple[int, int, str], Match] = {} # 完整对象,避免重复查询 # (NormalizedMatch, 原始 event) 成对保存:后续写 source_record_id 时 # 必须用配对的那条 event,不能依赖外层循环变量残留值。 normalized_matches: list[tuple] = [] if raw_events: # 一次遍历: 收集球队名 + 规范化 all_team_names = set() for raw in raw_events: nm = normalize_bzzoiro(raw, code) if nm is not None: try: nm.validate() except Exception as e: # P1-3: 统一使用 warning,不追加到 errors(仅运行时错误入 errors) logger.warning("normalize skip: %s", e) continue normalized_matches.append((nm, raw)) all_team_names.add(nm.home_team) all_team_names.add(nm.away_team) if all_team_names: stmt = select(Team).where(Team.name.in_(all_team_names)) teams = (await db.execute(stmt)).scalars().all() team_name_to_id = {t.name: t.id for t in teams} # P1-2: 按需加载,只加载 raw_events 涉及日期范围的比赛(加 30 天缓冲) # 避免加载联赛全部历史比赛到内存(多赛季采集时内存溢出) if normalized_matches: from datetime import timedelta dates = [nm.date for nm in normalized_matches if nm.date is not None] if dates: min_dt = min(dates) - timedelta(days=30) max_dt = max(dates) + timedelta(days=30) stmt = ( select(Match) .where(Match.league_id == league.id) .where(Match.match_date >= min_dt) .where(Match.match_date <= max_dt) ) existing_matches = { _match_key(m.home_team_id, m.away_team_id, m.match_date_date): m for m in (await db.execute(stmt)).scalars() } # else: existing_matches 保持空 dict(全量新比赛) for nm, raw in normalized_matches: # 球队: 内存查找 + 按需创建 home_team_id = team_name_to_id.get(nm.home_team) if home_team_id is None: home = Team(name=nm.home_team) db.add(home) await db.flush() home_team_id = home.id team_name_to_id[nm.home_team] = home_team_id away_team_id = team_name_to_id.get(nm.away_team) if away_team_id is None: away = Team(name=nm.away_team) db.add(away) await db.flush() away_team_id = away.id team_name_to_id[nm.away_team] = away_team_id # 查找已有比赛: 内存查找 match_key = _match_key(home_team_id, away_team_id, nm.date) existing_match = existing_matches.get(match_key) if existing_match is None: m = Match( league_id=league.id, season=nm.season_label or None, home_team_id=home_team_id, away_team_id=away_team_id, match_date=nm.date, match_date_date=_to_date(nm.date), match_status=nm.match_status, home_goals=nm.home_goals, away_goals=nm.away_goals, home_ht_goals=nm.home_ht_goals, away_ht_goals=nm.away_ht_goals, match_stage=nm.match_stage, ) db.add(m) await db.flush() existing_matches[match_key] = m # 防止同批重复 if nm.home_xg is not None or nm.away_xg is not None: now = datetime.now(timezone.utc) stats = MatchStats( match_id=m.id, home_xg=nm.home_xg, away_xg=nm.away_xg, home_shots=nm.home_shots, away_shots=nm.away_shots, home_shots_on_target=nm.home_shots_on_target, away_shots_on_target=nm.away_shots_on_target, home_corners=nm.home_corners, away_corners=nm.away_corners, home_possession=nm.home_possession, home_yellow_cards=nm.home_yellow_cards, away_yellow_cards=nm.away_yellow_cards, home_red_cards=nm.home_red_cards, away_red_cards=nm.away_red_cards, source="bzzoiro", source_event_id=str(raw.get("id", "")), retrieved_at=now, available_at=now, ) db.add(stats) league_r["inserted"] += 1 else: # 已有比赛: 直接从内存获取对象更新(无需再查询) changed = False if existing_match.match_status != nm.match_status and nm.match_status == "finished": existing_match.match_status = nm.match_status changed = True if existing_match.home_goals is None and nm.home_goals is not None: existing_match.home_goals = nm.home_goals existing_match.away_goals = nm.away_goals existing_match.home_ht_goals = nm.home_ht_goals existing_match.away_ht_goals = nm.away_ht_goals changed = True if existing_match.match_stage is None and nm.match_stage: existing_match.match_stage = nm.match_stage changed = True if existing_match.stats is None and (nm.home_xg is not None or nm.away_xg is not None): now = datetime.now(timezone.utc) existing_match.stats = MatchStats( match_id=existing_match.id, source="bzzoiro", source_event_id=str(raw.get("id", "")), retrieved_at=now, available_at=now, ) db.add(existing_match.stats) await db.flush() if existing_match.stats is not None: for fld in ("home_xg", "away_xg", "home_shots", "away_shots", "home_shots_on_target", "away_shots_on_target", "home_corners", "away_corners", "home_possession", "home_yellow_cards", "away_yellow_cards", "home_red_cards", "away_red_cards"): if getattr(existing_match.stats, fld, None) is None: v = getattr(nm, fld, None) if v is not None: setattr(existing_match.stats, fld, v) changed = True if changed: league_r["updated"] += 1 # 注意: 不在此处 commit,由调用方 UnitOfWork 控制事务 result["leagues"][code] = league_r result["total_inserted"] += league_r["inserted"] result["total_updated"] += league_r["updated"] return result