返回项目列表

基于 Flink 的电商实时交易指标平台

免费
Flink,实时数仓,ClickHouse,数据治理

项目描述

一句话

基于 Kafka + Flink + ClickHouse 构建电商实时交易指标链路,实现秒级 GMV、订单数、UV、支付成功率计算及支付失败率风控告警,覆盖乱序处理、状态去重与 Checkpoint Exactly-Once。

业务背景

电商经营侧需要:

  • 支付成功后 3~5 秒内 看到 GMV、订单量、支付成功率
  • 按类目实时下钻盯盘
  • 支付失败率短时突增要告警

离线 Hive/Spark T+1 无法满足。本项目建设近实时(秒级)指标平台,指标口径与离线对齐,便于对账与讲解「流批一致性」。

技术栈
层级技术讲述口径
消息队列Kafka 3.x(KRaft)分区、Consumer Group、Offset
流计算Flink 1.18 DataStream + SQL / Java 11Watermark、窗口、状态、Checkpoint
OLAP / 缓存ClickHouse、Redis分钟预聚合、热点、幂等表引擎
查询展示FastAPI、Grafana下钻 API、延迟可观测
模拟数据Python Generator故意注入乱序 / 重复 / 脏数据
系统架构
text
1 Kafka Flink
2 key=
3 page_view user_id trade_events OdsCleanJob +order ClickHouse ODS
4 place_order (6 partitions) WindowMetricsJob dws_trade_1min
5 pay_success RiskAlertJob Redis
6 // FastAPI/Grafana
7
分层对齐离线心智
实时层对应离线本项目落点
ODS贴源Kafka 原始 Topic + CH ods_trade_event
DWD明细清洗去重后的支付 / 浏览明细
DWS汇总1 分钟滚动窗口 dws_trade_1min
ADS应用风控告警 + Redis / API
核心指标口径
指标定义窗口关键注意
GMVpay_success.amount 求和1min 滚动order_id 去重;用 event_time
订单数去重后成功订单数1min 滚动与 GMV 共用去重语义
UVpage_view 独立 user_id1min 滚动演示精确 Set;生产可用 HLL
支付成功率success / (success+fail)5min 滑动 / 1min 步长分母为 0 跳过
转化率(进阶)下单 UV / 浏览 UV会话窗或双流 Join需跨事件对齐,见实现章节
项目亮点(简历可写)
  1. 有界乱序 Watermark(20s)+ idleness + 迟到侧输出
  2. Keyed State 按订单去重 + State TTL,区分「消息重复」与「业务成功一次」
  3. 滚动 / 滑动窗口增量聚合,避免窗口存全量导致 OOM
  4. Checkpoint Exactly-Once;写 CH 采用 at-least-once + 表引擎幂等(说清折中)
  5. 脏数据 / 重复 / 迟到三类侧输出,支撑对账与排障
  6. 梳理实时与离线差异的可解释因素(迟到、日界、TTL、时区)
设计目标
  • 端到端延迟 P95 < 5s(设计目标)
  • 支持 Checkpoint 故障恢复后指标口径可解释
  • 完整可讲解代码骨架 + 面试问答(不必本地真实跑通全套中间件)

项目实现

1. 总体实现路径
text
1Generator + //
2 Kafka trade_eventskey=user_idpartitions=6
3 Flink
4 OdsCleanJoborder_id
5 WindowMetricsJobEvent Time + Watermark + /
6 RiskAlertJob
7 ClickHouseODS + DWS + / Redis
8 FastAPI / Grafana

工程结构(讲述用):

text
1realtime-trade-metrics/
2 generator/event_generator.py
3 flink-jobs/src/main/java/com/rtm/
4 phase2/OdsCleanJob.java
5 phase3/WindowMetricsJob.java
6 phase5/RiskAlertJob.java
7 util/FlinkJobConfig.java
8 sql/clickhouse_ddl.sql
9 sql/flink_sql_metrics.sql
10 api/main.py
11 scripts/chaos_kill_tm.sh

2. 事件模拟器(制造「脏现实」)

职责:按漏斗生成 page_view → add_cart → place_order → pay_success/fail,并故意注入问题数据,逼出 Watermark / 去重 / 质量监控能力。

注入比例 / 行为对应考点
乱序\~15%,event_time 回拨 5~30sWatermark、迟到
重复\~5%,同 event_id 再发一次状态去重 / Sink 幂等
脏数据\~1%,缺字段或负金额侧输出 dirty-events
失败尖峰--fail-spike风控滑动窗告警

