"""Bzzoiro 数据源:抓取 + 入库(单一数据源)。 三条管线: 1. events — 比赛日程/比分(/events/),含 source_event_id 血缘 2. standings— 联赛积分榜快照(/leagues/{id}/standings/) 3. stats — 已完赛比赛详细统计回填(/events/{id}/stats/) 使用 Repository 模式进行数据访问,不直接控制事务(由调用方 UnitOfWork 控制)。 """ from __future__ import annotations import asyncio import logging import random from collections.abc import Iterable from datetime import datetime, timedelta, timezone from sqlalchemy import select import httpx from src.core.runtime_config import get_runtime_value 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.key_ring import get_key_ring from src.data.normalize import normalize_bzzoiro from src.data.team_names_zh import zh_name from src.data.sources import register from src.db.models import League, Match, MatchStats, Standing, 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 _to_int_or_none(value) -> int | None: """宽松转 int(用于上游 ID 解析,失败返回 None 不抛错)。""" if value is None: return None try: return int(str(value).strip()) except (TypeError, ValueError): return None 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,不再阻塞事件循环线程池)。 多 key 轮换:遇到 429 自动切换到下一个 key;全部 key 冷却时等待最早恢复。 """ base = (await get_runtime_value("BZZOIRO_BASE")).rstrip("/") raw_keys = await get_runtime_value("BZZOIRO_KEY") ring = get_key_ring(base, raw_keys) url = f"{base}/{path.lstrip('/')}" key = ring.get() if not key: raise RuntimeError("BZZOIRO_KEY 未设置") last_exc: Exception | None = None for attempt in range(max_retries): headers = { "Authorization": f"Token {key}", "Accept": "application/json", } try: client = get_client() # 整请求兜底: httpx 无 total 超时,用 wait_for 防「滴水式」限速挂死 resp = await asyncio.wait_for( client.get( url, headers=headers, params=params, timeout=httpx.Timeout(connect=10.0, read=30.0, write=10.0, pool=10.0), ), timeout=60.0, ) 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: # 限流:标记当前 key 冷却,切换到下一个 new_key = ring.report_rate_limited(key) if new_key and new_key != key: logger.info("bzzoiro 429 → 切换 key: %s → %s,立即重试", _km(key), _km(new_key)) key = new_key continue # 立即重试,不等待 # 单 key 或全部冷却:等待最早恢复的 key wait = ring.wait_if_all_blocked() if wait > 0: logger.warning("bzzoiro 全部 key 冷却,等待 %.1fs 后重试", wait) await asyncio.sleep(min(wait, 30.0)) else: delay = min(2 ** attempt, 16) + random.uniform(0, 1) logger.warning("bzzoiro 429, retry %d in %.1fs", attempt + 1, delay) await asyncio.sleep(delay) key = ring.get() or key 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}") def _km(key: str) -> str: """key 脱敏缩写(用于日志)。""" if len(key) <= 8: return key[:2] + "***" return key[:4] + "..." + key[-4:] 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 # normalized_matches 存的是 (nm, raw) 元组,遍历需解包 dates = [nm.date for nm, _raw 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, name_zh=zh_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, name_zh=zh_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, source_event_id=_to_int_or_none(raw.get("id")), ) db.add(m) await db.flush() existing_matches[match_key] = m # 防止同批重复 # 统计字段不在 /events/ 载荷中(单独由 stats 管线回填), # 此处不再创建 MatchStats。 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.source_event_id is None: eid = _to_int_or_none(raw.get("id")) if eid is not None: existing_match.source_event_id = eid 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 # ============================================================ # 积分榜管线:/leagues/{id}/standings/ → standings 表 # ============================================================ async def fetch_bzzoiro_standings(league_code: str, season: str | None = None) -> dict: """抓取联赛积分榜(纯抓取,不入库)。season 为 None 时取当前赛季。""" league_id = BZZOIRO_LEAGUE_IDS.get(league_code) if league_id is None: raise ValueError(f"未知联赛代码: {league_code}") params: dict = {} if season: params["season"] = season return await _fetch_json_async(f"/leagues/{league_id}/standings/", params) def _season_label_from_dates(start_date, end_date) -> str: """从赛季起止日期推导赛季标签(与 derive_season_label 语义一致)。""" try: if isinstance(start_date, str): start = datetime.fromisoformat(start_date[:10]) else: start = start_date y = start.year return f"{y}-{y + 1}" if start.month >= 8 else f"{y - 1}-{y}" except (TypeError, ValueError): return "?" async def ingest_bzzoiro_standings(db, *, leagues: Iterable[str], season: str | None = None) -> dict: """采集积分榜 → upsert standings 表。 season 为 None 时采集当前赛季(bzzoiro 默认返回 is_current 赛季)。 球队名与 events 管线使用同一 normalize 规则,保证 Team 匹配。 """ from src.data.team_names import normalize as normalize_name result: dict = {"leagues": {}, "total_upserted": 0, "errors": []} for code in leagues: league_r: dict = {"upserted": 0, "teams_created": 0, "rows": 0} try: payload = await fetch_bzzoiro_standings(code, season=season) except Exception as e: logger.exception("bzzoiro standings fetch failed for %s", code) result["leagues"][code] = {"error": str(e)} result["errors"].append(f"{code}: {e}") continue rows = payload.get("standings") or [] if not rows: result["leagues"][code] = {"error": "无积分榜数据(赛季未开始或未提供)"} result["errors"].append(f"{code}: 无积分榜数据") 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() # 赛季标签:优先用返回的 season 对象推导 season_obj = payload.get("season") or {} season_label = _season_label_from_dates( season_obj.get("start_date"), season_obj.get("end_date") ) if season_label == "?": season_label = season or "" # 批量预载球队 names = {normalize_name(str(r.get("team_name", ""))) for r in rows} names.discard("") team_map: dict[str, Team] = {} if names: stmt = select(Team).where(Team.name.in_(names)) for t in (await db.execute(stmt)).scalars(): team_map[t.name] = t # 预加载该 league+season 已有的 standings(避免同批内重复 INSERT 导致 UniqueViolation) existing_standings: set[int] = set() stmt = select(Standing.team_id).where( Standing.league_id == league.id, Standing.season == season_label, ) for (tid,) in (await db.execute(stmt)).all(): existing_standings.add(tid) now = datetime.now(timezone.utc) seen_teams: set[int] = set() # 同批内去重:同一 team 只处理一次 for r in rows: team_name = normalize_name(str(r.get("team_name", ""))) if not team_name: continue team = team_map.get(team_name) if team is None: team = Team(name=team_name, name_zh=zh_name(team_name)) db.add(team) await db.flush() team_map[team_name] = team league_r["teams_created"] += 1 # 同批内同一 team 仅处理第一次 if team.id in seen_teams: continue seen_teams.add(team.id) zone = r.get("zone") or {} values = dict( position=_to_int_or_none(r.get("position")) or 0, played=_to_int_or_none(r.get("played")) or 0, won=_to_int_or_none(r.get("won")) or 0, drawn=_to_int_or_none(r.get("drawn")) or 0, lost=_to_int_or_none(r.get("lost")) or 0, goals_for=_to_int_or_none(r.get("gf")) or 0, goals_against=_to_int_or_none(r.get("ga")) or 0, goal_diff=_to_int_or_none(r.get("gd")) or 0, points=_to_int_or_none(r.get("pts")) or 0, xg_for=_to_float_or_none(r.get("xgf")), xg_against=_to_float_or_none(r.get("xga")), form=r.get("form") or None, zone=zone.get("label") or zone.get("key") or None, updated_at=now, retrieved_at=now, ) if team.id in existing_standings: # 已有 → 仍需 UPDATE:回查对象(少量,可接受) stmt = select(Standing).where( Standing.league_id == league.id, Standing.season == season_label, Standing.team_id == team.id, ) standing = (await db.execute(stmt)).scalar_one_or_none() if standing: for k, v in values.items(): setattr(standing, k, v) else: db.add(Standing( league_id=league.id, season=season_label, team_id=team.id, **values, )) existing_standings.add(team.id) # 防止同批内重复 league_r["upserted"] += 1 league_r["rows"] = len(rows) result["leagues"][code] = league_r result["total_upserted"] += league_r["upserted"] logger.info( "bzzoiro standings 采集完成: %s 赛季 %s, upsert %d/%d", code, season_label, league_r["upserted"], league_r["rows"], ) return result # ============================================================ # 统计回填管线:/events/{id}/stats/ → match_stats 表 # ============================================================ # bzzoiro stats 字段 → MatchStats 字段映射(stats.home / stats.away 下) _STATS_FIELD_MAP = { "xg": ("home_xg", "away_xg"), # 回退 expected_goals "ball_possession": ("home_possession", None), # 只取主队值,客队=100-home "total_shots": ("home_shots", "away_shots"), "shots_on_target": ("home_shots_on_target", "away_shots_on_target"), "corner_kicks": ("home_corners", "away_corners"), "yellow_cards": ("home_yellow_cards", "away_yellow_cards"), "red_cards": ("home_red_cards", "away_red_cards"), "big_chances": ("home_big_chances", "away_big_chances"), "fouls": ("home_fouls", "away_fouls"), } def _pick(d: dict, *keys): """按优先级取第一个非空字段值。""" for k in keys: v = d.get(k) if v is not None: return v return None def _stats_from_payload(payload: dict) -> dict: """把 /events/{id}/stats/ 响应映射成 MatchStats 字段 dict。 响应结构: {"event_id": ..., "stats": {"home": {...}, "away": {...}}} """ stats = (payload or {}).get("stats") or {} home = stats.get("home") or {} away = stats.get("away") or {} out: dict = {} xg_h = _pick(home, "xg", "expected_goals") xg_a = _pick(away, "xg", "expected_goals") if xg_h is not None: out["home_xg"] = _to_float_or_none(xg_h) if xg_a is not None: out["away_xg"] = _to_float_or_none(xg_a) poss = home.get("ball_possession") if poss is not None: p = _to_float_or_none(poss) if p is not None: out["home_possession"] = p out["away_possession"] = round(100 - p, 1) if 0 <= p <= 100 else None for src, (h_fld, a_fld) in _STATS_FIELD_MAP.items(): if src in ("xg", "ball_possession"): continue # 已处理 hv = home.get(src) av = away.get(src) if hv is not None and h_fld: out[h_fld] = _to_int_or_none(hv) if av is not None and a_fld: out[a_fld] = _to_int_or_none(av) return out def _to_float_or_none(value) -> float | None: if value is None: return None try: return float(str(value).strip()) except (TypeError, ValueError): return None async def ingest_bzzoiro_event_stats( db, *, leagues: Iterable[str], limit: int = 100, only_missing: bool = True, ) -> dict: """回填已完赛比赛的详细统计(逐场调 /events/{id}/stats/)。 筛选条件: match_status=finished 且 source_event_id 非空。 only_missing=True 时跳过已有统计的比赛(增量);False 则全量刷新。 limit 控制单次最多处理的比赛数(上游限速 1.2s/请求,大批量需分次触发)。 """ result: dict = {"fetched": 0, "created": 0, "updated": 0, "skipped": 0, "errors": []} league_ids = [BZZOIRO_LEAGUE_IDS[c] for c in leagues if c in BZZOIRO_LEAGUE_IDS] if not league_ids: result["errors"].append("无有效联赛代码") return result stmt = ( select(Match) .options(select(Match.stats)) .where(Match.match_status == "finished") .where(Match.source_event_id.is_not(None)) .where(Match.league_id.in_(league_ids)) .order_by(Match.match_date.desc()) .limit(limit * 3 if only_missing else limit) ) matches = (await db.execute(stmt)).scalars().all() now = datetime.now(timezone.utc) processed = 0 for m in matches: if processed >= limit: break if only_missing and m.stats is not None and m.stats.home_shots is not None: result["skipped"] += 1 continue processed += 1 try: payload = await _fetch_json_async(f"/events/{m.source_event_id}/stats/") except Exception as e: logger.warning("stats fetch failed match=%s event=%s: %s", m.id, m.source_event_id, e) result["errors"].append(f"match {m.id}: {e}") await asyncio.sleep(REQUEST_INTERVAL) continue result["fetched"] += 1 fields = _stats_from_payload(payload) if not fields: result["skipped"] += 1 await asyncio.sleep(REQUEST_INTERVAL) continue if m.stats is None: # available_at 语义:完赛统计最早在开球+2h 可用(回测防泄漏) available_at = m.match_date + timedelta(hours=2) if m.match_date else now m.stats = MatchStats( match_id=m.id, source="bzzoiro", source_record_id=str(m.source_event_id), retrieved_at=now, available_at=available_at, ) db.add(m.stats) result["created"] += 1 else: result["updated"] += 1 if m.stats.source is None: m.stats.source = "bzzoiro" m.stats.source_record_id = str(m.source_event_id) if m.stats.retrieved_at is None: m.stats.retrieved_at = now if m.stats.available_at is None and m.match_date: m.stats.available_at = m.match_date + timedelta(hours=2) for fld, v in fields.items(): # away_possession 为计算字段,模型无此列,跳过 if hasattr(m.stats, fld): setattr(m.stats, fld, v) await asyncio.sleep(REQUEST_INTERVAL) logger.info( "bzzoiro stats 回填完成: 抓取 %d, 新建 %d, 更新 %d, 跳过 %d, 错误 %d", result["fetched"], result["created"], result["updated"], result["skipped"], len(result["errors"]), ) return result