基于 Flink 的电商实时交易指标平台
免费项目描述
基于 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 11 | Watermark、窗口、状态、Checkpoint |
| OLAP / 缓存 | ClickHouse、Redis | 分钟预聚合、热点、幂等表引擎 |
| 查询展示 | FastAPI、Grafana | 下钻 API、延迟可观测 |
| 模拟数据 | Python Generator | 故意注入乱序 / 重复 / 脏数据 |
| 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 |
| 指标 | 定义 | 窗口 | 关键注意 |
|---|---|---|---|
| GMV | pay_success.amount 求和 | 1min 滚动 | 按 order_id 去重;用 event_time |
| 订单数 | 去重后成功订单数 | 1min 滚动 | 与 GMV 共用去重语义 |
| UV | page_view 独立 user_id | 1min 滚动 | 演示精确 Set;生产可用 HLL |
| 支付成功率 | success / (success+fail) | 5min 滑动 / 1min 步长 | 分母为 0 跳过 |
| 转化率(进阶) | 下单 UV / 浏览 UV | 会话窗或双流 Join | 需跨事件对齐,见实现章节 |
- 有界乱序 Watermark(20s)+ idleness + 迟到侧输出
- Keyed State 按订单去重 + State TTL,区分「消息重复」与「业务成功一次」
- 滚动 / 滑动窗口增量聚合,避免窗口存全量导致 OOM
- Checkpoint Exactly-Once;写 CH 采用 at-least-once + 表引擎幂等(说清折中)
- 脏数据 / 重复 / 迟到三类侧输出,支撑对账与排障
- 梳理实时与离线差异的可解释因素(迟到、日界、TTL、时区)
- 端到端延迟 P95 < 5s(设计目标)
- 支持 Checkpoint 故障恢复后指标口径可解释
- 完整可讲解代码骨架 + 面试问答(不必本地真实跑通全套中间件)
项目实现
| 1 | Generator(漏斗事件 + 乱序/重复/脏数据) |
| 2 | → Kafka trade_events(key=user_id,partitions=6) |
| 3 | → Flink |
| 4 | ├─ OdsCleanJob:解析、脏数据侧输出、order_id 状态去重 |
| 5 | ├─ WindowMetricsJob:Event Time + Watermark + 滚动/滑动窗口 |
| 6 | └─ RiskAlertJob:滑动窗失败率阈值告警 |
| 7 | → ClickHouse(ODS + 分钟 DWS + 告警)/ Redis 热点 |
| 8 | → FastAPI / Grafana 查询与展示 |
工程结构(讲述用):
| 1 | realtime-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 |
职责:按漏斗生成 page_view → add_cart → place_order → pay_success/fail,并故意注入问题数据,逼出 Watermark / 去重 / 质量监控能力。
| 注入 | 比例 / 行为 | 对应考点 |
|---|---|---|
| 乱序 | \~15%,event_time 回拨 5~30s | Watermark、迟到 |
| 重复 | \~5%,同 event_id 再发一次 | 状态去重 / Sink 幂等 |
| 脏数据 | \~1%,缺字段或负金额 | 侧输出 dirty-events |
| 失败尖峰 | --fail-spike | 风控滑动窗告警 |
Kafka 发送时 key = user_id,保证同用户事件进入同一分区、分区内有序。
| Topic | Key | 分区建议 | 用途 |
|---|---|---|---|
trade_events | user_id | 6(= Source 并行度) | 主事件流 |
trade_metrics | category / 窗起点 | 3 | 指标输出(也可直写 CH) |
risk_alerts | rule_id | 1~3 | 风控告警 |
要点:
- 分区数 ≥ Flink Source 并行度,否则部分 subtask 空闲
- 同 key 有序;跨分区无全局有序
- 多 Job 使用不同
group.id,互不影响(解耦)
KafkaSource读trade_events,offset 纳入 Checkpoint- JSON →
TradeEvent - 脏数据过滤 → 侧输出
dirty-events - 支付类事件
keyBy(order_id)→ValueState去重(TTL 24h) - 重复 → 侧输出
duplicate-pay-events;主流写 ODS(演示 print / 生产写 CH)
- 解析失败 / 必填字段缺失(
event_id、event_type、user_id、event_time) - 支付事件缺少
order_id pay_success且amount <= 0
| 1 | StateTtlConfig ttl = StateTtlConfig |
| 2 | .newBuilder(Time.hours(24)) |
| 3 | .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) |
| 4 | .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) |
| 5 | .cleanupFullSnapshot() |
| 6 | .build(); |
| 7 | |
| 8 | ValueStateDescriptor<Boolean> desc = |
| 9 | new ValueStateDescriptor<>("order-seen", Types.BOOLEAN); |
| 10 | desc.enableTimeToLive(ttl); |
| 11 | seen = getRuntimeContext().getState(desc); |
| 12 | |
| 13 | // processElement: |
| 14 | if (Boolean.TRUE.equals(seen.value())) { |
| 15 | ctx.output(DUPLICATE_TAG, value); // 重复侧输出 |
| 16 | return; |
| 17 | } |
| 18 | seen.update(Boolean.TRUE); |
| 19 | out.collect(value); // 首次放行 |
为什么必须去重: 业务上一单只成功一次,但至少一次投递、回调重试、Checkpoint 前故障重放,都会产生多条 pay_success。\
Exactly-Once ≠ 业务幂等: EOS 防「同一条输入算两次」;输入里本来就有两条业务重复消息时,仍要按 order_id 去重。
边界(加分): 若同一 order_id 先 fail 后 success,简单 seen(order_id) 可能误伤;生产可改为只对 pay_success 去重,或存终态状态机。
| 1 | WatermarkStrategy<TradeEvent> wm = WatermarkStrategy |
| 2 | .<TradeEvent>forBoundedOutOfOrderness(Duration.ofSeconds(20)) |
| 3 | .withIdleness(Duration.ofSeconds(30)) |
| 4 | .withTimestampAssigner((e, ts) -> e.getEventTime()); |
| 配置 | 取值 | 含义 |
|---|---|---|
| 时间语义 | Event Time | 指标按业务发生时间入窗 |
| 乱序 bound | 20s | WM ≈ max(event_time) − 20s |
| Idleness | 30s | 冷分区不拖死全局 WM |
| Allowed Lateness | 5s | 关窗后短暂可更新 |
| 更晚数据 | sideOutputLateData | late-events 可观测 |
多并行度水位线取 最小值;调参看 ingest_time - event_time 的 P95/P99。
| 1 | events |
| 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
| 1 | events |
| 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。
| 1 | Acc: gmv, orderCount, paySuccessCnt, payFailCnt, Set<userId> |
| 2 | on pay_success → gmv+=amount, orderCount++, paySuccessCnt++ |
| 3 | on pay_fail → payFailCnt++ |
| 4 | on page_view → users.add(userId) |
| 5 | merge → 字段相加 + set 合并 |
| 项 | 规则 |
|---|---|
| 窗口 | 滑动 5min / 步长 1min |
| 指标 | failRate = fail / (success + fail) |
| 触发 | total ≥ 50 且 failRate ≥ 0.30 |
| 输出 | RiskAlert → 告警 Topic / CH ads_risk_alert |
滑动窗比滚动窗对「刚跨边界的突增」更敏感。阈值与最小样本量可做成广播配置流动态下发(进阶口述)。
| 1 | env.enableCheckpointing(60_000L); |
| 2 | env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); |
| 3 | env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000L); |
| 4 | env.getCheckpointConfig().setCheckpointTimeout(120_000L); |
| 5 | // 演示:HashMapStateBackend;生产:EmbeddedRocksDBStateBackend(增量) |
| 6 | // 关键算子显式 uid,保证 Savepoint 可恢复 |
故障演练叙事:记录某分钟 GMV → kill TaskManager → 从最近成功 Checkpoint 恢复 → 与干净重放对比;CH 侧靠幂等表收敛短暂重复。
ODS 明细:
| 1 | ENGINE = MergeTree |
| 2 | PARTITION BY dt |
| 3 | ORDER BY (event_type, category, event_time, event_id) |
| 4 | TTL dt + INTERVAL 30 DAY; |
分钟指标(Flink 窗口输出):
| 1 | ENGINE = ReplacingMergeTree(version) |
| 2 | PARTITION BY toDate(window_start) |
| 3 | ORDER BY (window_start, category); |
| 4 | -- 业务键:(window_start, category);version 可用写入时间或 checkpoint id |
查询演示可用 FINAL;生产注意 merge 成本,也可查询侧聚合去重。
风控告警表: 同样 ReplacingMergeTree(version),按 rule_id + window_end + category 排序。
| 接口 | 行为 |
|---|---|
GET /metrics/realtime | 先读 Redis 当前分钟大盘,未命中回源 CH |
GET /metrics/category | 按类目下钻近 N 分钟分钟表 |
GET /alerts/recent | 最近风控告警 |
Redis Key 示例:rtm:dashboard:{yyyyMMddHHmm},TTL 约 120s。
| 1 | WATERMARK FOR event_time AS event_time - INTERVAL '20' SECOND; |
| 2 | |
| 3 | -- 滚动 1 分钟 GMV |
| 4 | SELECT window_start, window_end, category, |
| 5 | SUM(amount) AS gmv, COUNT(*) AS order_count |
| 6 | FROM TABLE(TUMBLE(TABLE trade_events_src, DESCRIPTOR(event_time), INTERVAL '1' MINUTE)) |
| 7 | WHERE event_type = 'pay_success' |
| 8 | GROUP BY window_start, window_end, category; |
| 9 | |
| 10 | -- 滑动:HOP(table, time, slide, size) 注意 slide 在前 |
选型: DataStream 适合侧输出、复杂去重 TTL;SQL 迭代快、易与数仓对齐。生产常见 SQL 为主、DataStream/UDF 补复杂逻辑。
| 1 | 采集→Kafka + Kafka 停留 + Watermark 等待(主要下界≈bound) |
| 2 | + Flink 计算 + Sink 批量 + 查询/缓存 |
| 3 | 设计目标合计 P95 < 5s |
削峰填谷:常态仍秒级;尖峰短暂 Lag 换稳定;死守 SLA 靠扩容而非只靠堆积。
| 差异来源 | 说明 |
|---|---|
| 迟到超 bound | 实时进侧输出;离线批可算全量 |
| 日界 / 时区 | 分钟窗 vs 业务日闭区间 |
| 去重 TTL | 状态过期 vs 离线全量 distinct |
| Sink 幂等合并前 | CH 可能短暂重复行 |
目标不是永远 0 差异,而是 差异可解释、可度量、有阈值告警。
CVR = 下单 UV / 浏览 UV 不能两个独立窗口直接相除。
- 会话窗:
keyBy(user)+ Session Window(gap 如 30min),会话内同时有 view 与 order 算转化 - Interval Join: 浏览流 ⋈ 订单流,下单落在浏览后 N 分钟内算转化
对齐后再按浏览时间进入大盘窗口。Join/会话必须设状态 TTL,防止状态爆炸。
- 背景:T+1 不够,要秒级看板与风控
- 架构:Kafka → Flink 三 Job → CH/Redis
- 难点 1:乱序与 Watermark / 迟到侧输出
- 难点 2:order_id 去重 + Checkpoint;CH 幂等折中
- 难点 3:UV 精确 vs HLL;流批对账
- 收尾:生产可补 CDC 维表、倾斜治理、SLA 大盘
问题讲解
登录后可以选中正文添加批注(仅自己可见)。
评论 (0)
登录后参与评论。
还没有评论,来做第一个。