Kafka 发送时 key = user_id,保证同用户事件进入同一分区、分区内有序。


3. Kafka Topic 设计
TopicKey分区建议用途
trade_eventsuser_id6(= Source 并行度)主事件流
trade_metricscategory / 窗起点3指标输出(也可直写 CH)
risk_alertsrule_id1~3风控告警

要点:

  • 分区数 ≥ Flink Source 并行度,否则部分 subtask 空闲
  • 同 key 有序;跨分区无全局有序
  • 多 Job 使用不同 group.id,互不影响(解耦)

4. Phase2 · OdsCleanJob(清洗 + 去重)
4.1 处理链路
  1. KafkaSourcetrade_events,offset 纳入 Checkpoint
  2. JSON → TradeEvent
  3. 脏数据过滤 → 侧输出 dirty-events
  4. 支付类事件 keyBy(order_id)ValueState 去重(TTL 24h)
  5. 重复 → 侧输出 duplicate-pay-events;主流写 ODS(演示 print / 生产写 CH)
4.2 脏数据规则(示例)
  • 解析失败 / 必填字段缺失(event_idevent_typeuser_idevent_time
  • 支付事件缺少 order_id
  • pay_successamount <= 0
4.3 去重核心代码(面试可讲)
java
1StateTtlConfig ttl = StateTtlConfig
2 .newBuilder(Time.hours(24))
3 .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
4 .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
5 .cleanupFullSnapshot()
6 .build();
7
8ValueStateDescriptor<Boolean> desc =
9 new ValueStateDescriptor<>("order-seen", Types.BOOLEAN);
10desc.enableTimeToLive(ttl);
11seen = getRuntimeContext().getState(desc);
12
13// processElement:
14if (Boolean.TRUE.equals(seen.value())) {
15 ctx.output(DUPLICATE_TAG, value); // 重复侧输出
16 return;
17}
18seen.update(Boolean.TRUE);
19out.collect(value); // 首次放行

为什么必须去重: 业务上一单只成功一次,但至少一次投递、回调重试、Checkpoint 前故障重放,都会产生多条 pay_success。\
Exactly-Once ≠ 业务幂等: EOS 防「同一条输入算两次」;输入里本来就有两条业务重复消息时,仍要按 order_id 去重。

边界(加分): 若同一 order_id 先 fail 后 success,简单 seen(order_id) 可能误伤;生产可改为只对 pay_success 去重,或存终态状态机。


5. Phase3 · WindowMetricsJob(核心)
5.1 时间与水位线
java
1WatermarkStrategy<TradeEvent> wm = WatermarkStrategy
2 .<TradeEvent>forBoundedOutOfOrderness(Duration.ofSeconds(20))
3 .withIdleness(Duration.ofSeconds(30))
4 .withTimestampAssigner((e, ts) -> e.getEventTime());
配置取值含义
时间语义Event Time指标按业务发生时间入窗
乱序 bound20sWM ≈ max(event_time) − 20s
Idleness30s冷分区不拖死全局 WM
Allowed Lateness5s关窗后短暂可更新
更晚数据sideOutputLateDatalate-events 可观测

多并行度水位线取 最小值;调参看 ingest_time - event_time 的 P95/P99。

5.2 滚动 1 分钟:GMV / 订单 / UV
java
1events
2 .keyBy(TradeEvent::getCategory)
3 .window(TumblingEventTimeWindows.of(Time.minutes(1)))
4 .allowedLateness(Time.seconds(5))
5 .sideOutputLateData(LATE_TAG)
6 .aggregate(new TradeAggregate(), new MetricWindowFunction());
  • GMV / 订单: 仅统计 pay_success,金额用「分」的 long
  • UV: 累加器内 Set<user_id>(演示精确);生产大盘改 HLL
  • 增量聚合: AggregateFunction 只存中间结果,避免 List 存全量事件 OOM
5.3 滑动 5 分钟 / 1 分钟步长:支付成功率
java
1events
2 .filter(e -> e.isPaySuccess() || e.isPayFail())
3 .keyBy(TradeEvent::getCategory)
4 .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1)))
5 .aggregate(...);

成功率 = success / (success + fail);分母为 0 则跳过或输出 null。

5.4 累加器逻辑(摘要)
text
1Acc: gmv, orderCount, paySuccessCnt, payFailCnt, Set<userId>
2on pay_success gmv+=amount, orderCount++, paySuccessCnt++
3on pay_fail payFailCnt++
4on page_view users.add(userId)
5merge + set

