项目背景
某电商运营团队原来的实时大屏是基于「每小时跑一次离线 Spark + 前端轮询接口」实现的,数据延迟普遍在 1 小时以上,大促时运营想看「最近 5 分钟成交」根本看不到。技术债还在于离线表越攒越大,跑批窗口越拖越长,夜里经常跑不完。
本次目标是搭一套真正的实时看板:成交、流量、转化等核心指标延迟控制在 5 秒内,看板支持秒级刷新,且保证数据不重不丢(Exactly-Once)。技术栈选 Flink 做流计算、Kafka 做数据总线、ClickHouse 做查询存储、WebSocket 做前端推送。
实时看板最难的不是「跑起来」,而是「数据准」。延迟可以做小,但丢一条订单数据就是事故,Exactly-Once 是底线。
架构设计
整体走 Lambda 的简化版——只用流链路,离线仅做兜底对账。埋点经网关写入 Kafka,Flink 消费后做窗口聚合与去重,结果双写 ClickHouse(实时查询)与 Kafka(下游二次消费),前端通过 WebSocket 订阅聚合结果。
埋点 → 网关 → Kafka(原始 Topic) → Flink(清洗/去重/窗口聚合)
├→ ClickHouse(看板查询,物化视图加速)
└→ Kafka(聚合 Topic) → WebSocket 推前端
关键设计:Flink 用 Checkpoint + 两阶段提交保证 Exactly-Once;ClickHouse 用 ReplacingMergeTree + 版本号去重,配合物化视图预聚合,把查询压到亚秒级;前端不直接查 ClickHouse,而是订阅 WebSocket,由后端定时推增量。
核心实现
1)Flink 窗口聚合:按「1 分钟滚动窗口 + 维度(省份+类目)」聚合成交额,用 ProcessWindowFunction 输出结果,并通过 Checkpointed 的 Offset 提交保证 Exactly-Once。
// OrderAggregJob.java —— 1 分钟窗口聚合
DataStream<OrderEvent> orders = env
.addSource(kafkaSource("orders")) // 精确一次的 KafkaSource
.keyBy(o -> o.getProvince() + "|" + o.getCategory());
orders
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.process(new OrderAggFunction()) // 窗口内聚合 GMV/订单数
.addSink(clickHouseSink()); // 两阶段提交 Sink
// ProcessWindowFunction:窗口结束时输出
public class OrderAggFunction extends ProcessWindowFunction<
OrderEvent, AggResult, String, TimeWindow> {
@Override
public void process(String key, Context ctx,
Iterable<OrderEvent> in, Collector<AggResult> out) {
long gmv = 0; int cnt = 0;
for (OrderEvent o : in) { gmv += o.getAmount(); cnt++; }
out.collect(new AggResult(key, ctx.window().getEnd(), gmv, cnt));
}
}
2)ClickHouse 表设计与查询:用 ReplacingMergeTree 按 (dim, window_end) 去重,版本号取 event_time,配合物化视图把分钟级数据预聚合到小时级,看板查小时图直接命中物化视图。
-- ClickHouse 建表:去重 + 物化视图预聚合
CREATE TABLE order_agg_min (
dim String, -- province|category
window_end DateTime,
gmv UInt64,
order_cnt UInt32,
event_time DateTime,
version UInt64
) ENGINE = ReplacingMergeTree(version)
PARTITION BY toYYYYMMDD(window_end)
ORDER BY (dim, window_end);
-- 物化视图:分钟 → 小时 自动预聚合
CREATE MATERIALIZED VIEW order_agg_hour
ENGINE = SummingMergeTree()
ORDER BY (dim, hour)
AS SELECT dim, toStartOfHour(window_end) AS hour,
sum(gmv) AS gmv, sum(order_cnt) AS order_cnt
FROM order_agg_min GROUP BY dim, hour;
3)WebSocket 增量推送:后端每 3 秒查询 ClickHouse 最近窗口的增量,推给前端订阅者,前端拿到后局部更新图表,避免全量重渲染。
- 数据延迟:1 小时 → 4 秒(端到端)
- 看板查询 P95:3.2s → 280ms(物化视图命中)
- 数据准确性:Exactly-Once,对账差异 0 条
- 看板并发:支撑 500+ 运营同时在线
技术难点
难点一:Exactly-Once 落地。Flink 的 Exactly-Once 需要 Source 与 Sink 都支持两阶段提交。Kafka Source 已内置,但 ClickHouse Sink 默认是 at-least-once。我们的方案是「Sink 先写带版本号,ClickHouse 用 ReplacingMergeTree 去重」——即使重复写入也会被合并,等价于 Exactly-Once 效果,但要保证 version 单调,用 CheckpointId 充当版本号。
难点二:数据倾斜。某头部类目的订单量是其他类目的 20 倍,直接 keyBy 会导致某个 TaskManager 处理不过来背压。解决方式是「Key 加盐」——把热 key 拆成 category_N(N=0~9)十个子 key 分散到不同 slot,窗口聚合后再做一次二次合并。
// 倾斜处理:热 key 加盐打散,聚合后二次合并
// 第一次 keyBy:加盐打散
orders
.keyBy(o -> o.getCategory() + "_"
+ (o.hashCode() % 10))
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.process(new FirstStageAgg()) // 输出 (category_N, gmv, cnt)
// 第二次 keyBy:按原 category 合并
.keyBy(r -> r.getRealCategory())
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.process(new SecondStageAgg()); // 输出最终聚合
难点三: watermark 与乱序。埋点上报存在网络延迟,事件时间乱序严重。我们把 watermark 设为「最大事件时间 - 30 秒」允许迟到,并对迟到数据走侧路输出单独补算,避免丢数据又能保证窗口及时关闭。
踩坑复盘
坑 1:Checkpoint 间隔过短导致吞吐掉。最初设 10 秒一次 Checkpoint,发现每次 Checkpoint 时吞吐骤降。调到 60 秒后吞吐恢复,但延迟略增。最终折中 30 秒 + 增量 Checkpoint,平衡了 Exactly-Once 与吞吐。
坑 2:ClickHouse 大查询把内存打爆。看板查询没做时间范围限制,运营选了「全部时间」直接 OOM。加了「单查询最多扫描 7 天 + 强制命中分区」的护栏后稳定。教训:ClickHouse 查询必须有强制时间范围与扫描量限制。
坑 3:物化视图没同步更新。改了底表字段但忘了重建物化视图,看板数据突然少了一列。规范了「改表必须连带重建物化视图」的发布流程并加 CI 校验后杜绝。教训:物化视图是隐式依赖,要纳入变更检查。
实时数仓的三道命门:Exactly-Once 保证准、Key 加盐保证不倾斜、物化视图保证查得快——少一道都会在峰值翻车。
项目总结
大屏上线后数据延迟从 1 小时压到 4 秒,大促期间运营第一次能「看着实时数据调投放」。整套链路在双 11 峰值 12 万订单/分钟下稳定运行,对账零差异。项目最大的经验是:实时不是堆组件,而是把 Exactly-Once、倾斜、查询加速三件事各自做扎实,再靠压测与对账反复验证。
- 数据延迟:1 小时 → 4 秒
- 查询 P95:3.2s → 280ms
- 峰值吞吐:12 万订单/分钟稳定
- 数据准确性:对账差异 0 条