feat: 数据管线架构优化 + 前端部署配置

- Bronze 层 RawEvent 表(原始事件存档)
- IngestFailure 死信表(失败持久化+重试)
- DataQualityCheck 数据质量监控
- DataLineage 血缘追踪
- MatchStats xG 追踪字段
- Nginx 前端服务配置
This commit is contained in:
shangfangjian
2026-09-17 01:15:51 +08:00
parent e89ab1a0c9
commit 312778d995
3 changed files with 174 additions and 0 deletions
@@ -0,0 +1,107 @@
"""创建 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')
@@ -0,0 +1,29 @@
"""为 MatchStats 添加 xG 追踪字段
Revision ID: 0009_match_stats_xg_fields
Revises: 0008_raw_event_and_ingest_failure
Create Date: 2026-09-17
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
revision: str = '0009_match_stats_xg_fields'
down_revision: Union[str, None] = '0008_raw_event_and_ingest_failure'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.add_column('match_stats', sa.Column('xg_source', sa.String(30), nullable=True))
op.add_column('match_stats', sa.Column('xg_updated_at', sa.DateTime(timezone=True), nullable=True))
op.add_column('match_stats', sa.Column('xg_source_record_id', sa.String(100), nullable=True))
def downgrade() -> None:
op.drop_column('match_stats', 'xg_source_record_id')
op.drop_column('match_stats', 'xg_updated_at')
op.drop_column('match_stats', 'xg_source')
+38
View File
@@ -0,0 +1,38 @@
server {
listen 80;
server_name localhost;
# Serve static files
root /usr/share/nginx/html;
index index.html;
# API proxy
location /api/ {
proxy_pass http://api:8000/api/;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection 'upgrade';
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
proxy_cache_bypass $http_upgrade;
}
# Health check proxy
location /health {
proxy_pass http://api:8000/health;
proxy_http_version 1.1;
proxy_set_header Host $host;
}
# SPA routing - serve index.html for all routes
location / {
try_files $uri $uri/ /index.html;
}
# Gzip compression
gzip on;
gzip_types text/plain text/css application/json application/javascript text/xml application/xml application/xml+rss text/javascript;
gzip_min_length 1000;
}