单文件 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)。
216 lines
8.7 KiB
Python
216 lines
8.7 KiB
Python
"""bzzoiro standings 管线:联赛积分榜快照(/leagues/{id}/standings/)→ standings 表。
|
|
|
|
从 bzzoiro.py 拆出。同一联赛同一赛季只保留最新快照(按 (league, season, team)
|
|
upsert);球队名与 events 管线使用同一 normalize 规则,保证 Team 匹配。
|
|
|
|
可替换协作者(抓取函数 / Bronze 写入助手)在运行期经聚合门面
|
|
src.data.bzzoiro 解析 —— 与拆分前的单文件 monkeypatch 语义一致。
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from collections.abc import Iterable
|
|
from datetime import datetime, timezone
|
|
|
|
from sqlalchemy import select
|
|
|
|
from src.data.bzzoiro_common import _to_float_or_none, _to_int_or_none
|
|
from src.data.config import BZZOIRO_LEAGUE_IDS, LEAGUE_COUNTRIES, LEAGUE_NAMES
|
|
from src.data.team_names_zh import zh_name
|
|
from src.db.models import Standing, Team
|
|
from src.db.repositories import LeagueRepository, TeamRepository
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
async def fetch_bzzoiro_standings(league_code: str, season: str | None = None) -> dict:
|
|
"""抓取联赛积分榜(纯抓取,不入库)。season 为 None 时取当前赛季。"""
|
|
from src.data import bzzoiro as bz
|
|
|
|
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 bz._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
|
|
if start is None:
|
|
return "?"
|
|
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 import bzzoiro as bz
|
|
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, "errors": []}
|
|
try:
|
|
payload = await bz.fetch_bzzoiro_standings(code, season=season)
|
|
except Exception as e:
|
|
logger.exception("bzzoiro standings fetch failed for %s", code)
|
|
league_r["errors"].append(str(e))
|
|
await bz._safe_write_ingest_failure(
|
|
db,
|
|
entity_type="standings",
|
|
source_record_id=None,
|
|
error=e,
|
|
raw_payload={"league": code, "season": season},
|
|
)
|
|
result["leagues"][code] = league_r
|
|
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
|
|
|
|
# 联赛(get-or-create,D4: 经 LeagueRepository)
|
|
league = await LeagueRepository(db).get_or_create(
|
|
code, LEAGUE_NAMES.get(code, code), LEAGUE_COUNTRIES.get(code)
|
|
)
|
|
team_r = TeamRepository(db)
|
|
|
|
# 赛季标签:优先用返回的 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 ""
|
|
|
|
# 批量预载球队(与 events 管线使用同一 normalize 规则,保证 Team 匹配)
|
|
names = {normalize_name(str(r.get("team_name", ""))) for r in rows}
|
|
names.discard("")
|
|
team_map: dict[str, Team] = await team_r.get_all_by_names(list(names))
|
|
|
|
now = datetime.now(timezone.utc)
|
|
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 = await team_r.get_or_create(team_name, name_zh=zh_name(team_name))
|
|
team_map[team_name] = team
|
|
league_r["teams_created"] += 1
|
|
|
|
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,
|
|
)
|
|
|
|
# 同一联赛同一赛季只保留最新快照:按 (league, season, team) upsert
|
|
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 is None:
|
|
standing = Standing(
|
|
league_id=league.id, season=season_label, team_id=team.id, **values
|
|
)
|
|
db.add(standing)
|
|
else:
|
|
for k, v in values.items():
|
|
setattr(standing, k, v)
|
|
league_r["upserted"] += 1
|
|
|
|
league_r["rows"] = len(rows)
|
|
|
|
# D1(对称 events/stats 管线): 联赛成功 upsert → 补写 Bronze 层。
|
|
# 幂等键 standings:{league}:{season}:积分榜是联赛级快照,一次成功
|
|
# 采集写一条 RawEvent(整份原始载荷)+ 一条血缘。season 用实际入库的
|
|
# 标签(由载荷推导,与 Standing.season 同口径),不依赖调用方传参,
|
|
# 保证不同调用方(season=None 或显式传参)对同一赛季命中同一条 RawEvent。
|
|
if league_r["upserted"] > 0:
|
|
bronze_batch_id = f"bzzoiro-standings-{code}-{now:%Y%m%d%H%M%S}"
|
|
await _write_standings_bronze(
|
|
db,
|
|
source_record_id=f"standings:{code}:{season_label}",
|
|
raw_payload=payload,
|
|
league_id=league.id,
|
|
league_code=code,
|
|
season_label=season_label,
|
|
rows_upserted=league_r["upserted"],
|
|
batch_id=bronze_batch_id,
|
|
)
|
|
|
|
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
|
|
|
|
|
|
async def _write_standings_bronze(
|
|
db,
|
|
*,
|
|
source_record_id: str,
|
|
raw_payload: dict,
|
|
league_id: int | None,
|
|
league_code: str,
|
|
season_label: str,
|
|
rows_upserted: int,
|
|
batch_id: str,
|
|
) -> None:
|
|
"""standings 成功 upsert 一个联赛后的 Bronze 层补写:RawEvent(幂等) + DataLineage。
|
|
|
|
与 _write_events_bronze 同级约束:幂等性由 _write_raw_event 的
|
|
source_record_id 查重保证(积分榜是联赛级快照,同联赛同赛季重复采集
|
|
命中同一条 RawEvent);best-effort:基础设施写入失败只记 warning,
|
|
绝不拖垮采集主流程。
|
|
"""
|
|
from src.data import bzzoiro as bz
|
|
|
|
try:
|
|
await bz._write_raw_event(db, "bzzoiro", source_record_id, raw_payload, batch_id)
|
|
await bz._write_lineage(
|
|
db, "bzzoiro", source_record_id,
|
|
"standings", league_id, "standings_ingest",
|
|
{"league": league_code, "season": season_label, "rows_upserted": rows_upserted},
|
|
batch_id,
|
|
)
|
|
except Exception:
|
|
logger.warning(
|
|
"standings Bronze 写入失败(record=%s, league=%s),不影响采集主流程",
|
|
source_record_id, league_code, exc_info=True,
|
|
)
|