"""创建 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')