6. Phase5 · RiskAlertJob(风控)
规则
窗口滑动 5min / 步长 1min
指标failRate = fail / (success + fail)
触发total ≥ 50failRate ≥ 0.30
输出RiskAlert → 告警 Topic / CH ads_risk_alert

滑动窗比滚动窗对「刚跨边界的突增」更敏感。阈值与最小样本量可做成广播配置流动态下发(进阶口述)。


7. 容错:FlinkJobConfig
java
1env.enableCheckpointing(60_000L);
2env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
3env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000L);
4env.getCheckpointConfig().setCheckpointTimeout(120_000L);
5// 演示:HashMapStateBackend;生产:EmbeddedRocksDBStateBackend(增量)
6// 关键算子显式 uid,保证 Savepoint 可恢复

故障演练叙事:记录某分钟 GMV → kill TaskManager → 从最近成功 Checkpoint 恢复 → 与干净重放对比;CH 侧靠幂等表收敛短暂重复。


8. ClickHouse 表设计(幂等兜底)

ODS 明细:

sql
1ENGINE = MergeTree
2PARTITION BY dt
3ORDER BY (event_type, category, event_time, event_id)
4TTL dt + INTERVAL 30 DAY;

分钟指标(Flink 窗口输出):

sql
1ENGINE = ReplacingMergeTree(version)
2PARTITION BY toDate(window_start)
3ORDER BY (window_start, category);
4-- 业务键:(window_start, category);version 可用写入时间或 checkpoint id

查询演示可用 FINAL;生产注意 merge 成本,也可查询侧聚合去重。

风控告警表: 同样 ReplacingMergeTree(version),按 rule_id + window_end + category 排序。


9. 服务层:Redis + API
接口行为
GET /metrics/realtime先读 Redis 当前分钟大盘,未命中回源 CH
GET /metrics/category按类目下钻近 N 分钟分钟表
GET /alerts/recent最近风控告警

Redis Key 示例:rtm:dashboard:{yyyyMMddHHmm},TTL 约 120s。


sql
1WATERMARK FOR event_time AS event_time - INTERVAL '20' SECOND;
2
3-- 滚动 1 分钟 GMV
4SELECT window_start, window_end, category,
5 SUM(amount) AS gmv, COUNT(*) AS order_count
6FROM TABLE(TUMBLE(TABLE trade_events_src, DESCRIPTOR(event_time), INTERVAL '1' MINUTE))
7WHERE event_type = 'pay_success'
8GROUP BY window_start, window_end, category;
9
10-- 滑动:HOP(table, time, slide, size) 注意 slide 在前

选型: DataStream 适合侧输出、复杂去重 TTL;SQL 迭代快、易与数仓对齐。生产常见 SQL 为主、DataStream/UDF 补复杂逻辑。


11. 端到端延迟预算
text
1Kafka + Kafka + Watermark (bound)
2 + Flink + Sink + /
3 P95 < 5s

削峰填谷:常态仍秒级;尖峰短暂 Lag 换稳定;死守 SLA 靠扩容而非只靠堆积。


12. 流批对账(可解释差异)
差异来源说明
迟到超 bound实时进侧输出;离线批可算全量
日界 / 时区分钟窗 vs 业务日闭区间
去重 TTL状态过期 vs 离线全量 distinct
Sink 幂等合并前CH 可能短暂重复行

目标不是永远 0 差异,而是 差异可解释、可度量、有阈值告警


13. 转化率进阶(Phase6 口述)

CVR = 下单 UV / 浏览 UV 不能两个独立窗口直接相除。

  • 会话窗: keyBy(user) + Session Window(gap 如 30min),会话内同时有 view 与 order 算转化
  • Interval Join: 浏览流 ⋈ 订单流,下单落在浏览后 N 分钟内算转化

对齐后再按浏览时间进入大盘窗口。Join/会话必须设状态 TTL,防止状态爆炸。


14. 面试讲述节奏(5~8 分钟)
  1. 背景:T+1 不够,要秒级看板与风控
  2. 架构:Kafka → Flink 三 Job → CH/Redis
  3. 难点 1:乱序与 Watermark / 迟到侧输出
  4. 难点 2:order_id 去重 + Checkpoint;CH 幂等折中
  5. 难点 3:UV 精确 vs HLL;流批对账
  6. 收尾:生产可补 CDC 维表、倾斜治理、SLA 大盘

问题讲解

12

评论 (0)

登录后参与评论。

还没有评论,来做第一个。

登录后可以选中正文添加批注(仅自己可见)。

基于 Flink 的电商实时交易指标平台