Files
Profeto/alembic/versions/0008_raw_event_and_ingest_failure.py
shangfangjian 312778d995 feat: 数据管线架构优化 + 前端部署配置
- Bronze 层 RawEvent 表(原始事件存档)
- IngestFailure 死信表(失败持久化+重试)
- DataQualityCheck 数据质量监控
- DataLineage 血缘追踪
- MatchStats xG 追踪字段
- Nginx 前端服务配置
2026-09-17 01:15:51 +08:00

108 lines
4.9 KiB
Python

"""创建 Bronze 层、死信表、质量监控、血缘追踪 4 张新表
Revision ID: 0008_raw_event_and_ingest_failure
Revises: 0007_predictions_unique_constraint
Create Date: 2026-09-17
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
from sqlalchemy.dialects.postgresql import JSONB
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:
# 1. RawEvent - Bronze 层原始事件存档
op.create_table(
'raw_events',
sa.Column('id', sa.BigInteger(), primary_key=True),
sa.Column('source_system', sa.String(50), 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), server_default=sa.func.now()),
sa.Column('ingest_batch_id', sa.String(36), nullable=True),
sa.UniqueConstraint('source_system', 'source_record_id', name='uq_raw_event'),
)
op.create_index('ix_raw_event_batch', 'raw_events', ['ingest_batch_id'])
# 2. IngestFailure - 采集失败死信表
op.create_table(
'ingest_failures',
sa.Column('id', sa.BigInteger(), primary_key=True),
sa.Column('source_system', sa.String(50), nullable=False),
sa.Column('entity_type', sa.String(50), 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(), server_default='0'),
sa.Column('next_retry_at', sa.DateTime(timezone=True), nullable=True),
sa.Column('status', sa.String(20), server_default='pending'),
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
sa.Column('resolved_at', sa.DateTime(timezone=True), nullable=True),
)
op.create_index('ix_ingest_failure_status', 'ingest_failures', ['status', 'next_retry_at'])
op.create_check_constraint(
'ck_ingest_failure_status',
'ingest_failures',
"status IN ('pending', 'retrying', 'resolved', 'abandoned')",
)
# 3. DataQualityCheck - 数据质量监控
op.create_table(
'data_quality_checks',
sa.Column('id', sa.BigInteger(), primary_key=True),
sa.Column('check_name', sa.String(100), nullable=False),
sa.Column('entity_type', sa.String(50), nullable=False),
sa.Column('entity_id', sa.Integer(), nullable=True),
sa.Column('expected_value', sa.Float(), nullable=True),
sa.Column('actual_value', sa.Float(), nullable=False),
sa.Column('passed', sa.Boolean(), nullable=False),
sa.Column('severity', sa.String(10), server_default='warning'),
sa.Column('detail', JSONB(), nullable=True),
sa.Column('checked_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
)
op.create_index('ix_dqc_checked_at', 'data_quality_checks', ['checked_at'])
op.create_index('ix_dqc_entity', 'data_quality_checks', ['entity_type', 'entity_id'])
# 4. DataLineage - ETL 血缘追踪
op.create_table(
'data_lineage',
sa.Column('id', sa.BigInteger(), primary_key=True),
sa.Column('source_system', sa.String(50), 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(50), nullable=False),
sa.Column('transform_detail', JSONB(), nullable=True),
sa.Column('batch_id', sa.String(36), nullable=True),
sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now()),
)
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_entity', table_name='data_quality_checks')
op.drop_index('ix_dqc_checked_at', table_name='data_quality_checks')
op.drop_table('data_quality_checks')
op.drop_index('ix_ingest_failure_status', table_name='ingest_failures')
op.drop_table('ingest_failures')
op.drop_index('ix_raw_event_batch', table_name='raw_events')
op.drop_table('raw_events')