Skip to content

4.10 数据中台

一句话总结:数据中台是业务中台的核心——没有数据,决策就是拍脑袋。

📊 学习进度

  • 状态:⬜ 未开始
  • 上次实时更新:2026-07-03
  • 预计时长:3-4 小时
  • 已完成:0/3 个模块
  • 在整体流程中的位置:预测市场实战·第 7 阶段

📍 本章定位

  • 服务方案:方案 3(核心 90%)
  • 学习方式:🔥 推荐
  • 在流程中的作用:数据采集、分析、展示
  • 核心知识点:数据管道、分析工具、可视化
  • 预计时长:3-4 小时
  • 完成后能做什么:能搭建数据中台

人机分工

环节谁做重要度说明
数据需求定义🧑 人⭐⭐⭐⭐⭐决定采集什么数据
数据管道搭建🤖 AI⭐⭐⭐⭐AI 生成代码
数据分析🤖 AI + 🧑 人⭐⭐⭐⭐AI 分析,人解读
可视化设计🧑 人 + 🤖 AI⭐⭐⭐⭐AI 生成初稿

1. 数据中台概述

1.1 什么是数据中台

数据中台是连接数据源和业务系统的中间层,负责数据的采集、清洗、存储、分析和展示。在预测市场中,数据中台是"大脑",支撑交易决策、做市策略、风控监控和产品优化。

1.2 为什么预测市场需要数据中台

没有数据中台有数据中台
决策靠直觉决策靠数据
问题发现滞后实时告警
无法优化策略数据驱动优化
用户流失不知原因用户行为分析
做市靠经验做市靠算法

1.3 数据中台架构


2. 数据采集

2.1 数据源分类

数据源类型采集方式频率说明
链上交易结构化事件监听实时合约事件
订单簿结构化日志采集实时交易记录
用户行为半结构化埋点实时页面访问、点击
市场行情结构化API 拉取分钟级价格、交易量
新闻舆情非结构化爬虫小时级新闻、社交媒体

2.2 链上数据监听

typescript
import { ethers } from 'ethers';

/**
 * 链上事件监听器
 * 监听预测市场合约事件,写入数据管道
 */
class ChainEventListener {
  private provider: ethers.Provider;
  private contracts: Map<string, ethers.Contract>;
  
  constructor(rpcUrl: string) {
    this.provider = new ethers.JsonRpcProvider(rpcUrl);
    this.contracts = new Map();
  }
  
  /**
   * 监听事件创建
   */
  async listenEventCreated(callback: (event: EventCreated) => void): Promise<void> {
    const contract = this.getContract('EventFactory');
    
    contract.on('EventCreated', (eventId, eventAddress, title, event) => {
      callback({
        eventId: eventId.toString(),
        eventAddress,
        title,
        blockNumber: event.blockNumber,
        transactionHash: event.transactionHash,
        timestamp: Date.now()
      });
    });
  }
  
  /**
   * 监听交易撮合
   */
  async listenTrade(callback: (trade: TradeEvent) => void): Promise<void> {
    const contract = this.getContract('OrderBook');
    
    contract.on('Trade', (eventId, trader, isBuy, isYes, amount, price, fee, event) => {
      callback({
        eventId: eventId.toString(),
        trader,
        isBuy,
        isYes,
        amount: ethers.formatUnits(amount, 6),
        price: ethers.formatUnits(price, 4),
        fee: ethers.formatUnits(fee, 6),
        blockNumber: event.blockNumber,
        transactionHash: event.transactionHash,
        timestamp: Date.now()
      });
    });
  }
  
  /**
   * 监听结算结果
   */
  async listenResolution(callback: (resolution: ResolutionEvent) => void): Promise<void> {
    const contract = this.getContract('OracleManager');
    
    contract.on('ResultFinalized', (eventId, outcome, event) => {
      callback({
        eventId: eventId.toString(),
        outcome,
        blockNumber: event.blockNumber,
        transactionHash: event.transactionHash,
        timestamp: Date.now()
      });
    });
  }
}

