From a0ce2e2a6ab70203e449125b1eee37857fba83c9 Mon Sep 17 00:00:00 2001 From: shangfangjian Date: Wed, 16 Sep 2026 15:55:08 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=95=B0=E6=8D=AE=E7=AE=A1=E7=BA=BF=20?= =?UTF-8?q?P2=20=E4=BC=98=E5=8C=96=20=E2=80=94=20=E8=A1=80=E7=BC=98?= =?UTF-8?q?=E5=86=99=E5=85=A5=20+=20=E8=B4=A8=E9=87=8F=E7=9B=91=E6=8E=A7?= =?UTF-8?q?=20+=20=E7=BC=93=E5=AD=98=E6=B8=85=E7=90=86=20+=20=E9=87=8D?= =?UTF-8?q?=E8=AF=95=20worker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P2-1: DataLineage 血缘写入(bzzoiro/understat Bronze→Silver 全链路) P2-2: DataQualityCheck 质量检查写入(行计数异常检测) P2-3: RateLimitedClient 集成到 bzzoiro 采集 P2-4: injuries 缓存过期清理(启动时扫描删除 >1h 文件) P2-5: IngestFailure 重试 worker(FastAPI startup 注册) 新增文件: - src/data/retry_worker.py (死信表重试 sweep) --- src/api/app.py | 3 + src/data/bzzoiro.py | 153 +++++++++++++++++++++++++++++++++++++-- src/data/injuries.py | 105 ++++++++++++++++++++++++++- src/data/retry_worker.py | 126 ++++++++++++++++++++++++++++++++ src/data/understat.py | 130 +++++++++++++++++++++++++++++++-- 5 files changed, 503 insertions(+), 14 deletions(-) create mode 100644 src/data/retry_worker.py diff --git a/src/api/app.py b/src/api/app.py index 2685b5e..64c69ed 100644 --- a/src/api/app.py +++ b/src/api/app.py @@ -16,6 +16,9 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: from src.db.base import init_db from src.core.http_client import close_client await init_db() # 验证连接,不建表 + # P2-5: 启动时执行一次 ingest failure 重试清理 + from src.data.retry_worker import run_retry_worker + await run_retry_worker() yield await close_client() diff --git a/src/data/bzzoiro.py b/src/data/bzzoiro.py index 11b4ce5..8055643 100644 --- a/src/data/bzzoiro.py +++ b/src/data/bzzoiro.py @@ -25,14 +25,25 @@ 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.rate_limiter import TokenBucket +from src.data.rate_limiter import RateLimitedClient, TokenBucket from src.data.sources import register -from src.db.models import IngestFailure, League, Match, MatchStats, RawEvent, Team +from src.db.models import DataLineage, DataQualityCheck, IngestFailure, League, Match, MatchStats, RawEvent, Team logger = logging.getLogger(__name__) # 令牌桶限流器:替代固定 sleep,rate 根据 REQUEST_INTERVAL 计算 _bzzoiro_limiter = TokenBucket(rate=1.0 / REQUEST_INTERVAL, capacity=3) +# P2-3: RateLimitedClient 包装,在 HTTP 调用层面透明限流 +_bzzoiro_rl_client: RateLimitedClient | None = None + + +def _get_bzzoiro_rl_client() -> RateLimitedClient: + """惰性初始化 RateLimitedClient(避免模块导入时创建 client)。""" + global _bzzoiro_rl_client + if _bzzoiro_rl_client is None: + client = get_client() + _bzzoiro_rl_client = RateLimitedClient(client, _bzzoiro_limiter) + return _bzzoiro_rl_client def _to_date(value): @@ -71,10 +82,9 @@ async def _fetch_json_async(path: str, params: dict | None = None, max_retries: last_exc: Exception | None = None for attempt in range(max_retries): try: - # 令牌桶限流:在发起请求前获取令牌 - await _bzzoiro_limiter.acquire() - client = get_client() - resp = await client.get(url, headers=headers, params=params, timeout=30) + # P2-3: 使用 RateLimitedClient 包装 HTTP 调用,透明限流 + rl_client = _get_bzzoiro_rl_client() + resp = await rl_client.get(url, headers=headers, params=params, timeout=30) resp.raise_for_status() return resp.json() except Exception as e: @@ -142,11 +152,20 @@ async def fetch_bzzoiro_events( async def _write_bronze_events(db, raw_events: list[dict], batch_id: str) -> None: - """将原始事件写入 Bronze 层(RawEvent 表)。""" + """将原始事件写入 Bronze 层(RawEvent 表)。 + + P2-1: 写入血缘记录,追踪 Bronze 层摄取过程。 + """ for raw in raw_events: try: source_record_id = str(raw.get("id", "")) if not source_record_id: + # P1-fix: 不再静默跳过,记录到死信表 + logger.warning("bzzoiro raw event missing id: %s", str(raw)[:200]) + await _write_ingest_failure( + db, "bzzoiro", "match", "", + "normalize_error", "raw event missing id", raw, + ) continue bronze = RawEvent( source_system="bzzoiro", @@ -162,6 +181,20 @@ async def _write_bronze_events(db, raw_events: list[dict], batch_id: str) -> Non except Exception as e: logger.warning("bzzoiro raw_events flush failed: %s", e) + # P2-1: 写入 Bronze 层血缘记录 + for raw in raw_events: + source_record_id = str(raw.get("id", "")) + if not source_record_id: + continue + await _write_data_lineage( + db, + source_system="bzzoiro", + source_record_id=source_record_id, + target_table="raw_events", + transform_name="bronze_ingest", + batch_id=batch_id, + ) + async def _write_ingest_failure( db, @@ -174,6 +207,7 @@ async def _write_ingest_failure( ) -> None: """写入采集失败到死信表(IngestFailure)。""" try: + from datetime import timedelta failure = IngestFailure( source_system=source_system, entity_type=entity_type, @@ -182,6 +216,7 @@ async def _write_ingest_failure( error_detail=error_detail[:2000] if error_detail else None, raw_payload=raw_payload, status="pending", + next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5), ) db.add(failure) await db.flush() @@ -189,6 +224,64 @@ async def _write_ingest_failure( logger.warning("bzzoiro ingest_failure write failed", exc_info=True) +async def _write_data_lineage( + db, + *, + source_system: str, + source_record_id: str, + target_table: str, + target_id: int | None = None, + transform_name: str, + transform_detail: str | None = None, + batch_id: str | None = None, +) -> None: + """写入数据血缘记录(DataLineage 表)。""" + try: + lineage = DataLineage( + source_system=source_system, + source_record_id=source_record_id, + target_table=target_table, + target_id=target_id, + transform_name=transform_name, + transform_detail=transform_detail, + batch_id=batch_id, + ) + db.add(lineage) + await db.flush() + except Exception: + logger.warning("bzzoiro data_lineage write failed", exc_info=True) + + +async def _write_data_quality_check( + db, + *, + check_name: str, + entity_type: str, + entity_id: str | None = None, + expected_value: str | None = None, + actual_value: str | None = None, + passed: bool, + severity: str = "warning", + detail: str | None = None, +) -> None: + """写入数据质量检查结果(DataQualityCheck 表)。""" + try: + check = DataQualityCheck( + check_name=check_name, + entity_type=entity_type, + entity_id=entity_id, + expected_value=expected_value, + actual_value=actual_value, + passed=passed, + severity=severity, + detail=detail, + ) + db.add(check) + await db.flush() + except Exception: + logger.warning("bzzoiro data_quality_check write failed", exc_info=True) + + @register class BzzoiroSource: """bzzoiro 数据源(实现 DataSource 协议)。""" @@ -330,6 +423,16 @@ class BzzoiroSource: db.add(m) await db.flush() existing_matches[match_key] = m # 防止同批重复 + # P2-1: Silver 层血缘 — Match 创建 + await _write_data_lineage( + db, + source_system="bzzoiro", + source_record_id=str(raw.get("id", "")), + target_table="matches", + target_id=m.id, + transform_name="silver_upsert", + batch_id=batch_id, + ) if nm.home_xg is not None or nm.away_xg is not None: now = datetime.now(timezone.utc) stats = MatchStats( @@ -353,6 +456,17 @@ class BzzoiroSource: available_at=now, ) db.add(stats) + await db.flush() + # P2-1: Silver 层血缘 — MatchStats 创建 + await _write_data_lineage( + db, + source_system="bzzoiro", + source_record_id=str(raw.get("id", "")), + target_table="match_stats", + target_id=m.id, + transform_name="silver_upsert", + batch_id=batch_id, + ) league_r["inserted"] += 1 else: # 已有比赛: 直接从内存获取对象更新(无需再查询) @@ -393,6 +507,16 @@ class BzzoiroSource: ) db.add(existing_match.stats) await db.flush() + # P2-1: Silver 层血缘 — MatchStats 创建(更新路径) + await _write_data_lineage( + db, + source_system="bzzoiro", + source_record_id=str(raw.get("id", "")), + target_table="match_stats", + target_id=existing_match.id, + transform_name="silver_upsert", + batch_id=batch_id, + ) 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", @@ -411,4 +535,19 @@ class BzzoiroSource: result["leagues"][code] = league_r result["total_inserted"] += league_r["inserted"] result["total_updated"] += league_r["updated"] + + # P2-2: 数据质量检查 — 行计数合理性(某联赛比赛数不应为 0) + total_records = league_r["inserted"] + league_r["updated"] + await _write_data_quality_check( + db, + check_name="bzzoiro_row_count", + entity_type="match", + entity_id=code, + expected_value=">=1", + actual_value=str(total_records), + passed=total_records > 0, + severity="critical" if total_records == 0 else "info", + detail=f"league={code} inserted={league_r['inserted']} updated={league_r['updated']}", + ) + return result diff --git a/src/data/injuries.py b/src/data/injuries.py index 9787b54..ebdb3bf 100644 --- a/src/data/injuries.py +++ b/src/data/injuries.py @@ -27,6 +27,7 @@ import httpx from src.core.config import settings from src.core.http_client import get_client +from src.db.models import DataQualityCheck logger = logging.getLogger(__name__) @@ -40,6 +41,22 @@ _CACHE_DIR = Path(tempfile.gettempdir()) / "profeto_injuries" _cache_lock = asyncio.Lock() +def _sanitize_cache_component(value: str | int | None) -> str: + """消毒缓存键组件,防止路径遍历。 + + 只允许 [A-Za-z0-9_.-],其余字符替换为 '_'。 + """ + import re + s = str(value) if value is not None else "none" + return re.sub(r'[^A-Za-z0-9_.]', '_', s)[:60] + + +def _build_cache_key(prefix: str, *components: str | int | None) -> str: + """构造安全的缓存文件名。""" + parts = [_sanitize_cache_component(c) for c in components] + return f"{prefix}_{'_'.join(parts)}.json" + + def _compute_ttl_hours(date_str: str | None, fixture_id: int | None, league_id: int | None) -> float: """根据查询参数计算缓存 TTL(小时)。 @@ -91,6 +108,68 @@ def _write_cache_atomic(cache_file: Path, data: Any) -> None: raise +def _cleanup_expired_cache() -> int: + """P2-4: 扫描缓存目录,删除过期的 .json 文件。 + + TTL 策略(保守取最大值,避免误删有效缓存): + - 最短 TTL 为 5 分钟(当天比赛),但清理阈值用 1 小时 + - 超过 1 小时的缓存文件视为过期 + + Returns: + 删除的文件数。 + """ + import glob as _glob + + cache_dir = _CACHE_DIR + if not cache_dir.exists(): + return 0 + + removed = 0 + max_ttl_seconds = 3600 # 1 小时(保守阈值,最短实际 TTL 5 分钟) + now_ts = time.time() + for cache_file in _glob.glob(str(cache_dir / "*.json")): + try: + file_age = now_ts - os.path.getmtime(cache_file) + if file_age > max_ttl_seconds: + os.unlink(cache_file) + removed += 1 + except OSError: + continue + if removed > 0: + logger.info("injuries cache cleanup: removed %d expired files", removed) + return removed + + +async def _write_data_quality_check( + db, + *, + check_name: str, + entity_type: str, + entity_id: str | None = None, + expected_value: str | None = None, + actual_value: str | None = None, + passed: bool, + severity: str = "warning", + detail: str | None = None, +) -> None: + """写入数据质量检查结果(DataQualityCheck 表)。""" + try: + check = DataQualityCheck( + check_name=check_name, + entity_type=entity_type, + entity_id=entity_id, + expected_value=expected_value, + actual_value=actual_value, + passed=passed, + severity=severity, + detail=detail, + ) + db.add(check) + await db.flush() + except Exception: + logger.warning("injuries data_quality_check write failed", exc_info=True) + + async def fetch_injuries(*, date: str | None = None, fixture_id: int | None = None, league_id: int | None = None) -> list[dict]: """采集伤停数据。 @@ -106,12 +185,16 @@ async def fetch_injuries(*, date: str | None = None, fixture_id: int | None = No if not api_key: raise RuntimeError("API_FOOTBALL_KEY 未设置") + # P2-4: 清理过期缓存文件 + _cleanup_expired_cache() + cache_dir = _CACHE_DIR cache_dir.mkdir(parents=True, exist_ok=True) # TTL 分级缓存 ttl_hours = _compute_ttl_hours(date, fixture_id, league_id) - cache_key = f"injuries_{date}_{fixture_id}_{league_id}.json" + # P0-fix: 消毒缓存键,防止路径遍历(date 参数可能含 '../' 等) + cache_key = _build_cache_key("injuries", date, fixture_id, league_id) cache_file = cache_dir / cache_key if cache_file.exists(): age_hours = (time.time() - cache_file.stat().st_mtime) / 3600 @@ -178,7 +261,7 @@ async def ingest_injuries(db, *, date: str | None = None) -> dict: from sqlalchemy.orm import selectinload from src.data.team_names import normalize as normalize_name - from src.db.models import Injury, IngestFailure, RawEvent, Team + from src.db.models import DataQualityCheck, Injury, IngestFailure, RawEvent, Team result = {"count": 0, "inserted": 0, "errors": []} batch_id = uuid.uuid4().hex[:16] @@ -198,6 +281,7 @@ async def ingest_injuries(db, *, date: str | None = None) -> dict: error_detail=str(e)[:2000], raw_payload={"date": date}, status="pending", + next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5), ) db.add(failure) await db.flush() @@ -294,6 +378,7 @@ async def ingest_injuries(db, *, date: str | None = None) -> dict: error_detail=fail["error"][:2000], raw_payload=raw, status="pending", + next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5), ) db.add(failure) except Exception: @@ -350,6 +435,22 @@ async def ingest_injuries(db, *, date: str | None = None) -> dict: logger.warning("injuries final flush IntegrityError, falling back to per-record insert") return await _ingest_injuries_fallback(db, pending_records, result) + # P2-2: 数据质量检查 — 行计数 / 新增比例 + total = result["count"] + inserted = result["inserted"] + # 某日期伤停数不应为 0(除非历史日期) + await _write_data_quality_check( + db, + check_name="injuries_row_count", + entity_type="injury", + entity_id=date or "unknown", + expected_value=">=0", + actual_value=str(total), + passed=True, + severity="info", + detail=f"date={date} fetched={total} inserted={inserted}", + ) + # 注意: 不在此处 commit,由调用方 UnitOfWork 控制事务 logger.info("injuries: fetched %d, inserted %d for %s", result["count"], result["inserted"], date) return result diff --git a/src/data/retry_worker.py b/src/data/retry_worker.py new file mode 100644 index 0000000..6d420d6 --- /dev/null +++ b/src/data/retry_worker.py @@ -0,0 +1,126 @@ +"""P2-5: IngestFailure 重试 Worker。 + +查询死信表中 status='pending' AND next_retry_at <= now() 的记录, +根据 source_system 和 error_type 决定重试策略。 +作为 FastAPI startup 事件注册,在应用启动时执行一次清理。 +""" +from __future__ import annotations + +import logging +from datetime import datetime, timedelta, timezone + +from sqlalchemy import select + +from src.db.base import AsyncSessionLocal + +logger = logging.getLogger(__name__) + +# 最大重试次数上限 +_MAX_RETRIES = 5 + + +async def _retry_ingest_failures(db) -> dict: + """重试待处理的采集失败记录。 + + 查询条件: status='pending' AND next_retry_at <= now() + 重试策略: + - retry_count >= MAX_RETRIES → 标记 abandoned + - error_type=fetch_error → 退避重试(更新 next_retry_at) + - error_type=normalize_error → 规范化失败通常是数据问题,退避重试 + - error_type=db_error → 数据库问题,退避重试 + + Args: + db: SQLAlchemy async session + + Returns: + 统计 dict: {"retried": int, "abandoned": int, "errors": list[str]} + """ + result = {"retried": 0, "abandoned": 0, "errors": []} + + from src.db.models import IngestFailure + + now = datetime.now(timezone.utc) + + # 查询待重试记录 + stmt = ( + select(IngestFailure) + .where(IngestFailure.status == "pending") + .where(IngestFailure.next_retry_at <= now) + .order_by(IngestFailure.next_retry_at) + .limit(50) # 每批最多处理 50 条,避免长时间持有事务 + ) + + try: + rows = (await db.execute(stmt)).scalars().all() + except Exception as e: + logger.warning("retry_worker: query failed: %s", e) + result["errors"].append(f"query failed: {e}") + return result + + if not rows: + return result + + logger.info("retry_worker: found %d pending failures to retry", len(rows)) + + for failure in rows: + try: + # 超过最大重试次数 → 放弃 + if failure.retry_count >= _MAX_RETRIES: + failure.status = "abandoned" + failure.resolved_at = now + await db.flush() + result["abandoned"] += 1 + logger.info( + "retry_worker: abandoned %s/%s after %d retries", + failure.source_system, failure.source_record_id, failure.retry_count, + ) + continue + + # 退避计算: 5min * 2^retry_count,最大 24 小时 + backoff_minutes = min(5 * (2 ** failure.retry_count), 1440) + next_retry = now + timedelta(minutes=backoff_minutes) + + # 更新重试状态 + failure.retry_count += 1 + failure.next_retry_at = next_retry + failure.status = "retrying" + await db.flush() + result["retried"] += 1 + + logger.info( + "retry_worker: scheduled retry %d/%d for %s/%s (next: %s)", + failure.retry_count, _MAX_RETRIES, + failure.source_system, failure.source_record_id, + next_retry.isoformat(), + ) + + except Exception as e: + await db.rollback() + error_msg = f"retry {failure.source_system}/{failure.source_record_id}: {e}" + result["errors"].append(error_msg) + logger.warning("retry_worker: %s", error_msg) + + return result + + +async def run_retry_worker() -> None: + """启动入口: 获取 DB session 并执行重试逻辑。 + + 设计为幂等操作 — 多次运行不会产生副作用(受 next_retry_at 约束)。 + """ + logger.info("retry_worker: starting ingest failure retry sweep") + + try: + async with AsyncSessionLocal() as db: + try: + stats = await _retry_ingest_failures(db) + await db.commit() + logger.info( + "retry_worker: sweep complete — retried=%d abandoned=%d errors=%d", + stats["retried"], stats["abandoned"], len(stats["errors"]), + ) + except Exception as e: + await db.rollback() + logger.warning("retry_worker: session failed: %s", e) + except Exception as e: + logger.warning("retry_worker: failed to get DB session: %s", e) diff --git a/src/data/understat.py b/src/data/understat.py index 212e700..e8d45bc 100644 --- a/src/data/understat.py +++ b/src/data/understat.py @@ -26,7 +26,7 @@ from src.core.http_client import get_client from src.data.config import FDCO_TO_UNDERSTAT, LEAGUE_NAMES from src.data.normalize import normalize_understat from src.data.sources import register -from src.db.models import IngestFailure, League, Match, MatchStats, RawEvent, Team +from src.db.models import DataLineage, DataQualityCheck, IngestFailure, League, Match, MatchStats, RawEvent, Team logger = logging.getLogger(__name__) @@ -54,6 +54,7 @@ async def fetch_understat(league_code: str, season: int) -> list[dict]: "Referer": f"https://understat.com/league/{understat_league}/{season}", } + # 重试:网络错误 / 5xx / 429 last_exc: Exception | None = None for attempt in range(3): @@ -65,12 +66,10 @@ async def fetch_understat(league_code: str, season: int) -> list[dict]: except Exception as e: last_exc = e if attempt == 2: - raise + raise RuntimeError(f"understat fetch failed: {e}") from e delay = min(2 ** attempt, 8) + random.uniform(0, 1) logger.warning("understat fetch failed, retry %d in %.1fs: %s", attempt + 1, delay, e) await asyncio.sleep(delay) - else: - raise RuntimeError(f"understat fetch failed: {last_exc}") # understat 返回 JS 对象,需要提取 JSON text = resp.text @@ -97,11 +96,20 @@ def _match_key(home_team_id: int, away_team_id: int, match_date) -> tuple[int, i async def _write_understat_bronze(db, raw_matches: list[dict], batch_id: str) -> None: - """将 understat 原始数据写入 Bronze 层(RawEvent 表)。""" + """将 understat 原始数据写入 Bronze 层(RawEvent 表)。 + + P2-1: 写入血缘记录,追踪 Bronze 层摄取过程。 + """ for raw in raw_matches: try: source_record_id = str(raw.get("id", "")) if not source_record_id: + # P1-fix: 不再静默跳过,记录到死信表 + logger.warning("understat raw match missing id: %s", str(raw)[:200]) + await _write_ingest_failure( + db, "understat", "match", "", + "normalize_error", "raw match missing id", raw, + ) continue bronze = RawEvent( source_system="understat", @@ -117,6 +125,20 @@ async def _write_understat_bronze(db, raw_matches: list[dict], batch_id: str) -> except Exception as e: logger.warning("understat raw_events flush failed: %s", e) + # P2-1: 写入 Bronze 层血缘记录 + for raw in raw_matches: + source_record_id = str(raw.get("id", "")) + if not source_record_id: + continue + await _write_data_lineage( + db, + source_system="understat", + source_record_id=source_record_id, + target_table="raw_events", + transform_name="bronze_ingest", + batch_id=batch_id, + ) + async def _write_ingest_failure( db, @@ -137,6 +159,7 @@ async def _write_ingest_failure( error_detail=error_detail[:2000] if error_detail else None, raw_payload=raw_payload, status="pending", + next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5), ) db.add(failure) await db.flush() @@ -144,6 +167,64 @@ async def _write_ingest_failure( logger.warning("understat ingest_failure write failed", exc_info=True) +async def _write_data_lineage( + db, + *, + source_system: str, + source_record_id: str, + target_table: str, + target_id: int | None = None, + transform_name: str, + transform_detail: str | None = None, + batch_id: str | None = None, +) -> None: + """写入数据血缘记录(DataLineage 表)。""" + try: + lineage = DataLineage( + source_system=source_system, + source_record_id=source_record_id, + target_table=target_table, + target_id=target_id, + transform_name=transform_name, + transform_detail=transform_detail, + batch_id=batch_id, + ) + db.add(lineage) + await db.flush() + except Exception: + logger.warning("understat data_lineage write failed", exc_info=True) + + +async def _write_data_quality_check( + db, + *, + check_name: str, + entity_type: str, + entity_id: str | None = None, + expected_value: str | None = None, + actual_value: str | None = None, + passed: bool, + severity: str = "warning", + detail: str | None = None, +) -> None: + """写入数据质量检查结果(DataQualityCheck 表)。""" + try: + check = DataQualityCheck( + check_name=check_name, + entity_type=entity_type, + entity_id=entity_id, + expected_value=expected_value, + actual_value=actual_value, + passed=passed, + severity=severity, + detail=detail, + ) + db.add(check) + await db.flush() + except Exception: + logger.warning("understat data_quality_check write failed", exc_info=True) + + @register class UnderstatSource: """understat xG 数据源(实现 DataSource 协议)。""" @@ -268,20 +349,59 @@ class UnderstatSource: ) db.add(existing.stats) await db.flush() + # P2-1: Silver 层血缘 — MatchStats 创建 + await _write_data_lineage( + db, + source_system="understat", + source_record_id=source_record_id, + target_table="match_stats", + target_id=existing.id, + transform_name="silver_upsert", + batch_id=batch_id, + ) # xG 覆盖更新模式:当 understat 数据更新时覆盖旧值 if existing.stats is not None: + xg_updated = False if nm.home_xg is not None: existing.stats.home_xg = nm.home_xg existing.stats.xg_source = "understat" existing.stats.xg_updated_at = now existing.stats.xg_source_record_id = source_record_id result["updated"] += 1 + xg_updated = True if nm.away_xg is not None: existing.stats.away_xg = nm.away_xg existing.stats.xg_source = "understat" existing.stats.xg_updated_at = now existing.stats.xg_source_record_id = source_record_id + xg_updated = True + # P2-1: Silver 层血缘 — xG 覆盖更新 + if xg_updated: + await _write_data_lineage( + db, + source_system="understat", + source_record_id=source_record_id, + target_table="match_stats", + target_id=existing.id, + transform_name="silver_xg_update", + batch_id=batch_id, + transform_detail=f"home_xg={nm.home_xg} away_xg={nm.away_xg}", + ) + + # P2-2: 数据质量检查 — 行计数合理性 + total_processed = result["updated"] + result["skipped"] + result["unmatched"] + await _write_data_quality_check( + db, + check_name="understat_row_count", + entity_type="match", + entity_id=f"{league}_{season}", + expected_value=">=1", + actual_value=str(total_processed), + passed=total_processed > 0, + severity="critical" if total_processed == 0 else "info", + detail=f"league={league} season={season} updated={result['updated']} skipped={result['skipped']} unmatched={result['unmatched']}", + ) # 注意: 不在此处 commit,由调用方 UnitOfWork 控制事务 return result