P0: injuries 缓存原子写(tempfile+os.replace) + TTL 分级 + asyncio.Lock 并发安全 P1: 引入 Bronze 层 RawEvent 表(原始事件存档,不可变) P1: IngestFailure 死信表(失败持久化,支持重试恢复) P1: understat 改为 xG 覆盖更新模式 + xG 来源追踪字段 P1: 令牌桶限流器(TokenBucket)替换固定 sleep P1: DataQualityCheck 数据质量监控表 P1: DataLineage 血缘追踪表 P2: 回测并发控制 Semaphore(3) 新增文件: - src/data/rate_limiter.py (令牌桶限流器) - alembic/versions/0008 (4 张新表迁移) - alembic/versions/0009 (MatchStats xG 字段迁移)
122 lines
6.1 KiB
Python
122 lines
6.1 KiB
Python
"""新增 Bronze 层 + 死信表 + 数据质量表 + 血缘表
|
|
|
|
Revision ID: 0008_raw_event_and_ingest_failure
|
|
Revises: 0007_predictions_unique_constraint
|
|
Create Date: 2026-09-17
|
|
|
|
架构审查报告 P1 实施:
|
|
- raw_events: Bronze 层,不可变原始采集记录
|
|
- ingest_failures: 采集失败死信表
|
|
- data_quality_checks: 数据质量检查结果记录
|
|
- data_lineage: ETL 全过程元数据血缘追踪
|
|
"""
|
|
|
|
from typing import Sequence, Union
|
|
|
|
from alembic import op
|
|
import sqlalchemy as sa
|
|
from sqlalchemy.dialects.postgresql import JSONB
|
|
|
|
# revision identifiers, used by Alembic.
|
|
revision: str = '0008_raw_event_and_ingest_failure'
|
|
down_revision: Union[str, None] = '0007_predictions_unique_constraint'
|
|
branch_labels: Union[str, Sequence[str], None] = None
|
|
depends_on: Union[str, Sequence[str], None] = None
|
|
|
|
|
|
def upgrade() -> None:
|
|
# ── raw_events: Bronze 层原始记录 ──
|
|
op.create_table(
|
|
"raw_events",
|
|
sa.Column("id", sa.Integer(), primary_key=True),
|
|
sa.Column("source_system", sa.String(30), nullable=False),
|
|
sa.Column("source_record_id", sa.String(100), nullable=False),
|
|
sa.Column("raw_payload", JSONB(), nullable=False),
|
|
sa.Column("ingested_at", sa.DateTime(timezone=True), nullable=False),
|
|
sa.Column("ingest_batch_id", sa.String(64), nullable=True),
|
|
)
|
|
op.create_index("ix_raw_events_batch", "raw_events", ["ingest_batch_id"])
|
|
op.create_index("ix_raw_events_source_ingested", "raw_events", ["source_system", "ingested_at"])
|
|
op.create_unique_constraint("uq_raw_events_source_record", "raw_events", ["source_system", "source_record_id"])
|
|
|
|
# ── ingest_failures: 采集失败死信表 ──
|
|
op.create_table(
|
|
"ingest_failures",
|
|
sa.Column("id", sa.Integer(), primary_key=True),
|
|
sa.Column("source_system", sa.String(30), nullable=False),
|
|
sa.Column("entity_type", sa.String(30), nullable=False),
|
|
sa.Column("source_record_id", sa.String(100), nullable=True),
|
|
sa.Column("error_type", sa.String(50), nullable=False),
|
|
sa.Column("error_detail", sa.Text(), nullable=True),
|
|
sa.Column("raw_payload", JSONB(), nullable=True),
|
|
sa.Column("retry_count", sa.Integer(), nullable=False, server_default="0"),
|
|
sa.Column("next_retry_at", sa.DateTime(timezone=True), nullable=True),
|
|
sa.Column("status", sa.String(20), nullable=False, server_default="pending"),
|
|
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
|
|
sa.Column("resolved_at", sa.DateTime(timezone=True), nullable=True),
|
|
)
|
|
op.create_index("ix_ingest_failures_status_next_retry", "ingest_failures", ["status", "next_retry_at"])
|
|
op.create_index("ix_ingest_failures_source", "ingest_failures", ["source_system", "entity_type"])
|
|
op.create_check_constraint("ck_ingest_failures_status", "ingest_failures",
|
|
"status IN ('pending', 'retrying', 'resolved', 'abandoned')")
|
|
op.create_check_constraint("ck_ingest_failures_retry_nonneg", "ingest_failures", "retry_count >= 0")
|
|
|
|
# ── data_quality_checks: 数据质量检查记录 ──
|
|
op.create_table(
|
|
"data_quality_checks",
|
|
sa.Column("id", sa.Integer(), primary_key=True),
|
|
sa.Column("check_name", sa.String(100), nullable=False),
|
|
sa.Column("entity_type", sa.String(30), nullable=False),
|
|
sa.Column("entity_id", sa.String(50), nullable=True),
|
|
sa.Column("expected_value", sa.Text(), nullable=True),
|
|
sa.Column("actual_value", sa.Text(), nullable=True),
|
|
sa.Column("passed", sa.Boolean(), nullable=False),
|
|
sa.Column("severity", sa.String(10), nullable=False, server_default="warning"),
|
|
sa.Column("detail", sa.Text(), nullable=True),
|
|
sa.Column("checked_at", sa.DateTime(timezone=True), nullable=False),
|
|
)
|
|
op.create_index("ix_dqc_check_time", "data_quality_checks", ["check_name", "checked_at"])
|
|
op.create_index("ix_dqc_entity", "data_quality_checks", ["entity_type", "entity_id"])
|
|
op.create_index("ix_dqc_severity_passed", "data_quality_checks", ["severity", "passed"])
|
|
op.create_check_constraint("ck_dqc_severity", "data_quality_checks",
|
|
"severity IN ('info', 'warning', 'critical')")
|
|
|
|
# ── data_lineage: ETL 血缘追踪 ──
|
|
op.create_table(
|
|
"data_lineage",
|
|
sa.Column("id", sa.Integer(), primary_key=True),
|
|
sa.Column("source_system", sa.String(30), nullable=False),
|
|
sa.Column("source_record_id", sa.String(100), nullable=False),
|
|
sa.Column("target_table", sa.String(50), nullable=False),
|
|
sa.Column("target_id", sa.Integer(), nullable=True),
|
|
sa.Column("transform_name", sa.String(100), nullable=False),
|
|
sa.Column("transform_detail", sa.Text(), nullable=True),
|
|
sa.Column("batch_id", sa.String(64), nullable=True),
|
|
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
|
|
)
|
|
op.create_index("ix_lineage_source", "data_lineage", ["source_system", "source_record_id"])
|
|
op.create_index("ix_lineage_target", "data_lineage", ["target_table", "target_id"])
|
|
op.create_index("ix_lineage_batch", "data_lineage", ["batch_id"])
|
|
|
|
|
|
def downgrade() -> None:
|
|
# 逆序删除
|
|
op.drop_index("ix_lineage_batch", table_name="data_lineage")
|
|
op.drop_index("ix_lineage_target", table_name="data_lineage")
|
|
op.drop_index("ix_lineage_source", table_name="data_lineage")
|
|
op.drop_table("data_lineage")
|
|
|
|
op.drop_index("ix_dqc_severity_passed", table_name="data_quality_checks")
|
|
op.drop_index("ix_dqc_entity", table_name="data_quality_checks")
|
|
op.drop_index("ix_dqc_check_time", table_name="data_quality_checks")
|
|
op.drop_table("data_quality_checks")
|
|
|
|
op.drop_index("ix_ingest_failures_source", table_name="ingest_failures")
|
|
op.drop_index("ix_ingest_failures_status_next_retry", table_name="ingest_failures")
|
|
op.drop_table("ingest_failures")
|
|
|
|
op.drop_constraint("uq_raw_events_source_record", table_name="raw_events", type_="unique")
|
|
op.drop_index("ix_raw_events_source_ingested", table_name="raw_events")
|
|
op.drop_index("ix_raw_events_batch", table_name="raw_events")
|
|
op.drop_table("raw_events")
|