2.3 用户行为埋点

typescript
/**
 * 用户行为埋点 SDK
 */
class Analytics {
  private endpoint: string;
  private buffer: Event[] = [];
  private flushInterval: number = 5000; // 5 秒批量发送
  
  constructor(endpoint: string) {
    this.endpoint = endpoint;
    this.startFlushTimer();
  }
  
  /**
   * 追踪页面访问
   */
  trackPageView(page: string, properties?: Record<string, any>): void {
    this.track('page_view', { page, ...properties });
  }
  
  /**
   * 追踪交易行为
   */
  trackTrade(eventId: string, side: string, amount: number, price: number): void {
    this.track('trade', { eventId, side, amount, price });
  }
  
  /**
   * 追踪事件创建
   */
  trackEventCreation(eventId: string, category: string): void {
    this.track('event_created', { eventId, category });
  }
  
  /**
   * 通用追踪方法
   */
  private track(eventName: string, properties: Record<string, any>): void {
    this.buffer.push({
      event: eventName,
      properties,
      timestamp: Date.now(),
      userId: this.getUserId(),
      sessionId: this.getSessionId()
    });
  }
  
  /**
   * 批量发送
   */
  private async flush(): Promise<void> {
    if (this.buffer.length === 0) return;
    
    const events = [...this.buffer];
    this.buffer = [];
    
    try {
      await fetch(this.endpoint, {
        method: 'POST',
        headers: { 'Content-Type': 'application/json' },
        body: JSON.stringify({ events })
      });
    } catch (error) {
      // 发送失败,放回缓冲区
      this.buffer.unshift(...events);
    }
  }
  
  private startFlushTimer(): void {
    setInterval(() => this.flush(), this.flushInterval);
  }
}

3. 数据存储

3.1 存储方案选型

存储用途数据量查询特点推荐
PostgreSQL业务数据中等事务、关联查询主库
Redis缓存、实时快速读写缓存层
ClickHouse分析数据聚合查询分析库
S3日志、备份超大批量读取归档

3.2 ClickHouse 分析表

