引言:选举与票房的奇妙交汇

在现代政治生态中,”辉煌大选”不仅仅是一场政治角逐,更是一场全民参与的”票房大戏”。实时票房追踪系统通过数据可视化和深度分析,将枯燥的选票数字转化为生动的政治风向标。本文将详细解析如何构建一个完整的实时票房追踪系统,涵盖数据采集、处理、可视化和深度分析的全过程。

一、实时票房追踪系统架构设计

1.1 系统核心组件

一个完整的实时票房追踪系统需要以下几个关键组件:

  • 数据采集层:负责从多个数据源实时获取选票数据
  • 数据处理层:清洗、验证和聚合原始数据
  • 存储层:高效存储历史数据和实时数据
  • 分析层:进行趋势预测和深度分析
  • 展示层:通过可视化界面呈现数据

1.2 技术栈选择

# 系统核心依赖
dependencies = {
    "数据采集": ["requests", "aiohttp", "BeautifulSoup", "Selenium"],
    "数据处理": ["pandas", "numpy", "scikit-learn"],
    "实时流处理": ["Kafka", "Redis", "WebSocket"],
    "数据存储": ["PostgreSQL", "MongoDB", "InfluxDB"],
    "可视化": ["Plotly", "Dash", "D3.js"],
    "后端服务": ["FastAPI", "Flask", "Celery"],
    "监控告警": ["Prometheus", "Grafana"]
}

二、数据采集与实时更新机制

2.1 多源数据采集策略

为了确保数据的准确性和实时性,我们需要从多个权威数据源同时采集:

import asyncio
import aiohttp
import time
from datetime import datetime
import json

class ElectionDataCollector:
    def __init__(self):
        self.sources = {
            'official_election_board': 'https://api.election.gov/v1/results',
            'major_news_networks': [
                'https://api.cnn.com/election',
                'https://api.foxnews.com/election',
                'https://api.nbcnews.com/election'
            ],
            'academic_institutions': 'https://data.harvard.edu/election-tracker'
        }
        self.data_cache = {}
        
    async def fetch_official_data(self, session, state_code):
        """从官方选举委员会获取数据"""
        url = f"{self.sources['official_election_board']}/{state_code}"
        try:
            async with session.get(url, timeout=10) as response:
                if response.status == 200:
                    data = await response.json()
                    return {
                        'source': 'official',
                        'timestamp': datetime.utcnow().isoformat(),
                        'state': state_code,
                        'data': data,
                        'confidence': 1.0
                    }
        except Exception as e:
            print(f"Error fetching official data: {e}")
            return None
    
    async def fetch_news_network_data(self, session, network_url):
        """从新闻网络获取预测数据"""
        try:
            async with session.get(network_url, timeout=5) as response:
                if response.status == 200:
                    data = await response.json()
                    return {
                        'source': 'news_network',
                        'timestamp': datetime.utcnow().isoformat(),
                        'data': data,
                        'confidence': 0.85
                    }
        except Exception as e:
            print(f"Error fetching news data: {e}")
            return None
    
    async def collect_all_sources(self, state_code):
        """并行采集所有数据源"""
        async with aiohttp.ClientSession() as session:
            tasks = [
                self.fetch_official_data(session, state_code),
                *[self.fetch_news_network_data(session, url) 
                  for url in self.sources['major_news_networks']]
            ]
            results = await asyncio.gather(*tasks, return_exceptions=True)
            
            # 过滤失败请求
            valid_results = [r for r in results if r and not isinstance(r, Exception)]
            return self.validate_and_merge(valid_results)
    
    def validate_and_merge(self, data_list):
        """数据验证与合并"""
        if not data_list:
            return None
            
        # 按置信度排序
        data_list.sort(key=lambda x: x['confidence'], reverse=True)
        
        # 选择置信度最高的数据作为主数据源
        primary_data = data_list[0]
        
        # 验证数据一致性
        for data in data_list[1:]:
            if not self.verify_consistency(primary_data, data):
                print(f"Warning: Data inconsistency detected from {data['source']}")
        
        return {
            'primary_source': primary_data['source'],
            'timestamp': primary_data['timestamp'],
            'state': primary_data['state'],
            'vote_counts': primary_data['data']['vote_counts'],
            'reporting_percentage': primary_data['data']['reporting_percentage'],
            'confidence_score': self.calculate_confidence(data_list),
            'all_sources': data_list
        }
    
    def verify_consistency(self, primary, secondary):
        """验证数据一致性"""
        # 检查投票数差异是否在合理范围内(5%阈值)
        primary_votes = primary['data']['vote_counts']['total']
        secondary_votes = secondary['data']['vote_counts']['total']
        diff_ratio = abs(primary_votes - secondary_votes) / max(primary_votes, 1)
        return diff_ratio < 0.05
    
    def calculate_confidence(self, data_list):
        """计算综合置信度"""
        if len(data_list) == 1:
            return data_list[0]['confidence']
        
        # 加权平均,官方数据权重更高
        weights = {'official': 0.6, 'news_network': 0.3, 'academic': 0.1}
        total_confidence = 0
        total_weight = 0
        
        for data in data_list:
            source_type = data['source']
            weight = weights.get(source_type, 0.1)
            total_confidence += data['confidence'] * weight
            total_weight += weight
        
        return total_confidence / total_weight if total_weight > 0 else 0

