feat: 数据管线 P2 优化 — 血缘写入 + 质量监控 + 缓存清理 + 重试 worker

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)
This commit is contained in:
shangfangjian
2026-09-16 15:55:08 +08:00
parent 5de8ffb09d
commit a0ce2e2a6a
5 changed files with 503 additions and 14 deletions
+3
View File
@@ -16,6 +16,9 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]:
from src.db.base import init_db from src.db.base import init_db
from src.core.http_client import close_client from src.core.http_client import close_client
await init_db() # 验证连接,不建表 await init_db() # 验证连接,不建表
# P2-5: 启动时执行一次 ingest failure 重试清理
from src.data.retry_worker import run_retry_worker
await run_retry_worker()
yield yield
await close_client() await close_client()
+146 -7
View File
@@ -25,14 +25,25 @@ from src.core.config import settings
from src.core.http_client import get_client 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.config import BZZOIRO_LEAGUE_IDS, LEAGUE_COUNTRIES, LEAGUE_NAMES, REQUEST_INTERVAL
from src.data.normalize import normalize_bzzoiro 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.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__) logger = logging.getLogger(__name__)
# 令牌桶限流器:替代固定 sleep,rate 根据 REQUEST_INTERVAL 计算 # 令牌桶限流器:替代固定 sleep,rate 根据 REQUEST_INTERVAL 计算
_bzzoiro_limiter = TokenBucket(rate=1.0 / REQUEST_INTERVAL, capacity=3) _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): 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 last_exc: Exception | None = None
for attempt in range(max_retries): for attempt in range(max_retries):
try: try:
# 令牌桶限流:在发起请求前获取令牌 # P2-3: 使用 RateLimitedClient 包装 HTTP 调用,透明限流
await _bzzoiro_limiter.acquire() rl_client = _get_bzzoiro_rl_client()
client = get_client() resp = await rl_client.get(url, headers=headers, params=params, timeout=30)
resp = await client.get(url, headers=headers, params=params, timeout=30)
resp.raise_for_status() resp.raise_for_status()
return resp.json() return resp.json()
except Exception as e: 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: 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: for raw in raw_events:
try: try:
source_record_id = str(raw.get("id", "")) source_record_id = str(raw.get("id", ""))
if not source_record_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 continue
bronze = RawEvent( bronze = RawEvent(
source_system="bzzoiro", 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: except Exception as e:
logger.warning("bzzoiro raw_events flush failed: %s", 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( async def _write_ingest_failure(
db, db,
@@ -174,6 +207,7 @@ async def _write_ingest_failure(
) -> None: ) -> None:
"""写入采集失败到死信表(IngestFailure)。""" """写入采集失败到死信表(IngestFailure)。"""
try: try:
from datetime import timedelta
failure = IngestFailure( failure = IngestFailure(
source_system=source_system, source_system=source_system,
entity_type=entity_type, entity_type=entity_type,
@@ -182,6 +216,7 @@ async def _write_ingest_failure(
error_detail=error_detail[:2000] if error_detail else None, error_detail=error_detail[:2000] if error_detail else None,
raw_payload=raw_payload, raw_payload=raw_payload,
status="pending", status="pending",
next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5),
) )
db.add(failure) db.add(failure)
await db.flush() await db.flush()
@@ -189,6 +224,64 @@ async def _write_ingest_failure(
logger.warning("bzzoiro ingest_failure write failed", exc_info=True) 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 @register
class BzzoiroSource: class BzzoiroSource:
"""bzzoiro 数据源(实现 DataSource 协议)。""" """bzzoiro 数据源(实现 DataSource 协议)。"""
@@ -330,6 +423,16 @@ class BzzoiroSource:
db.add(m) db.add(m)
await db.flush() await db.flush()
existing_matches[match_key] = m # 防止同批重复 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: if nm.home_xg is not None or nm.away_xg is not None:
now = datetime.now(timezone.utc) now = datetime.now(timezone.utc)
stats = MatchStats( stats = MatchStats(
@@ -353,6 +456,17 @@ class BzzoiroSource:
available_at=now, available_at=now,
) )
db.add(stats) 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 league_r["inserted"] += 1
else: else:
# 已有比赛: 直接从内存获取对象更新(无需再查询) # 已有比赛: 直接从内存获取对象更新(无需再查询)
@@ -393,6 +507,16 @@ class BzzoiroSource:
) )
db.add(existing_match.stats) db.add(existing_match.stats)
await db.flush() 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: if existing_match.stats is not None:
for fld in ("home_xg", "away_xg", "home_shots", "away_shots", for fld in ("home_xg", "away_xg", "home_shots", "away_shots",
"home_shots_on_target", "away_shots_on_target", "home_shots_on_target", "away_shots_on_target",
@@ -411,4 +535,19 @@ class BzzoiroSource:
result["leagues"][code] = league_r result["leagues"][code] = league_r
result["total_inserted"] += league_r["inserted"] result["total_inserted"] += league_r["inserted"]
result["total_updated"] += league_r["updated"] 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 return result
+103 -2
View File
@@ -27,6 +27,7 @@ import httpx
from src.core.config import settings from src.core.config import settings
from src.core.http_client import get_client from src.core.http_client import get_client
from src.db.models import DataQualityCheck
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -40,6 +41,22 @@ _CACHE_DIR = Path(tempfile.gettempdir()) / "profeto_injuries"
_cache_lock = asyncio.Lock() _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: def _compute_ttl_hours(date_str: str | None, fixture_id: int | None, league_id: int | None) -> float:
"""根据查询参数计算缓存 TTL(小时)。 """根据查询参数计算缓存 TTL(小时)。
@@ -91,6 +108,68 @@ def _write_cache_atomic(cache_file: Path, data: Any) -> None:
raise 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]: 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: if not api_key:
raise RuntimeError("API_FOOTBALL_KEY 未设置") raise RuntimeError("API_FOOTBALL_KEY 未设置")
# P2-4: 清理过期缓存文件
_cleanup_expired_cache()
cache_dir = _CACHE_DIR cache_dir = _CACHE_DIR
cache_dir.mkdir(parents=True, exist_ok=True) cache_dir.mkdir(parents=True, exist_ok=True)
# TTL 分级缓存 # TTL 分级缓存
ttl_hours = _compute_ttl_hours(date, fixture_id, league_id) 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 cache_file = cache_dir / cache_key
if cache_file.exists(): if cache_file.exists():
age_hours = (time.time() - cache_file.stat().st_mtime) / 3600 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 sqlalchemy.orm import selectinload
from src.data.team_names import normalize as normalize_name 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": []} result = {"count": 0, "inserted": 0, "errors": []}
batch_id = uuid.uuid4().hex[:16] 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], error_detail=str(e)[:2000],
raw_payload={"date": date}, raw_payload={"date": date},
status="pending", status="pending",
next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5),
) )
db.add(failure) db.add(failure)
await db.flush() await db.flush()
@@ -294,6 +378,7 @@ async def ingest_injuries(db, *, date: str | None = None) -> dict:
error_detail=fail["error"][:2000], error_detail=fail["error"][:2000],
raw_payload=raw, raw_payload=raw,
status="pending", status="pending",
next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5),
) )
db.add(failure) db.add(failure)
except Exception: 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") logger.warning("injuries final flush IntegrityError, falling back to per-record insert")
return await _ingest_injuries_fallback(db, pending_records, result) 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 控制事务 # 注意: 不在此处 commit,由调用方 UnitOfWork 控制事务
logger.info("injuries: fetched %d, inserted %d for %s", result["count"], result["inserted"], date) logger.info("injuries: fetched %d, inserted %d for %s", result["count"], result["inserted"], date)
return result return result
+126
View File
@@ -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)
+125 -5
View File
@@ -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.config import FDCO_TO_UNDERSTAT, LEAGUE_NAMES
from src.data.normalize import normalize_understat from src.data.normalize import normalize_understat
from src.data.sources import register 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__) 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}", "Referer": f"https://understat.com/league/{understat_league}/{season}",
} }
# 重试:网络错误 / 5xx / 429 # 重试:网络错误 / 5xx / 429
last_exc: Exception | None = None last_exc: Exception | None = None
for attempt in range(3): for attempt in range(3):
@@ -65,12 +66,10 @@ async def fetch_understat(league_code: str, season: int) -> list[dict]:
except Exception as e: except Exception as e:
last_exc = e last_exc = e
if attempt == 2: if attempt == 2:
raise raise RuntimeError(f"understat fetch failed: {e}") from e
delay = min(2 ** attempt, 8) + random.uniform(0, 1) delay = min(2 ** attempt, 8) + random.uniform(0, 1)
logger.warning("understat fetch failed, retry %d in %.1fs: %s", attempt + 1, delay, e) logger.warning("understat fetch failed, retry %d in %.1fs: %s", attempt + 1, delay, e)
await asyncio.sleep(delay) await asyncio.sleep(delay)
else:
raise RuntimeError(f"understat fetch failed: {last_exc}")
# understat 返回 JS 对象,需要提取 JSON # understat 返回 JS 对象,需要提取 JSON
text = resp.text 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: 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: for raw in raw_matches:
try: try:
source_record_id = str(raw.get("id", "")) source_record_id = str(raw.get("id", ""))
if not source_record_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 continue
bronze = RawEvent( bronze = RawEvent(
source_system="understat", source_system="understat",
@@ -117,6 +125,20 @@ async def _write_understat_bronze(db, raw_matches: list[dict], batch_id: str) ->
except Exception as e: except Exception as e:
logger.warning("understat raw_events flush failed: %s", 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( async def _write_ingest_failure(
db, db,
@@ -137,6 +159,7 @@ async def _write_ingest_failure(
error_detail=error_detail[:2000] if error_detail else None, error_detail=error_detail[:2000] if error_detail else None,
raw_payload=raw_payload, raw_payload=raw_payload,
status="pending", status="pending",
next_retry_at=datetime.now(timezone.utc) + timedelta(minutes=5),
) )
db.add(failure) db.add(failure)
await db.flush() await db.flush()
@@ -144,6 +167,64 @@ async def _write_ingest_failure(
logger.warning("understat ingest_failure write failed", exc_info=True) 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 @register
class UnderstatSource: class UnderstatSource:
"""understat xG 数据源(实现 DataSource 协议)。""" """understat xG 数据源(实现 DataSource 协议)。"""
@@ -268,20 +349,59 @@ class UnderstatSource:
) )
db.add(existing.stats) db.add(existing.stats)
await db.flush() 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 数据更新时覆盖旧值 # xG 覆盖更新模式:当 understat 数据更新时覆盖旧值
if existing.stats is not None: if existing.stats is not None:
xg_updated = False
if nm.home_xg is not None: if nm.home_xg is not None:
existing.stats.home_xg = nm.home_xg existing.stats.home_xg = nm.home_xg
existing.stats.xg_source = "understat" existing.stats.xg_source = "understat"
existing.stats.xg_updated_at = now existing.stats.xg_updated_at = now
existing.stats.xg_source_record_id = source_record_id existing.stats.xg_source_record_id = source_record_id
result["updated"] += 1 result["updated"] += 1
xg_updated = True
if nm.away_xg is not None: if nm.away_xg is not None:
existing.stats.away_xg = nm.away_xg existing.stats.away_xg = nm.away_xg
existing.stats.xg_source = "understat" existing.stats.xg_source = "understat"
existing.stats.xg_updated_at = now existing.stats.xg_updated_at = now
existing.stats.xg_source_record_id = source_record_id 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 控制事务 # 注意: 不在此处 commit,由调用方 UnitOfWork 控制事务
return result return result