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 链上数据监听
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 用户行为埋点
/**
* 用户行为埋点 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 分析表
-- 交易分析表
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) | 价格发现效率 | |
| 用户 | DAU | COUNT(DISTINCT user_id) WHERE date = today | 日活跃用户 |
| MAU | COUNT(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 交易分析
-- 每日交易量趋势
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 用户分析
-- 用户留存分析
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 市场分析
-- 市场效率分析
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 配置示例
{
"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/:eventId | K 线数据 |
7.2 API 响应示例
// 市场统计响应
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 实时聚合示例:
// 使用 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 | 不一致即告警 |
数据质量检查脚本:
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 的实时看板实现:
// 实时数据推送服务
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 天) | 冷存储(永久) |
|---|---|---|---|
| 交易数据 | PostgreSQL | ClickHouse | S3 Parquet |
| 用户行为 | Redis | ClickHouse | S3 Parquet |
| 链上事件 | PostgreSQL | ClickHouse | The Graph |
| 日志数据 | Loki | S3 | S3 Glacier |
| 市场快照 | Redis | ClickHouse | S3 |
归档脚本示例:
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 天后归档到 S3 | P1 |
| 数据质量 | Great Expectations 校验 | P1 |
数据质量检查 Prompt(用于 AI 生成数据校验规则):
请为预测市场数据平台生成数据质量检查规则,包括:
1. 交易数据:金额>0、价格在 0-1 之间、时间戳有效
2. 用户数据:地址格式正确、交易次数>=0
3. 市场数据:YES+NO 价格=1、深度>=0
4. 链上数据:交易哈希格式、区块号递增
请输出 SQL 格式的检查查询。数据告警配置示例:
// 数据告警规则
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 数据平台成本优化
随着数据量增长,存储和计算成本会快速上升。以下是成本优化策略:
// 数据生命周期管理
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% |
数据压缩配置示例:
-- 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,数据架构新范式