# 使用示例
async def main():
    collector = ElectionDataCollector()
    result = await collector.collect_all_sources('PA')
    print(json.dumps(result, indent=2))

# 运行
# asyncio.run(main())

2.2 实时流处理架构

from kafka import KafkaProducer, KafkaConsumer
import redis
import json

class RealTimeStreamProcessor:
    def __init__(self):
        self.kafka_bootstrap_servers = 'localhost:9092'
        self.redis_client = redis.Redis(host='localhost', port=6379, db=0)
        
    def setup_kafka_topics(self):
        """创建Kafka主题"""
        from kafka.admin import KafkaAdminClient, NewTopic
        
        admin = KafkaAdminClient(bootstrap_servers=self.kafka_bootstrap_servers)
        
        topics = [
            NewTopic(name="raw-election-data", num_partitions=3, replication_factor=1),
            NewTopic(name="processed-election-data", num_partitions=3, replication_factor=1),
            NewTopic(name="election-analytics", num_partitions=2, replication_factor=1)
        ]
        
        try:
            admin.create_topics(topics)
            print("Kafka topics created successfully")
        except Exception as e:
            print(f"Topic creation error: {e}")
    
    def produce_raw_data(self, data):
        """生产原始数据到Kafka"""
        producer = KafkaProducer(
            bootstrap_servers=self.kafka_bootstrap_servers,
            value_serializer=lambda v: json.dumps(v).encode('utf-8')
        )
        
        # 添加元数据
        message = {
            'timestamp': int(time.time()),
            'data': data,
            'type': 'raw_result'
        }
        
        producer.send('raw-election-data', message)
        producer.flush()
        print(f"Produced raw data: {data['state']}")
    
    def consume_and_process(self):
        """消费并处理数据"""
        consumer = KafkaConsumer(
            'raw-election-data',
            bootstrap_servers=self.kafka_bootstrap_servers,
            value_deserializer=lambda m: json.loads(m.decode('utf-8')),
            auto_offset_reset='latest',
            group_id='election-processor-group'
        )
        
        for message in consumer:
            raw_data = message.value
            processed_data = self.process_data(raw_data)
            
            # 发布到处理后的主题
            self.produce_processed_data(processed_data)
            
            # 更新Redis缓存
            self.update_redis_cache(processed_data)
    
    def process_data(self, raw_data):
        """数据处理逻辑"""
        data = raw_data['data']
        
        # 计算关键指标
        processed = {
            'state': data['state'],
            'timestamp': raw_data['timestamp'],
            'total_votes': data['vote_counts']['total'],
            'democrat_votes': data['vote_counts']['democrat'],
            'republican_votes': data['vote_counts']['republican'],
            'other_votes': data['vote_counts']['other'],
            'reporting_percentage': data['reporting_percentage'],
            'democrat_share': data['vote_counts']['democrat'] / data['vote_counts']['total'],
            'republican_share': data['vote_counts']['republican'] / data['vote_counts']['total'],
            'lead_margin': data['vote_counts']['democrat'] - data['vote_counts']['republican'],
            'projected_winner': self.project_winner(data)
        }
        
        return processed
    
    def project_winner(self, data):
        """预测胜者"""
        total = data['vote_counts']['total']
        if total == 0:
            return None
        
        dem_share = data['vote_counts']['democrat'] / total
        rep_share = data['vote_counts']['republican'] / total
        reporting = data['reporting_percentage']
        
        # 如果报告率超过95%且领先优势明显,预测胜者
        if reporting > 95:
            if dem_share > rep_share + 0.02:
                return 'democrat'
            elif rep_share > dem_share + 0.02:
                return 'republican'
        
        return 'too_close_to_call'
    
    def update_redis_cache(self, data):
        """更新Redis缓存"""
        key = f"election:{data['state']}"
        self.redis_client.setex(key, 3600, json.dumps(data))  # 1小时过期
        
        # 更新实时排行榜
        self.redis_client.zadd(
            "election:lead_margin",
            {data['state']: data['lead_margin']}
        )
    
    def produce_processed_data(self, data):
        """发送处理后的数据"""
        producer = KafkaProducer(
            bootstrap_servers=self.kafka_bootstrap_servers,
            value_serializer=lambda v: json.dumps(v).encode('utf-8')
        )
        
        producer.send('processed-election-data', data)
        producer.flush()

