我要提问
PROJECT CASE / 004

实时数据看板大数据实战

Flink + ClickHouse + Kafka,从埋点采集到秒级可视化的全链路实时数仓。

实时数据看板大数据实战:Flink + ClickHouse 秒级可视化全链路

实时数据看板大数据实战项目示意图

项目背景

某电商运营团队原来的实时大屏是基于「每小时跑一次离线 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 条
返回项目列表