流处理与实时管道:从 Kafka 到 Flink Exactly-Once
0. 元信息
- 主题路径:
docs/topics/data-engineering-basics/subtopics/streaming-and-realtime/README.md - 父主题:
data-engineering-basics - 适合对象:会写 Python/SQL、希望从”批处理思维”过渡到”实时数据流”的工程师
- 建议周期:2.5~3 周,每周 10~14 小时
- 前置知识:SQL、Python、Linux;建议先完成父主题第 3 阶段
- 最终目标:能为给定业务(点击流/订单/IoT)设计一条 Kafka → Flink → Iceberg 的实时管道,掌握 Exactly-Once、Watermark、State Backend、Window、迟到数据处理
1. 学习路线
Kafka 基础(broker / topic / partition / offset)→ Exactly-Once 语义 → Flink 架构(JobManager / TaskManager)→ DataStream API → Watermark 与 Window → State Backend(RocksDB)→ 两阶段提交 Sink → 流批一体 SQL
2. 阶段周数分配
| 阶段 | 2.5 周方案 | 3 周方案 | 备注 |
|---|---|---|---|
| 1. Kafka 基础 | 1.5 天 | 2 天 | broker + topic + consumer group |
| 2. Exactly-Once | 1.5 天 | 2 天 | 幂等 Producer + 事务 |
| 3. Flink 架构 | 1 天 | 1.5 天 | JobManager + TaskManager + Slot |
| 4. DataStream API | 1.5 天 | 2 天 | Kafka → Flink → Iceberg |
| 5. Watermark 与 Window | 1.5 天 | 2 天 | 乱序容忍 + 7 天滑窗 UV |
| 6. State Backend | 1 天 | 1.5 天 | HashMap vs RocksDB |
| 7. 两阶段提交 Sink | 0.5 天 | 1 天 | TwoPhaseCommitSinkFunction |
| 8. 流批一体 SQL | 0.5 天 | 1 天 | Flink SQL + Iceberg |
每天 1.5~2 小时。2.5 周方案聚焦核心 5 阶段;3 周方案多 2 天做 EOS 与流批一体。
3. 核心知识 / 产出 / 标准表
| 阶段 | 核心知识 | 实践产出 | 可观察学会标准 |
|---|---|---|---|
| 1. Kafka 基础 | broker / topic / partition / offset / consumer group / ISR | 一份 Kafka 集群(3 broker)+ 主题 | 能解释分区与副本机制 |
| 2. Exactly-Once | 幂等 Producer / 事务 / 幂等 Consumer | 一份事务消息样例 | 能解释 Kafka 与 Flink 的 EOS 区别 |
| 3. Flink 架构 | JobManager / TaskManager / Slot / Checkpoint | 一份 Flink Standalone 集群 | 能用 Web UI 跑 Job |
| 4. DataStream API | Source / Transform / Sink / Process Function | 一份 Kafka → Flink → Iceberg 的 Job | 能跑通简单计数 |
| 5. Watermark | 事件时间 vs 处理时间、Watermark 策略、迟到数据 | 一份乱序事件流处理(容忍 5s 乱序) | 能解释 Watermark 与乱序关系 |
| 6. Window | Tumbling / Sliding / Session / Global | 一份 7 天滑动窗口 UV 报表 | 能区分四类窗口与触发时机 |
| 7. State Backend | HashMapStateBackend / RocksDB / Embedded | 一份大状态 Job(>1GB) | 能解释 backend 选择与容错 |
| 8. 两阶段提交 Sink | TwoPhaseCommitSinkFunction / Iceberg / Kafka | 一份端到端 EOS 管道 | 能解释 barrier 对齐与 commit |
4. 第一周(每天 1.5~2 小时)
Day 1 约定:本计划使用 Docker Compose 起 Kafka(3 broker)+ Flink(Standalone)+ Iceberg on MinIO。Kafka 测试用
kafka-console-producer/kafka-console-consumer;Flink 跑flink run;SQL 客户端sql-client.sh。所有 Job 必须可重跑:清空 topic + 重置 offset 后跑出相同结果。
| 日 | 任务 | 当天交付 | 自检 |
|---|---|---|---|
| Day 1 | 装 Docker Compose + Kafka(3 broker,KRaft 模式)+ Flink Standalone 集群;起 docker compose up -d;kafka-topics --create --topic orders --partitions 3 --replication-factor 3 | docker-compose.yml + topic 列表 | kafka-topics --list 含 orders;Flink Web UI 8081 可见;docker logs kafka-1 无 error |
| Day 2 | Producer / Consumer 基础:用 kafka-console-producer 写 100 条订单;用 kafka-console-consumer --from-beginning 读出来;用 Python confluent-kafka 写一个 Producer 脚本 | Producer 脚本 + 100 条样例 | kafka-console-consumer --topic orders --from-beginning --max-messages 100 输出 100 条;kafka-consumer-groups --describe --group test-group 显示 offset |
| Day 3 | Exactly-Once:写一个事务 Producer(init_transactions + begin_transaction + commit_transaction);故意在中间抛异常触发 abort_transaction;验证没有部分写入 | 事务 Producer 脚本 | 异常后 kafka-console-consumer 看不到部分消息;transactional.id 配 acks=all + enable.idempotence=true |
| Day 4 | 第一个 Flink Job:Kafka Source → Map(解析 JSON)→ print;用 flink run 提交;故意让 TaskManager 挂掉,验证从 checkpoint 恢复 | Flink Job JAR + 恢复日志 | flink run 提交后 TaskManager 杀进程后重启,offset 继续推进;/tmp/flink-checkpoints/ 有新 checkpoint |
| Day 5 | Watermark:写一份乱序事件流(用 assignTimestampsAndWatermarks + forBoundedOutOfOrderness(Duration.ofSeconds(5)));用 TumblingEventTimeWindows.of(Time.minutes(1)) 聚合;故意晚到 30s 的事件走 Side Output | Watermark Job + Side Output 文件 | 乱序 ≤5s 的事件全部进窗口;晚 30s 的事件落 Side Output 目录,验证 OutputTag 配置正确 |
| Day 6 | State Backend:写一个 KeyedProcessFunction,状态用 ValueState<Long> 累加;先用 HashMapStateBackend 跑;改用 EmbeddedRocksDBStateBackend 跑 1GB 状态;对比 checkpoint 时长与吞吐 | State Job + 性能对比表 | HashMap 跑 100MB OK;RocksDB 跑 1GB 不 OOM;checkpoint 时长对比表(HashMap vs RocksDB) |
| Day 7 | 步骤 A:跑通 Kafka → Flink → Iceberg 的完整管道(订单事件 → 5 分钟滑窗汇总 → 写入 Iceberg DWS 表);步骤 B:补齐 4 类边界(Kafka broker 宕机、TaskManager OOM、迟到 1 小时、Schema 漂移) | 完整管道 + 故障演练报告 | pytest -k test_streaming_eos 4 个用例全过;每类故障复现 + 防护措施写到 notes/week1-day7.md |
Day 7 执行次序
步骤 A —— 最小可用流管道(60~90 分钟)
kafka-topics --create --topic orders;- Producer 写 1 万条订单到 Kafka;
- Flink Job:Kafka Source → keyBy(user_id) → TumblingEventTimeWindow(5min) → reduce(SumAmount) → IcebergSink;
SELECT * FROM iceberg.dws_user_order_5min查汇总结果;- 端到端 P95 < 30s。
步骤 B —— 4 类边界用例(30~45 分钟)
| # | 用例 | 期望行为 | 验证命令 |
|---|---|---|---|
| B1 | Kafka broker 宕机 1 个 | Producer 自动重试;Consumer 不丢数据 | docker stop kafka-1 后观察 |
| B2 | TaskManager OOM | JobManager 重启 Task;state 从 checkpoint 恢复 | docker kill taskmanager-1 |
| B3 | 迟到 1 小时事件 | 落 Side Output;不进窗口结果 | pytest -k test_late_event |
| B4 | Schema 漂移(新字段) | Flink Job 继续;警告日志 | 改 JSON 加新字段再发 |
Day 7 当天必完成步骤 A;步骤 B 至少完成 B1、B3。
5. 阶段通用验收
- 不看答案独立重写一条 Kafka → Flink → Iceberg 的端到端管道 + 一次 EOS 验证;
- 用自己的话解释”为什么流处理需要 Watermark""为什么 Flink 与 Kafka 各自管一份 EOS""为什么 State 必须可序列化”;
- 画一张图:流处理架构图 + Watermark 时间线 + State 后端选型决策树(三件套之一);
- 测试乱序、迟到、节点宕机、Schema 漂移、状态膨胀 5 类边界;
- 准备至少 3 组自定义事件流(正常 / 乱序 / 迟到)并贴出实际输出与端到端 P95;
- 记录端到端 P50 / P99 延迟、吞吐(条/秒)、checkpoint 时长、state 大小 4 个核心指标;
- 能修改已有 Flink Job(加 Window、改 Watermark 策略、换 State 后端)并用同样指标验证。
交付存放:第 3 项的图、第 5 项的输出、第 6 项的指标,统一存到
week1/notes/或week1/<job>/README.md。
6. 最终验收
-
独立设计并交付一条端到端实时管道(Kafka → Flink → Iceberg),覆盖 EOS、Watermark、State、Window、迟到处理;
-
至少完成 18 个实战项(分布建议):
子阶段 题目数量 难度 平台建议 Kafka 基础与 EOS 3 Easy / Medium Confluent Quickstart Flink DataStream API 4 Medium Flink Forward Talks Watermark 与 Window 4 Medium / Hard Streaming Systems 案例 State Backend 与容错 4 Medium Flink 官方训练 流批一体 SQL 3 Hard Flink SQL + Iceberg 案例 约束:至少 9 项达到 Medium,至少 3 项达到 Hard;每项必须留可运行 Job + checkpoint 恢复证据。
-
完成 1 个综合项目:CDC + 实时 dashboard(Debezium → Kafka → Flink → Iceberg → Superset);
-
能用 15 分钟讲清 Flink / Kafka 各自 EOS 实现、Watermark 与迟到数据处理、State 后端选型、checkpoint 与 savepoint 区别。
7. 综合项目
首选:CDC + 实时 dashboard(必做:MySQL → Kafka → Flink → Iceberg → Superset)。
备选:点击流实时管道(Nginx → Kafka → Flink → Iceberg → 实时 UV/PV 看板)。
备选:IoT 设备告警(MQTT → Kafka → Flink → Iceberg + 规则告警)。
CDC + 实时 dashboard 必做要求:
- 输入:MySQL
orders表(开启 binlog row 模式)+ 模拟写入工具(atlas / 自写脚本); - 输出:必输出(1)Debezium 抓取 binlog 到 Kafka topic
cdc.orders;(2)Flink Job 消费cdc.orders→ 写 Iceberg DWD + DWS(含 5 分钟汇总);(3)Superset 看板:实时订单量 / GMV / 转化率;(4)一次故障演练:故意延迟 1 小时、kill TaskManager、Schema 加列; - 算法 / 工程:用
flink-connector-debezium+IcebergSink;用WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5));Checkpoint 60 秒 + EOS; - 进阶可选:用 StarRocks / ClickHouse 替代 Superset 做 OLAP;接 OpenLineage 输出血缘事件。
任何综合项目都必须包含:
- 需求说明与数据契约(源 schema、目标 schema、SLA、迟到容忍);
- 流处理架构图与命名规范;
- 核心代码(Debezium 配置 + Flink Job + Iceberg DDL + Superset 看板配置);
- 边界测试(broker 宕机、TaskManager OOM、迟到 1 小时、Schema 漂移、状态膨胀);
- 可观测(Prometheus 抓 Flink metrics + Grafana 看板);
- README(设计取舍、已知限制、下一步);
- 复盘记录(模拟 1 次 P0 故障,写
retrospective.md); - notes/ 规范:
| 文件 | 内容 |
|---|---|
notes/design.md | 架构图、Watermark 策略、State 后端选型理由 |
notes/test.md | 每组测试的输入 / 期望 / 实际 / 通过情况 |
notes/retrospective.md | 用时、难点、收获、改进点 |
notes/runbook.md | 常见失败 + 处置 SOP(broker 宕机 / 状态丢失 / 反压) |
notes/flink-ui-screenshots/ | Flink Web UI 截图(Task 状态 / Checkpoint 时长 / 反压) |
notes/checkpoint-evidence/ | checkpoint 目录列表 + 恢复日志 |
notes/contract-yaml/ | 数据契约(Producers / Topics / Schemas) |
notes/sample-events/ | 至少 3 组测试事件流(正常 / 乱序 / 迟到) |
本主题贡献
实时管道的核心难点是”消息天然会丢会重 vs 业务要求 exactly-once”。本子主题专门讲清 Kafka → Flink 如何用两阶段提交 + 幂等 Producer + Watermark 把这条边界守住——以 Kafka 为消息总线,Flink 为计算引擎,KStream / DataStream API 为编程抽象,watermark 与 checkpoint 双驱动 exactly-once;不重复父主题的 Airflow 编排与 Kafka 基础。
3 职责
- 用 Kafka 事务 Producer(
transactional.id+enable.idempotence=true+acks=all)保证”发出去的消息要么全成功要么全失败”;Flink Kafka Source 用setStartOffsets(OffsetsInitializer.committedOffsets())+isolation.level=read_committed只消费已提交消息。 - 用 Flink Watermark(
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)))+TumblingEventTimeWindow+ Side Output 处理乱序与迟到;两阶段提交 Sink(TwoPhaseCommitSinkFunction+ IcebergIcebergSink)保证端到端 exactly-once。 - 用 Flink Checkpoint(60 秒,
CheckpointingMode.EXACTLY_ONCE)+EmbeddedRocksDBStateBackend把状态做可恢复;kill TaskManager 后从最近 checkpoint 续跑,Kafka offset 与 RocksDB state 双一致。
4 交付物
- 一条 Kafka → Flink → Iceberg 的端到端管道(
FlinkKafkaConsumer+keyBy(user_id)+TumblingEventTimeWindow(5min)+IcebergSink),含 Watermark 策略与迟到容忍 5 秒。 - 一份 exactly-once 验证脚本(事务 Producer 抛异常 →
abort_transaction→ 消费侧看不到部分消息;kill TaskManager → 自动重启 → offset 与 state 完全一致),含 checkpoint 恢复证据目录。 - 一份 RocksDB State Backend 配置(
state.backend=rocksdb+state.checkpoints.dir=s3://flink/checkpoints/),跑 1GB 状态不 OOM,checkpoint 时长 < 30 秒。 - 一份 Flink metrics + Prometheus 抓取配置(
latencyMarker/numRecordsIn/numRecordsOut/checkpointDuration),4 个核心指标基线:端到端 P50/P99 延迟 / 吞吐(条/秒) / checkpoint 时长 / state 大小。
3 指标
- exactly-once 验证 0 重复 0 丢失(事务 abort 后消费侧消息数为 0;checkpoint 恢复后 Kafka offset 与 RocksDB state 完全一致)。
- 端到端 P99 延迟 < 30 s(含 Kafka produce + Flink 处理 + Iceberg commit)。
- Checkpoint 时长 < 30 s(1GB 状态 + RocksDB 增量 checkpoint)。
8. 推荐开源资料
| 阶段 | 角色 | 资料 | 链接 | 用法 |
|---|---|---|---|---|
| 全部 | 主线书 | Akidau 等《Streaming Systems》 | https://www.oreilly.com/library/view/streaming-systems/9781491983867/ | 流处理三大问题(What/Where/When)必读 1~5 章 |
| 全部 | 概念 | Akidau 等《The Dataflow Model》 | https://www.vldb.org/pvldb/vol8/p1792-Akidau.pdf | Watermark / Trigger / Accumulation 起源 |
| 1~2 | 消息 | Apache Kafka 官方文档 | https://kafka.apache.org/documentation/ | Producer / Consumer / Connect / Streams |
| 1~2 | 消息 | Neha Narkhede《Kafka: The Definitive Guide》 | https://www.oreilly.com/library/view/kafka-the-definitive/9781492043089/ | Kafka 实战 1~6 章 |
| 3~4 | 流处理 | Apache Flink 官方文档 | https://nightlies.apache.org/flink/flink-docs-stable/ | DataStream / Table / Checkpoint |
| 3~4 | 流处理 | Flink Forward 大会视频 | https://www.flink-forward.org/ | 看 Best Practices 与 Exactly-Once 案例 |
| 5 | Watermark | Confluent Blog | https://www.confluent.io/blog/ | 看 ksqlDB 与 Schema Registry |
| 6 | State | Ververica Blog | https://www.ververica.com/blog | State 后端调优案例 |
| 7 | Sink | Iceberg Flink 集成 | https://iceberg.apache.org/docs/latest/flink/ | IcebergSink 用法 |
| 8 | 流批一体 | Flink SQL 官方文档 | https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/table/ | 流批 SQL 一致性 |
| 全部 | 故障演练 | Kleppmann《DDIA》Ch.11 | https://dataintensive.net/ | 流处理与一致性 |
许可证提示:Kafka、Flink、Iceberg 都是 Apache-2.0;Confluent 部分组件(ksqlDB、Schema Registry)有社区版与商业版之分,使用前查 LICENSE。默认做法是读思路后自己实现,不要直接复制 Ververica 商业版代码到生产。
默认使用顺序:先读 Akidau《Streaming Systems》第 13 章建立”流处理心智模型” → 装 Kafka 跑通 Producer/Consumer → 读 Narkhede 第 14 章 → 跑第一个 Flink Job → 读 Flink 官方文档 concepts → 加 Watermark + Window → 换 RocksDB State 后端 → 接 Iceberg Sink → 用 EXPLAIN 看 Flink SQL 计划 → 对比批流同一查询结果 → 写复盘到 notes/retrospective.md。
9. 学习资料汇聚(v0.3 自包含)
本节由本计划生成。链接指向原始材料或作者公开内容。规范会演进,记录时务必写明版本与日期。
9.1 背景与动机
Lambda 架构 2011 年由 Nathan Marz 提出,用”批处理层 + 速度层”双路处理实时数据,但维护两套代码。2014 年 Jay Kreps 提出 Kappa 架构,主张用单一日志流处理一切。Google 2015 年发表 Dataflow 论文(VLDB 2015),把 Watermark、Trigger、Accumulation 三大概念写进流处理理论。Apache Flink 2014 年从柏林大学开源,定位为”原生流处理引擎”;Spark Streaming 2013 年用微批(micro-batch)模拟流;Kafka 2011 年由 LinkedIn 开源,今天是事实标准消息层。三者共同构成实时数据流的事实标准技术栈。
9.2 概念地图
flowchart LR
Src[源 DB/日志/IoT] --> K[Kafka topic]
K --> F[Flink Job]
K --> SS[Spark Streaming]
F --> S1[Window]
F --> S2[State]
F --> S3[Watermark]
F --> Sink[Sink]
Sink --> Ice[Iceberg]
Sink --> K2[Kafka topic]
Sink --> DB[PG/MySQL]
F --> Checkpoint[Checkpoint]
F --> EOS[Exactly-Once]
9.3 基础知识讲解
9.3.1 经典论文
| 资料 | 贡献 | 读法 |
|---|---|---|
| Tyler Akidau et al., The Dataflow Model(VLDB 2015) | 流处理三大问题 | 必读 What/Where/When |
| Jay Kreps, Questioning the Lambda Architecture(2014) | Kappa 架构 | 必读 |
| Kostas Tzoumas et al., State Management in Apache Flink(2016) | State backend | 选读 |
| Paris Carbone et al., Apache Flink: Stream and Batch Processing in a Single Engine(2015) | Flink 架构 | 选读 |
9.3.2 经典书籍
| 书 | 侧重 | 用法 |
|---|---|---|
| Tyler Akidau et al., Streaming Systems(O’Reilly 2018) | 流处理三大问题 | 必读 1~5 章 |
| Holden Karau et al., Spark: The Definitive Guide(O’Reilly 2018) | Spark 全景 | 选读第 8~9 章 |
| Neha Narkhede et al., Kafka: The Definitive Guide(2nd ed., O’Reilly 2021) | Kafka 全景 | 必读 1~6 章 |
9.3.3 优秀博客与文档
| 资料 | 特点 | 用法 |
|---|---|---|
| Apache Kafka 官方文档 | Producer/Consumer/Streams | 必读 design + config |
| Apache Flink 官方文档 | DataStream/Table API | 必读 concepts + ops |
| Confluent Blog | Kafka 生态 | 看 Schema Registry 与 ksqlDB |
| Flink Forward | 大会视频 | 看 Best Practices |
| Ververica Blog | Flink 商业 | 看 Exactly-Once 案例 |
9.3.4 核心人物
| 人物 | 影响 | 材料 |
|---|---|---|
| Jay Kreps | Kafka / Kappa | Confluent 博客 |
| Tyler Akidau | Dataflow / Beam | Apache Beam 博客 |
| Kostas Tzoumas | Flink | Ververica 博客 |
| Stephan Ewen | Flink 起源 | Flink Forward |
9.3.5 开发方法
| 方法 | 动作 | 何时用 |
|---|---|---|
| Event-time first | 用事件时间,不用处理时间 | 任何实时分析 |
| Watermark with bounded out-of-orderness | 显式设乱序容忍 | 任何流处理 |
| State is checkpointed | State 必须可序列化 + 增量 | 任何有状态 Job |
| Idempotent Sink | 或用两阶段提交 | 任何写入外部存储 |
| Savepoint for migration | 升级前 Savepoint | 任何状态变更 |
9.4 经典问题与经典案例(≥5 道)
| # | 问题 | 重要性 | 最简答案 |
|---|---|---|---|
| 1 | Kafka Exactly-Once 怎么保证 | 一致性 | 幂等 Producer + 事务 + 幂等 Consumer |
| 2 | Flink 与 Kafka EOS 区别 | 概念混淆 | Flink 用 barrier + 两阶段提交 |
| 3 | Watermark 怎么设 | 乱序与迟到 | 用 bounded out-of-orderness |
| 4 | Window 三类怎么选 | 业务决定 | Tumbling 切片、Sliding 滑窗、Session 会话 |
| 5 | State Backend 怎么选 | 性能 | 小状态 HashMap、大状态 RocksDB |
| 6 | 迟到数据怎么处理 | 业务常见 | 用 Allowed Lateness + Side Output |
| 7 | Checkpoint 与 Savepoint 区别 | 运维 | Checkpoint 自动、Savepoint 手动 |
| 8 | 双流 Join 怎么写 | 关联场景 | 用 KeyedCoProcessFunction + 状态 TTL |
| 9 | 流处理重启怎么恢复 | 容错 | 配 checkpoint + 状态后端 |
| 10 | Kafka topic 分区数怎么定 | 吞吐 | 取最大并发 × 1.5 |
9.5 学习难点
概念难点
| 难点 | 突破路径 |
|---|---|
| 事件时间 vs 处理时间 | 画时间线 |
| Exactly-Once 边界 | 区分 Kafka、Flink、外部存储三层 |
| Watermark 与迟到 | 跑含乱序数据的实验 |
思维难点
| 难点 | 突破路径 |
|---|---|
| 状态管理 | 画状态机 |
| Checkpoint barrier 对齐 | 看 Flink 文档图示 |
工程难点
| 难点 | 突破路径 | |---|---|---| | RocksDB 调优 | 调 write buffer、block cache | | Checkpoint 失败 | 检查 state 大小、磁盘、网络 | | 反压 | 调 buffer pool、并行度 |
9.6 技术标准与接口
Entity
| 名称 | 版本 | 组织 | 状态 | 许可证 |
|---|---|---|---|---|
| Apache Kafka | 3.x | Apache | GA | Apache-2.0 |
| Apache Flink | 1.18+ | Apache | GA | Apache-2.0 |
| Apache Spark Streaming | 3.5+ | Apache | GA | Apache-2.0 |
| Apache Pulsar | 3.x | Apache | GA | Apache-2.0 |
| Apache Beam | 2.x | Apache | GA | Apache-2.0 |
Scope
Kafka 是消息层,Flink/Spark 是流处理引擎;Pulsar 是消息+存储一体;Beam 是多引擎抽象。
Structure
- Kafka:
Producer → Broker → Topic → Partition → Consumer Group。 - Flink:
JobManager + TaskManager + Slot + Operator Chain + Checkpoint + Savepoint。 - Spark Streaming:
DStream → Batch Interval → RDD/Dataset。
Ecosystem
- 消息:Kafka、Pulsar、NATS、RabbitMQ、Redis Streams。
- 流处理:Flink、Spark Streaming、Beam、ksqlDB。
- Sink:Iceberg、Delta、Hudi、PostgreSQL、Elasticsearch、Kafka。
- Schema:Schema Registry、Avro、Protobuf。
Depth Tiers
| 层级 | 能力 | 标准 |
|---|---|---|
| L0 | 知道存在 | 知道 Kafka/Flink 基本概念 |
| L1 | 看得懂示例 | 能读懂简单 Flink Job |
| L2 | 能正确调用 | 能用 Kafka/Flink 跑通简单管道 |
| L3 | 能解释与排错 | 能定位乱序、迟到、状态丢失 |
| L4 | 能设计与扩展 | 能设计 EOS 管道与状态管理 |
本子主题目标:L3。
Source
- Apache Kafka 官方文档:引用快照 2026-07-30。
- Apache Flink 官方文档:引用快照 2026-07-30。
- Streaming Systems:引用快照 2026-07-30。
9.7 关键代码
9.7.1 Kafka 幂等 Producer + 事务
# 9.7.1 Kafka 事务消息(Python confluent-kafka)
from confluent_kafka import Producer, Consumer
from confluent_kafka.admin import AdminClient
producer = Producer({
"bootstrap.servers": "kafka:9092",
"enable.idempotence": True,
"transactional.id": "tx-orders-1",
"acks": "all",
})
producer.init_transactions()
producer.begin_transaction()
try:
for order in orders:
producer.produce("orders", key=str(order["id"]), value=json.dumps(order))
producer.commit_transaction()
except Exception:
producer.abort_transaction()
raise
9.7.2 Flink DataStream Job(含 Watermark)
// 9.7.2 Flink Java DataStream:Kafka -> Iceberg with Watermark
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
DataStream<Order> orders = env
.addSource(new FlinkKafkaConsumer<>("orders", new OrderDeserializer(), kafkaProps))
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((o, ts) -> o.getEventTime())
);
orders
.keyBy(Order::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.reduce(new SumAmount())
.addSink(new IcebergSink.Builder<OrderSummary>()
.setTable(TableIdentifier.of("dw", "dws_user_order_5min"))
.setUpsert(true)
.build());
env.execute("orders-windowed-aggregation");
9.7.3 Flink SQL:流批一体查询
-- 9.7.3 Flink SQL:Kafka 表 + Iceberg 表 JOIN(流批一体)
CREATE TABLE kafka_orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10, 2),
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'format' = 'json',
'properties.bootstrap.servers' = 'kafka:9092'
);
CREATE TABLE iceberg_dwd_orders (...) WITH (
'connector' = 'iceberg',
'catalog-name' = 'iceberg_catalog',
'database-name' = 'dw'
);
INSERT INTO iceberg_dwd_orders
SELECT * FROM kafka_orders WHERE amount > 0;
10. 常见误区
- 流处理当批处理写(无 Watermark);
- Kafka 用
acks=0当 Exactly-Once; - Flink 用 HashMapState 处理大状态 OOM;
- 迟到数据不配 Allowed Lateness 直接丢;
- 双流 Join 不设状态 TTL,状态无限膨胀;
- Checkpoint 间隔设太短(< 1 分钟),影响吞吐;
- 反压靠加机器不调并行度;
- Flink Job 不设 RestartStrategy;
- Savepoint 与 Checkpoint 混用;
- Kafka topic 分区数拍脑袋;
- Stream SQL 当流式 ETL 写,不设状态;
- 不用 Schema Registry,事件结构漂移。
11. 所有知识点分类
- 编程语言
- 数据结构与算法
- 计算机基础
- 工程技术
- Web 与后端
- 前端与客户端
- 数据与人工智能
- 项目与职业能力
- 安全与可靠性
本计划归属:数据与人工智能 主 + 工程技术 辅。