# 使用示例
# processor = RealTimeStreamProcessor()
# processor.setup_kafka_topics()
# processor.consume_and_process()

三、数据存储与查询优化

3.1 混合存储策略

from sqlalchemy import create_engine, Column, Integer, String, Float, DateTime, JSON
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
from influxdb_client import InfluxDBClient, Point
from datetime import datetime
import pandas as pd

Base = declarative_base()

class ElectionResult(Base):
    """PostgreSQL关系型数据模型"""
    __tablename__ = 'election_results'
    
    id = Column(Integer, primary_key=True)
    state = Column(String(2), nullable=False, index=True)
    timestamp = Column(DateTime, nullable=False, index=True)
    democrat_votes = Column(Integer)
    republican_votes = Column(Integer)
    other_votes = Column(Integer)
    total_votes = Column(Integer)
    reporting_percentage = Column(Float)
    confidence_score = Column(Float)
    raw_data = Column(JSON)
    
    def __repr__(self):
        return f"<ElectionResult(state={self.state}, timestamp={self.timestamp})>"

class HybridStorage:
    def __init__(self):
        # PostgreSQL for transactional data
        self.pg_engine = create_engine(
            'postgresql://user:pass@localhost:5432/election_db',
            pool_size=10,
            max_overflow=20
        )
        Base.metadata.create_all(self.pg_engine)
        self.Session = sessionmaker(bind=self.pg_engine)
        
        # InfluxDB for time-series data
        self.influx_client = InfluxDBClient(
            url="http://localhost:8086",
            token="your-token",
            org="your-org"
        )
        self.influx_write_api = self.influx_client.write_api()
        
    def store_postgresql(self, data):
        """存储到PostgreSQL"""
        session = self.Session()
        try:
            result = ElectionResult(
                state=data['state'],
                timestamp=datetime.fromtimestamp(data['timestamp']),
                democrat_votes=data['democrat_votes'],
                republican_votes=data['republican_votes'],
                other_votes=data['other_votes'],
                total_votes=data['total_votes'],
                reporting_percentage=data['reporting_percentage'],
                confidence_score=data['confidence_score'],
                raw_data=data
            )
            session.add(result)
            session.commit()
            print(f"Stored in PostgreSQL: {data['state']}")
        except Exception as e:
            session.rollback()
            print(f"PostgreSQL storage error: {e}")
        finally:
            session.close()
    
    def store_influxdb(self, data):
        """存储到InfluxDB"""
        point = Point("election_results") \
            .tag("state", data['state']) \
            .field("democrat_votes", data['democrat_votes']) \
            .field("republican_votes", data['republican_votes']) \
            .field("other_votes", data['other_votes']) \
            .field("total_votes", data['total_votes']) \
            .field("reporting_percentage", data['reporting_percentage']) \
            .field("confidence_score", data['confidence_score']) \
            .field("lead_margin", data['lead_margin']) \
            .field("democrat_share", data['democrat_share']) \
            .field("republican_share", data['republican_share']) \
            .time(data['timestamp'], write_precision='s')
        
        self.influx_write_api.write(
            bucket="election_bucket",
            org="your-org",
            record=point
        )
        print(f"Stored in InfluxDB: {data['state']}")
    
    def store_all(self, data):
        """同时存储到所有系统"""
        self.store_postgresql(data)
        self.store_influxdb(data)
    
    def query_historical_trends(self, state, hours=24):
        """查询历史趋势"""
        # 从InfluxDB查询时间序列数据
        query_api = self.influx_client.query_api()
        
        query = f'''
        from(bucket: "election_bucket")
          |> range(start: -{hours}h)
          |> filter(fn: (r) => r._measurement == "election_results")
          |> filter(fn: (r) => r.state == "{state}")
          |> aggregateWindow(every: 5m, fn: mean, createEmpty: false)
        '''
        
        result = query_api.query(query)
        return result
    
    def get_latest_results(self, states=None):
        """获取最新结果"""
        session = self.Session()
        try:
            query = session.query(ElectionResult)
            
            if states:
                query = query.filter(ElectionResult.state.in_(states))
            
            # 获取每个州的最新记录
            subquery = session.query(
                ElectionResult.state,
                func.max(ElectionResult.timestamp).label('max_timestamp')
            ).group_by(ElectionResult.state).subquery()
            
            query = query.join(
                subquery,
                (ElectionResult.state == subquery.c.state) &
                (ElectionResult.timestamp == subquery.c.max_timestamp)
            )
            
            return query.all()
        finally:
            session.close()