sql
-- 交易分析表
CREATE TABLE trades_analytics (
    event_id String,
    trader String,
    side Enum8('YES' = 1, 'NO' = 2),
    price Decimal(5, 4),
    amount Decimal(18, 2),
    fee Decimal(18, 2),
    timestamp DateTime,
    date Date DEFAULT toDate(timestamp)
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(date)
ORDER BY (event_id, timestamp);

-- 用户行为表
CREATE TABLE user_events (
    user_id String,
    event_name String,
    properties String,
    timestamp DateTime,
    date Date DEFAULT toDate(timestamp)
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(date)
ORDER BY (user_id, event_name, timestamp);

-- 市场快照表
CREATE TABLE market_snapshots (
    event_id String,
    yes_price Decimal(5, 4),
    no_price Decimal(5, 4),
    volume_24h Decimal(18, 2),
    depth_yes Decimal(18, 2),
    depth_no Decimal(18, 2),
    timestamp DateTime
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(timestamp)
ORDER BY (event_id, timestamp);

3.3 数据流图


4. 核心数据指标

4.1 指标体系

指标类别指标计算公式说明
交易24h 交易量SUM(amount) WHERE timestamp > now() - 24h市场活跃度
成交笔数COUNT(trades)交易频率
平均成交价AVG(price)价格发现效率
用户DAUCOUNT(DISTINCT user_id) WHERE date = today日活跃用户
MAUCOUNT(DISTINCT user_id) WHERE month = current_month月活跃用户
新增用户COUNT(users) WHERE created_at > now() - 24h增长速度
留存率DAU_today / DAU_7days_ago用户粘性
做市盘口深度SUM(depth)流动性质量
价差ask_price - bid_price做市效率
库存风险ABS(position) / max_position做市风险
市场事件数量COUNT(events) WHERE status = 'active'供给能力
结算准确率resolved_correctly / total_resolved预言机质量

4.2 关键指标看板


5. 数据分析

5.1 交易分析

sql
-- 每日交易量趋势
SELECT 
    date,
    COUNT(*) as trade_count,
    SUM(amount) as total_volume,
    AVG(price) as avg_price,
    SUM(fee) as total_fee
FROM trades_analytics
WHERE date >= today() - 30
GROUP BY date
ORDER BY date;

-- 热门事件 TOP 10
SELECT 
    event_id,
    COUNT(*) as trade_count,
    SUM(amount) as total_volume,
    AVG(price) as avg_price
FROM trades_analytics
WHERE date >= today() - 7
GROUP BY event_id
ORDER BY total_volume DESC
LIMIT 10;

-- 用户交易分布
SELECT 
    CASE 
        WHEN total_volume < 100 THEN '< $100'
        WHEN total_volume < 1000 THEN '$100-$1k'
        WHEN total_volume < 10000 THEN '$1k-$10k'
        ELSE '> $10k'
    END as volume_bucket,
    COUNT(*) as user_count,
    SUM(total_volume) as total_volume
FROM (
    SELECT trader, SUM(amount) as total_volume
    FROM trades_analytics
    WHERE date >= today() - 30
    GROUP BY trader
)
GROUP BY volume_bucket
ORDER BY total_volume DESC;

5.2 用户分析

sql
-- 用户留存分析
WITH 
    cohort AS (
        SELECT user_id, MIN(date) as cohort_date
        FROM user_events
        WHERE event_name = 'page_view'
        GROUP BY user_id
    ),
    retention AS (
        SELECT 
            c.cohort_date,
            DATEDIFF(day, c.cohort_date, e.date) as day_number,
            COUNT(DISTINCT e.user_id) as retained_users
        FROM cohort c
        JOIN user_events e ON c.user_id = e.user_id
        WHERE e.date >= c.cohort_date
        GROUP BY c.cohort_date, day_number
    )
SELECT 
    cohort_date,
    day_number,
    retained_users,
    retained_users * 100.0 / FIRST_VALUE(retained_users) OVER (
        PARTITION BY cohort_date ORDER BY day_number
    ) as retention_rate
FROM retention
WHERE day_number IN (1, 3, 7, 14, 30)
ORDER BY cohort_date, day_number;

5.3 市场分析

sql
-- 市场效率分析
SELECT 
    event_id,
    AVG(yes_price) as avg_yes_price,
    STDDEV(yes_price) as price_volatility,
    AVG(depth_yes + depth_no) as avg_depth,
    AVG(yes_price + no_price - 1) as avg_spread
FROM market_snapshots
WHERE timestamp >= now() - INTERVAL 24 HOUR
GROUP BY event_id
HAVING COUNT(*) > 100
ORDER BY avg_volume DESC
LIMIT 20;

6. 数据可视化

6.1 Dashboard 设计

6.2 Grafana 配置示例

json
{
  "dashboard": {
    "title": "Prediction Market Dashboard",
    "panels": [
      {
        "title": "24h Trading Volume",
        "type": "timeseries",
        "targets": [{
          "query": "SELECT date, SUM(amount) FROM trades_analytics WHERE date >= today() - 7 GROUP BY date"
        }]
      },
      {
        "title": "Active Users",
        "type": "stat",
        "targets": [{
          "query": "SELECT COUNT(DISTINCT user_id) FROM user_events WHERE date = today()"
        }]
      },
      {
        "title": "Top Events",
        "type": "table",
        "targets": [{
          "query": "SELECT event_id, SUM(amount) as volume FROM trades_analytics WHERE date >= today() - 1 GROUP BY event_id ORDER BY volume DESC LIMIT 10"
        }]
      }
    ]
  }
}

7. 数据 API

7.1 API 接口

接口方法路径说明
获取市场统计GET/api/stats/market市场整体数据
获取事件统计GET/api/stats/events/:id单个事件数据
获取用户统计GET/api/stats/users/:address用户交易数据
获取排行榜GET/api/leaderboard交易量排行
获取价格历史GET/api/prices/:eventIdK 线数据

7.2 API 响应示例

typescript
// 市场统计响应
interface MarketStats {
  totalVolume24h: string;      // 24h 总交易量
  totalTrades24h: number;      // 24h 成交笔数
  activeEvents: number;        // 活跃事件数
  activeUsers24h: number;      // 24h 活跃用户
  averageSpread: number;       // 平均价差
  totalLiquidity: string;      // 总流动性
}

// 价格历史响应
interface PriceHistory {
  eventId: string;
  interval: '1m' | '5m' | '1h' | '1d';
  data: Array<{
    timestamp: number;
    open: number;
    high: number;
    low: number;
    close: number;
    volume: number;
  }>;
}

8. OPC 数据方案

8.1 轻量数据方案

作为 OPC,不需要复杂的数据中台。推荐轻量方案:

组件推荐方案成本说明
数据库PostgreSQL免费层业务+分析一体
缓存Redis (Upstash)免费层实时数据
分析Metabase免费BI 工具
监控Grafana Cloud免费层系统监控
告警Telegram Bot免费通知

8.2 OPC 数据流程


9. 数据管道进阶

9.1 实时数据管道架构

预测市场需要亚秒级的数据更新,传统的批处理无法满足需求:

Flink 实时聚合示例

typescript
// 使用 Node.js 模拟 Flink 的窗口聚合
class StreamProcessor {
  private windows: Map<string, WindowState> = new Map();
  
  /**
   * 滑动窗口聚合:每 10 秒计算一次交易量
   */
  processTrade(trade: TradeEvent): void {
    const windowKey = `volume_${Math.floor(trade.timestamp / 10000) * 10}`;
    
    if (!this.windows.has(windowKey)) {
      this.windows.set(windowKey, {
        volume: 0,
        count: 0,
        startTime: Math.floor(trade.timestamp / 10000) * 10000,
        endTime: Math.floor(trade.timestamp / 10000) * 10000 + 10000
      });
    }
    
    const window = this.windows.get(windowKey)!;
    window.volume += trade.amount;
    window.count += 1;
    
    // 触发聚合结果输出
    if (trade.timestamp >= window.endTime) {
      this.emitAggregation(window);
      this.windows.delete(windowKey);
    }
  }
  
  /**
   * 异常检测:交易量突增告警
   */
  detectAnomaly(eventId: string, currentVolume: number): boolean {
    const historicalAvg = this.getHistoricalAverage(eventId);
    const threshold = historicalAvg * 3; // 3 倍标准差
    
    if (currentVolume > threshold) {
      this.emitAlert({
        type: 'volume_spike',
        eventId,
        currentVolume,
        threshold,
        severity: 'warning'
      });
      return true;
    }
    return false;
  }
}

9.2 数据质量保障

质量维度检查项工具告警阈值
完整性数据缺失率Great Expectations>1%
准确性链上链下一致性自定义校验不一致即告警
及时性数据延迟监控脚本>30s
一致性跨表数据一致dbt tests不一致即告警

数据质量检查脚本

typescript
class DataQualityChecker {
  /**
   * 检查链上链下数据一致性
   */
  async checkConsistency(eventId: string): Promise<ConsistencyReport> {
    // 链上数据
    const onChainVolume = await this.getOnChainVolume(eventId);
    const onChainTrades = await this.getOnChainTradeCount(eventId);
    
    // 链下数据
    const offChainVolume = await this.db.getVolume(eventId);
    const offChainTrades = await this.db.getTradeCount(eventId);
    
    const volumeDiff = Math.abs(onChainVolume - offChainVolume) / onChainVolume;
    const tradesDiff = Math.abs(onChainTrades - offChainTrades);
    
    return {
      eventId,
      onChainVolume,
      offChainVolume,
      volumeDiff,
      onChainTrades,
      offChainTrades,
      tradesDiff,
      isConsistent: volumeDiff < 0.001 && tradesDiff === 0
    };
  }
}

10. 常见问题

问题原因解决方案
数据不准链上链下不同步双写+校验
查询慢没有索引优化索引+缓存
数据丢失管道故障冗余+监控
存储成本高数据量大分层存储+归档
实时性差批处理延迟流处理+缓存
数据不一致多数据源统一数据模型

10. 下一步

完成数据中台后,进入 阶段 8:测试运维

9.3 实时数据看板技术实现

预测市场需要亚秒级的数据推送,传统轮询无法满足需求。以下是基于 WebSocket 的实时看板实现:

typescript
// 实时数据推送服务
class RealtimeDashboard {
  private wss: WebSocketServer;
  private subscriptions: Map<string, Set<WebSocket>> = new Map();
  
  constructor(port: number) {
    this.wss = new WebSocketServer({ port });
    this.setupConnectionHandler();
    this.startDataStreams();
  }
  
  private setupConnectionHandler(): void {
    this.wss.on('connection', (ws) => {
      ws.on('message', (data) => {
        const msg = JSON.parse(data.toString());
        
        if (msg.type === 'subscribe') {
          this.subscribe(ws, msg.channel, msg.params);
        } else if (msg.type === 'unsubscribe') {
          this.unsubscribe(ws, msg.channel);
        }
      });
      
      ws.on('close', () => {
        this.removeAllSubscriptions(ws);
      });
    });
  }
  
  private subscribe(ws: WebSocket, channel: string, params: any): void {
    const key = `${channel}:${JSON.stringify(params)}`;
    
    if (!this.subscriptions.has(key)) {
      this.subscriptions.set(key, new Set());
    }
    this.subscriptions.get(key)!.add(ws);
  }
  
  private startDataStreams(): void {
    // 订单簿更新流:每 100ms 推送一次
    setInterval(() => {
      for (const [key, clients] of this.subscriptions) {
        if (key.startsWith('orderbook:')) {
          const eventId = JSON.parse(key.split(':')[1]).eventId;
          const snapshot = this.getOrderBookSnapshot(eventId);
          
          for (const client of clients) {
            if (client.readyState === WebSocket.OPEN) {
              client.send(JSON.stringify({
                type: 'orderbook',
                data: snapshot,
                timestamp: Date.now()
              }));
            }
          }
        }
      }
    }, 100);
    
    // 成交流:实时推送
    // 由撮合引擎事件触发
  }
}

实时看板数据流架构

9.4 数据归档与冷热分离

随着数据量增长,需要实施冷热分离策略以控制成本:

数据类型热存储(30 天)温存储(90 天)冷存储(永久)
交易数据PostgreSQLClickHouseS3 Parquet
用户行为RedisClickHouseS3 Parquet
链上事件PostgreSQLClickHouseThe Graph
日志数据LokiS3S3 Glacier
市场快照RedisClickHouseS3

归档脚本示例

typescript
class DataArchiver {
  /**
   * 将过期数据从热存储迁移到冷存储
   */
  async archiveOldData(retentionDays: number = 30): Promise<void> {
    const cutoffDate = new Date();
    cutoffDate.setDate(cutoffDate.getDate() - retentionDays);
    
    // 1. 导出 PostgreSQL 数据到 Parquet
    const trades = await this.pg.query(
      'SELECT * FROM trades WHERE created_at < $1',
      [cutoffDate]
    );
    
    const parquetData = this.convertToParquet(trades.rows);
    
    // 2. 上传到 S3
    await this.s3.putObject({
      Bucket: 'prediction-market-archive',
      Key: `trades/${cutoffDate.toISOString().split('T')[0]}.parquet`,
      Body: parquetData
    });
    
    // 3. 从 PostgreSQL 删除
    await this.pg.query(
      'DELETE FROM trades WHERE created_at < $1',
      [cutoffDate]
    );
    
    console.log(`Archived ${trades.rows.length} trades before ${cutoffDate}`);
  }
}

9.5 数据平台最佳实践

实践说明优先级
数据分层热/温/冷三层存储P0
链上对账每日链上链下数据校验P0
实时告警异常数据即时通知P0
数据归档30 天后归档到 S3P1
数据质量Great Expectations 校验P1

数据质量检查 Prompt(用于 AI 生成数据校验规则):

text
请为预测市场数据平台生成数据质量检查规则,包括:
1. 交易数据:金额>0、价格在 0-1 之间、时间戳有效
2. 用户数据:地址格式正确、交易次数>=0
3. 市场数据:YES+NO 价格=1、深度>=0
4. 链上数据:交易哈希格式、区块号递增

请输出 SQL 格式的检查查询。

数据告警配置示例

typescript
// 数据告警规则
const dataAlertRules = [
  {
    name: '交易量异常',
    condition: 'volume_1h > avg_volume_24h * 3',
    severity: 'warning',
    message: '过去 1 小时交易量超过 24 小时平均值的 3 倍'
  },
  {
    name: '链上链下不一致',
    condition: 'abs(onchain_volume - offchain_volume) / onchain_volume > 0.001',
    severity: 'critical',
    message: '链上链下交易量差异超过 0.1%'
  },
  {
    name: '数据延迟',
    condition: 'now() - latest_data_timestamp > 30',
    severity: 'warning',
    message: '数据更新延迟超过 30 秒'
  }
];

9.6 数据平台成本优化

随着数据量增长,存储和计算成本会快速上升。以下是成本优化策略:

typescript
// 数据生命周期管理
interface DataLifecyclePolicy {
  hotStorage: {
    duration: number;      // 热存储时长(天)
    maxSize: number;       // 最大存储容量
    storage: 'PostgreSQL'; // 存储类型
  };
  warmStorage: {
    duration: number;      // 温存储时长(天)
    storage: 'ClickHouse'; // 存储类型
  };
  coldStorage: {
    duration: number;      // 冷存储时长(天)
    storage: 'S3';         // 存储类型
  };
}

// 默认策略
const defaultPolicy: DataLifecyclePolicy = {
  hotStorage: { duration: 30, maxSize: 100, storage: 'PostgreSQL' },
  warmStorage: { duration: 90, storage: 'ClickHouse' },
  coldStorage: { duration: 365, storage: 'S3' }
};

成本优化清单

优化项说明节省比例
冷热分离过期数据迁移到便宜存储60-80%
数据压缩使用列式存储压缩70-90%
查询优化索引+物化视图50-70%
采样分析大数据集使用采样80-90%
定时任务非实时分析延迟执行30-50%

数据压缩配置示例

sql
-- ClickHouse 压缩配置
CREATE TABLE trades_compressed (
    event_id String,
    price Decimal(5, 4),
    amount Decimal(18, 2),
    timestamp DateTime
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(timestamp)
ORDER BY (event_id, timestamp)
SETTINGS index_granularity = 8192;

-- 使用 ZSTD 压缩
ALTER TABLE trades_compressed MODIFY SETTING compression_codec = 'ZSTD(3)';

参考与延伸

[11] Data Engineering Best Practices(2025)— 数据工程最佳实践

[12] Apache Parquet(2025)— 列式存储格式

[1] ClickHouse(2025)— 列式分析数据库

[2] Metabase(2025)— 开源 BI 工具

[3] Grafana(2025)— 监控可视化平台

[4] Apache Kafka(2025)— 消息队列

[5] dbt(2025)— 数据转换工具

[6] Apache Flink(2025)— 实时流处理框架

[7] Great Expectations(2025)— 数据质量保障工具

[8] Apache Spark(2025)— 大数据处理框架

[9] Snowflake(2025)— 云数据仓库

[10] Data Mesh(2022)— Zhamak Dehghani,数据架构新范式

OPC 超级个体实战指南