# 使用示例
# storage = HybridStorage()
# storage.store_all(processed_data)

四、深度分析与预测模型

4.1 趋势预测算法

import numpy as np
from sklearn.linear_model import LinearRegression
from sklearn.ensemble import RandomForestRegressor
from scipy import stats
import warnings
warnings.filterwarnings('ignore')

class ElectionPredictor:
    def __init__(self):
        self.models = {}
        self.historical_data = None
        
    def load_historical_data(self, storage):
        """加载历史数据"""
        # 从数据库获取过去选举数据
        # 这里使用模拟数据
        self.historical_data = self._generate_mock_historical_data()
        
    def _generate_mock_historical_data(self):
        """生成模拟历史数据用于演示"""
        np.random.seed(42)
        states = ['PA', 'MI', 'WI', 'AZ', 'GA', 'NV', 'NC']
        
        data = []
        for state in states:
            for election_year in [2016, 2020]:
                for reporting_pct in range(10, 101, 10):
                    # 模拟投票模式
                    base_dem = 0.48 if election_year == 2020 else 0.47
                    base_rep = 0.51 if election_year == 2020 else 0.50
                    
                    # 添加随机波动
                    dem_share = base_dem + np.random.normal(0, 0.02)
                    rep_share = base_rep + np.random.normal(0, 0.02)
                    
                    # 报告率越高,数据越稳定
                    volatility = 0.05 * (100 - reporting_pct) / 100
                    
                    final_dem = dem_share + np.random.normal(0, volatility)
                    final_rep = rep_share + np.random.normal(0, volatility)
                    
                    data.append({
                        'state': state,
                        'year': election_year,
                        'reporting_pct': reporting_pct,
                        'dem_share': final_dem,
                        'rep_share': final_rep,
                        'final_dem_share': final_dem,
                        'final_rep_share': final_rep,
                        'lead_margin': (final_dem - final_rep) * 10000
                    })
        
        return pd.DataFrame(data)
    
    def train_state_models(self):
        """为每个州训练预测模型"""
        if self.historical_data is None:
            raise ValueError("Historical data not loaded")
        
        for state in self.historical_data['state'].unique():
            state_data = self.historical_data[self.historical_data['state'] == state]
            
            # 特征工程
            X = state_data[['reporting_pct', 'year']].values
            y = state_data['lead_margin'].values
            
            # 使用随机森林模型
            model = RandomForestRegressor(
                n_estimators=100,
                max_depth=5,
                random_state=42
            )
            model.fit(X, y)
            self.models[state] = model
            
            print(f"Model trained for {state}: R² = {model.score(X, y):.3f}")
    
    def predict_current_election(self, current_data):
        """预测当前选举结果"""
        predictions = {}
        
        for state, data in current_data.items():
            if state not in self.models:
                # 如果没有该州模型,使用简单平均
                predictions[state] = {
                    'predicted_dem_share': data['democrat_share'],
                    'predicted_rep_share': data['republican_share'],
                    'confidence': 0.5,
                    'model_used': 'fallback'
                }
                continue
            
            # 准备特征
            X = np.array([[data['reporting_percentage'], 2024]])
            
            # 预测领先优势
            predicted_margin = self.models[state].predict(X)[0]
            
            # 转换为得票率
            total_share = data['democrat_share'] + data['republican_share']
            predicted_dem_share = (total_share + predicted_margin / 10000) / 2
            predicted_rep_share = total_share - predicted_dem_share
            
            # 计算置信度
            confidence = self.calculate_confidence_interval(
                state, 
                data['reporting_percentage']
            )
            
            predictions[state] = {
                'predicted_dem_share': predicted_dem_share,
                'predicted_rep_share': predicted_rep_share,
                'predicted_margin': predicted_margin,
                'confidence': confidence,
                'model_used': 'random_forest',
                'current_reporting': data['reporting_percentage']
            }
        
        return predictions
    
    def calculate_confidence_interval(self, state, reporting_pct):
        """计算预测置信区间"""
        # 报告率越高,置信度越高
        base_confidence = min(reporting_pct / 100, 1.0)
        
        # 如果有历史模型,考虑模型稳定性
        if state in self.models:
            model = self.models[state]
            # 使用模型在验证集上的表现作为稳定性指标
            # 这里简化处理
            stability = 0.85
            return base_confidence * stability
        
        return base_confidence * 0.6  # 较低的置信度
    
    def detect_anomalies(self, current_data):
        """检测异常数据"""
        anomalies = []
        
        for state, data in current_data.items():
            # 检查报告率是否异常
            if data['reporting_percentage'] > 100:
                anomalies.append({
                    'state': state,
                    'type': 'reporting_over_100',
                    'value': data['reporting_percentage']
                })
            
            # 检查投票数是否异常增长
            if data['total_votes'] > 10000000:  # 假设上限
                anomalies.append({
                    'state': state,
                    'type': 'vote_count_anomaly',
                    'value': data['total_votes']
                })
            
            # 检查得票率是否合理
            total_share = data['democrat_share'] + data['republican_share']
            if total_share < 0.95 or total_share > 1.05:
                anomalies.append({
                    'state': state,
                    'type': 'vote_share_anomaly',
                    'value': total_share
                })
        
        return anomalies

# 使用示例
# predictor = ElectionPredictor()
# predictor.load_historical_data(storage)
# predictor.train_state_models()
# predictions = predictor.predict_current_election(current_data)

5. 可视化与实时展示

5.1 交互式仪表板

import plotly.graph_objects as go
import plotly.express as px
from plotly.subplots import make_subplots
import dash
from dash import dcc, html
from dash.dependencies import Input, Output
import threading

class ElectionDashboard:
    def __init__(self):
        self.app = dash.Dash(__name__)
        self.setup_layout()
        self.setup_callbacks()
        
    def setup_layout(self):
        """设置仪表板布局"""
        self.app.layout = html.Div([
            html.H1("辉煌大选实时票房追踪", 
                   style={'textAlign': 'center', 'color': '#2c3e50'}),
            
            # 关键指标卡片
            html.Div([
                html.Div([
                    html.H3("总投票数"),
                    html.Div(id='total-votes', 
                            style={'fontSize': '24px', 'color': '#3498db'})
                ], className='metric-card'),
                
                html.Div([
                    html.H3("领先州数"),
                    html.Div(id='lead-states',
                            style={'fontSize': '24px', 'color': '#e74c3c'})
                ], className='metric-card'),
                
                html.Div([
                    html.H3("预测胜率"),
                    html.Div(id='win-probability',
                            style={'fontSize': '24px', 'color': '#27ae60'})
                ], className='metric-card')
            ], style={'display': 'flex', 'justifyContent': 'space-around'}),
            
            # 实时地图
            html.Div([
                html.H3("实时选举地图"),
                dcc.Graph(id='election-map')
            ]),
            
            # 趋势图表
            html.Div([
                html.H3("关键州趋势"),
                dcc.Dropdown(
                    id='state-selector',
                    options=[],
                    value='PA'
                ),
                dcc.Graph(id='trend-chart')
            ]),
            
            # 数据表格
            html.Div([
                html.H3("详细数据"),
                html.Div(id='data-table')
            ]),
            
            # 自动刷新控制
            dcc.Interval(
                id='interval-component',
                interval=5*1000,  # 每5秒更新
                n_intervals=0
            )
        ], style={'padding': '20px', 'fontFamily': 'Arial'})
    
    def setup_callbacks(self):
        """设置回调函数"""
        
        @self.app.callback(
            [Output('total-votes', 'children'),
             Output('lead-states', 'children'),
             Output('win-probability', 'children'),
             Output('election-map', 'figure'),
             Output('trend-chart', 'figure'),
             Output('data-table', 'children'),
             Output('state-selector', 'options')],
            [Input('interval-component', 'n_intervals')]
        )
        def update_dashboard(n):
            # 获取最新数据
            current_data = self.get_current_data()
            
            # 计算关键指标
            total_votes = sum(d['total_votes'] for d in current_data.values())
            lead_states = len([s for s, d in current_data.items() 
                             if d['lead_margin'] > 0])
            
            # 计算预测胜率
            win_prob = self.calculate_win_probability(current_data)
            
            # 生成地图
            map_fig = self.create_election_map(current_data)
            
            # 生成趋势图
            state = 'PA'  # 默认选择
            trend_fig = self.create_trend_chart(state)
            
            # 生成数据表格
            table = self.create_data_table(current_data)
            
            # 州选择器选项
            options = [{'label': s, 'value': s} for s in current_data.keys()]
            
            return [
                f"{total_votes:,}",
                f"{lead_states}",
                f"{win_prob:.1%}",
                map_fig,
                trend_fig,
                table,
                options
            ]
        
        @self.app.callback(
            Output('trend-chart', 'figure'),
            [Input('state-selector', 'value')]
        )
        def update_trend_chart(selected_state):
            return self.create_trend_chart(selected_state)
    
    def get_current_data(self):
        """获取当前数据(模拟)"""
        # 实际应用中从Redis或数据库获取
        return {
            'PA': {'total_votes': 6500000, 'lead_margin': 150000, 
                   'democrat_share': 0.52, 'republican_share': 0.48,
                   'reporting_percentage': 85},
            'MI': {'total_votes': 5200000, 'lead_margin': -80000,
                   'democrat_share': 0.49, 'republican_share': 0.51,
                   'reporting_percentage': 78},
            'AZ': {'total_votes': 3200000, 'lead_margin': 120000,
                   'democrat_share': 0.53, 'republican_share': 0.47,
                   'reporting_percentage': 92},
            'GA': {'total_votes': 4800000, 'lead_margin': -50000,
                   'democrat_share': 0.485, 'republican_share': 0.515,
                   'reporting_percentage': 88}
        }
    
    def calculate_win_probability(self, data):
        """计算获胜概率"""
        total_prob = 0
        count = 0
        
        for state, info in data.items():
            if info['lead_margin'] > 0:
                # 根据领先优势和报告率加权
                margin_weight = min(abs(info['lead_margin']) / 100000, 2.0)
                reporting_weight = info['reporting_percentage'] / 100
                prob = 0.5 + (margin_weight * reporting_weight * 0.2)
                total_prob += prob
                count += 1
        
        return total_prob / max(count, 1)
    
    def create_election_map(self, data):
        """创建选举地图"""
        states = list(data.keys())
        margins = [data[s]['lead_margin'] for s in states]
        colors = ['#e74c3c' if m < 0 else '#3498db' for m in margins]
        
        fig = go.Figure(data=[
            go.Bar(
                x=states,
                y=margins,
                marker_color=colors,
                text=[f'{m/1000:.1f}k' for m in margins],
                textposition='outside'
            )
        ])
        
        fig.update_layout(
            title="各州领先优势(单位:千票)",
            xaxis_title="州",
            yaxis_title="领先优势",
            showlegend=False
        )
        
        return fig
    
    def create_trend_chart(self, state):
        """创建趋势图表"""
        # 模拟历史数据
        times = list(range(0, 25, 2))
        dem_trend = [50 + np.sin(t/5)*2 + np.random.normal(0, 0.5) for t in times]
        rep_trend = [50 - np.sin(t/5)*2 + np.random.normal(0, 0.5) for t in times]
        
        fig = go.Figure()
        fig.add_trace(go.Scatter(
            x=times,
            y=dem_trend,
            mode='lines+markers',
            name='民主党',
            line=dict(color='#3498db', width=3)
        ))
        fig.add_trace(go.Scatter(
            x=times,
            y=rep_trend,
            mode='lines+markers',
            name='共和党',
            line=dict(color='#e74c3c', width=3)
        ))
        
        fig.update_layout(
            title=f"{state}州得票率趋势",
            xaxis_title="时间(小时)",
            yaxis_title="得票率(%)",
            hovermode='x unified'
        )
        
        return fig
    
    def create_data_table(self, data):
        """创建数据表格"""
        table_data = []
        for state, info in data.items():
            table_data.append({
                '州': state,
                '报告率': f"{info['reporting_percentage']}%",
                '民主党': f"{info['democrat_share']:.1%}",
                '共和党': f"{info['republican_share']:.1%}",
                '领先': f"{info['lead_margin']:,}"
            })
        
        # 转换为HTML表格
        df = pd.DataFrame(table_data)
        return html.Table([
            html.Thead(
                html.Tr([html.Th(col) for col in df.columns])
            ),
            html.Tbody([
                html.Tr([
                    html.Td(df.iloc[i][col]) for col in df.columns
                ]) for i in range(len(df))
            ])
        ], style={'width': '100%', 'border': '1px solid #ddd'})
    
    def run(self, debug=False):
        """运行仪表板"""
        self.app.run_server(debug=debug, host='0.0.0.0', port=8050)

# 使用示例
# dashboard = ElectionDashboard()
# dashboard.run(debug=True)

6. 监控与告警系统

6.1 系统健康监控

import psutil
import requests
from prometheus_client import Counter, Gauge, Histogram, start_http_server
import logging

class MonitoringSystem:
    def __init__(self):
        # Prometheus指标
        self.data_freshness = Gauge('election_data_freshness_seconds', 
                                   'Time since last data update')
        self.api_latency = Histogram('election_api_latency_seconds', 
                                    'API response time')
        self.data_errors = Counter('election_data_errors_total', 
                                  'Total data errors', ['source', 'error_type'])
        self.system_cpu = Gauge('election_system_cpu_percent', 
                               'System CPU usage')
        self.system_memory = Gauge('election_system_memory_percent', 
                                  'System memory usage')
        
        # 日志配置
        logging.basicConfig(
            level=logging.INFO,
            format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
            handlers=[
                logging.FileHandler('election_monitor.log'),
                logging.StreamHandler()
            ]
        )
        self.logger = logging.getLogger('ElectionMonitor')
        
        # 启动Prometheus指标服务器
        start_http_server(8000)
        self.logger.info("Monitoring server started on port 8000")
    
    def check_data_freshness(self, last_update_timestamp):
        """检查数据新鲜度"""
        current_time = time.time()
        freshness = current_time - last_update_timestamp
        
        self.data_freshness.set(freshness)
        
        if freshness > 300:  # 5分钟无更新
            self.logger.warning(f"Data freshness alert: {freshness:.0f} seconds")
            self.send_alert("数据更新延迟", f"最后更新时间: {freshness:.0f}秒前")
            return False
        
        return True
    
    def monitor_system_resources(self):
        """监控系统资源"""
        cpu_percent = psutil.cpu_percent(interval=1)
        memory = psutil.virtual_memory()
        
        self.system_cpu.set(cpu_percent)
        self.system_memory.set(memory.percent)
        
        # 告警阈值
        if cpu_percent > 80:
            self.logger.warning(f"High CPU usage: {cpu_percent}%")
            self.send_alert("CPU使用率过高", f"当前: {cpu_percent}%")
        
        if memory.percent > 85:
            self.logger.warning(f"High memory usage: {memory.percent}%")
            self.send_alert("内存使用率过高", f"当前: {memory.percent}%")
    
    def monitor_api_health(self, api_url):
        """监控API健康状态"""
        try:
            start_time = time.time()
            response = requests.get(api_url, timeout=5)
            latency = time.time() - start_time
            
            self.api_latency.observe(latency)
            
            if response.status_code != 200:
                self.data_errors.labels(
                    source='api', 
                    error_type=f'status_{response.status_code}'
                ).inc()
                self.logger.error(f"API error: {response.status_code}")
                return False
            
            if latency > 2.0:
                self.logger.warning(f"Slow API response: {latency:.2f}s")
            
            return True
            
        except requests.exceptions.RequestException as e:
            self.data_errors.labels(source='api', error_type='connection').inc()
            self.logger.error(f"API connection error: {e}")
            return False
    
    def send_alert(self, title, message):
        """发送告警(模拟)"""
        # 实际应用中可以集成Slack、邮件、短信等
        alert_data = {
            'title': title,
            'message': message,
            'timestamp': datetime.utcnow().isoformat(),
            'severity': 'WARNING'
        }
        
        # 模拟发送到Webhook
        print(f"ALERT: {title} - {message}")
        
        # 可以集成到实际的告警系统
        # requests.post('https://hooks.slack.com/services/...', json=alert_data)
    
    def run_health_checks(self, api_url, last_update_timestamp):
        """运行完整健康检查"""
        self.monitor_system_resources()
        
        data_ok = self.check_data_freshness(last_update_timestamp)
        api_ok = self.monitor_api_health(api_url)
        
        return data_ok and api_ok

# 使用示例
# monitor = MonitoringSystem()
# while True:
#     monitor.run_health_checks('http://localhost:8050', time.time())
#     time.sleep(60)

7. 完整部署架构

7.1 Docker Compose 配置

# docker-compose.yml
version: '3.8'

services:
  # 数据库服务
  postgres:
    image: postgres:14
    environment:
      POSTGRES_DB: election_db
      POSTGRES_USER: user
      POSTGRES_PASSWORD: pass
    ports:
      - "5432:5432"
    volumes:
      - postgres_data:/var/lib/postgresql/data
  
  # 时序数据库
  influxdb:
    image: influxdb:2.0
    ports:
      - "8086:8086"
    environment:
      DOCKER_INFLUXDB_INIT_MODE: setup
      DOCKER_INFLUXDB_INIT_USERNAME: admin
      DOCKER_INFLUXDB_INIT_PASSWORD: admin123
      DOCKER_INFLUXDB_INIT_ORG: election_org
      DOCKER_INFLUXDB_INIT_BUCKET: election_bucket
    volumes:
      - influx_data:/var/lib/influxdb2
  
  # Redis缓存
  redis:
    image: redis:7-alpine
    ports:
      - "6379:6379"
    command: redis-server --appendonly yes
    volumes:
      - redis_data:/data
  
  # Kafka消息队列
  zookeeper:
    image: confluentinc/cp-zookeeper:7.0.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
  
  kafka:
    image: confluentinc/cp-kafka:7.0.0
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
  
  # 数据采集服务
  data-collector:
    build: ./services/collector
    depends_on:
      - kafka
      - redis
    environment:
      KAFKA_BOOTSTRAP: kafka:9092
      REDIS_HOST: redis
    restart: unless-stopped
  
  # 处理服务
  data-processor:
    build: ./services/processor
    depends_on:
      - kafka
      - postgres
      - influxdb
    environment:
      KAFKA_BOOTSTRAP: kafka:9092
      PG_URL: postgresql://user:pass@postgres:5432/election_db
      INFLUX_URL: http://influxdb:8086
    restart: unless-stopped
  
  # API服务
  api:
    build: ./services/api
    depends_on:
      - postgres
      - redis
    ports:
      - "8000:8000"
    environment:
      PG_URL: postgresql://user:pass@postgres:5432/election_db
      REDIS_HOST: redis
    restart: unless-stopped
  
  # 前端仪表板
  dashboard:
    build: ./services/dashboard
    depends_on:
      - api
    ports:
      - "8050:8050"
    environment:
      API_URL: http://api:8000
    restart: unless-stopped
  
  # 监控服务
  prometheus:
    image: prom/prometheus
    ports:
      - "9090:9090"
    volumes:
      - ./monitoring/prometheus.yml:/etc/prometheus/prometheus.yml
  
  grafana:
    image: grafana/grafana
    ports:
      - "3000:3000"
    environment:
      GF_SECURITY_ADMIN_PASSWORD: admin123
    volumes:
      - grafana_data:/var/lib/grafana

volumes:
  postgres_data:
  influx_data:
  redis_data:
  grafana_data:

8. 安全与合规考虑

8.1 数据安全实现

from cryptography.fernet import Fernet
from cryptography.hazmat.primitives import hashes
from cryptography.hazmat.primitives.kdf.pbkdf2 import PBKDF2HMAC
import base64
import os

class DataSecurity:
    def __init__(self, master_key=None):
        """初始化加密系统"""
        if master_key is None:
            master_key = os.urandom(32)
        
        self.master_key = master_key
        
        # 派生加密密钥
        kdf = PBKDF2HMAC(
            algorithm=hashes.SHA256(),
            length=32,
            salt=b'election_salt_2024',
            iterations=100000,
        )
        self.key = base64.urlsafe_b64encode(kdf.derive(master_key))
        self.cipher = Fernet(self.key)
    
    def encrypt_sensitive_data(self, data):
        """加密敏感数据"""
        if isinstance(data, dict):
            data_str = json.dumps(data, sort_keys=True)
        else:
            data_str = str(data)
        
        encrypted = self.cipher.encrypt(data_str.encode())
        return base64.urlsafe_b64encode(encrypted).decode()
    
    def decrypt_sensitive_data(self, encrypted_data):
        """解密数据"""
        try:
            encrypted_bytes = base64.urlsafe_b64decode(encrypted_data.encode())
            decrypted_bytes = self.cipher.decrypt(encrypted_bytes)
            return json.loads(decrypted_bytes.decode())
        except Exception as e:
            print(f"Decryption error: {e}")
            return None
    
    def hash_voter_id(self, voter_id):
        """对选民ID进行单向哈希"""
        digest = hashes.Hash(hashes.SHA256())
        digest.update(voter_id.encode())
        digest.update(b'election_2024_salt')
        return digest.finalize().hex()
    
    def audit_log(self, action, user, data_hash):
        """记录审计日志"""
        log_entry = {
            'timestamp': datetime.utcnow().isoformat(),
            'action': action,
            'user': user,
            'data_hash': data_hash,
            'ip_address': '127.0.0.1'  # 实际获取真实IP
        }
        
        # 写入不可篡改的日志存储
        with open('audit_log.jsonl', 'a') as f:
            f.write(json.dumps(log_entry) + '\n')

# 使用示例
# security = DataSecurity()
# encrypted = security.encrypt_sensitive_data({'voter_id': '12345', 'vote': 'candidate_A'})
# hashed_id = security.hash_voter_id('voter_12345')

9. 总结与最佳实践

9.1 关键成功因素

  1. 数据质量优先:多源验证确保准确性
  2. 实时性保障:流处理架构实现秒级更新
  3. 可扩展性:微服务架构支持高并发
  4. 安全性:端到端加密和审计日志
  5. 监控告警:全方位系统健康监控

9.2 性能优化建议

  • 使用Redis缓存热点数据
  • 采用列式存储优化时序查询
  • 实施数据分片策略
  • 使用CDN加速静态资源
  • 实施限流和熔断机制

9.3 合规性检查清单

  • [ ] 数据来源合法性验证
  • [ ] 选民隐私保护措施
  • [ ] 数据保留期限合规
  • [ ] 审计日志完整性
  • [ ] 访问权限最小化原则

通过以上完整的系统架构和实现方案,您可以构建一个专业、可靠、实时的”辉煌大选票房追踪系统”,为选民和分析师提供准确、及时的选举数据服务。