CalcGuide · 技术博客主页 / 一页纸学习计划
🟠

流处理与实时管道:从 Kafka 到 Flink Exactly-Once

分类:数据与人工智能 · 路径:docs/topics/streaming-and-realtime/README.md

#kafka#flink#spark-streaming#stream-processing#exactly-once

用 Kafka + Flink/Spark Streaming 搭实时管道,掌握 Exactly-Once、Watermark、State、Window

父主题

数据工程基础:从 ETL 到数仓与流处理

子主题(0)

流处理与实时管道:从 Kafka 到 Flink Exactly-Once

0. 元信息

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-Once1.5 天2 天幂等 Producer + 事务
3. Flink 架构1 天1.5 天JobManager + TaskManager + Slot
4. DataStream API1.5 天2 天Kafka → Flink → Iceberg
5. Watermark 与 Window1.5 天2 天乱序容忍 + 7 天滑窗 UV
6. State Backend1 天1.5 天HashMap vs RocksDB
7. 两阶段提交 Sink0.5 天1 天TwoPhaseCommitSinkFunction
8. 流批一体 SQL0.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 APISource / Transform / Sink / Process Function一份 Kafka → Flink → Iceberg 的 Job能跑通简单计数
5. Watermark事件时间 vs 处理时间、Watermark 策略、迟到数据一份乱序事件流处理(容忍 5s 乱序)能解释 Watermark 与乱序关系
6. WindowTumbling / Sliding / Session / Global一份 7 天滑动窗口 UV 报表能区分四类窗口与触发时机
7. State BackendHashMapStateBackend / RocksDB / Embedded一份大状态 Job(>1GB)能解释 backend 选择与容错
8. 两阶段提交 SinkTwoPhaseCommitSinkFunction / 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 -dkafka-topics --create --topic orders --partitions 3 --replication-factor 3docker-compose.yml + topic 列表kafka-topics --listorders;Flink Web UI 8081 可见;docker logs kafka-1 无 error
Day 2Producer / 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 3Exactly-Once:写一个事务 Producer(init_transactions + begin_transaction + commit_transaction);故意在中间抛异常触发 abort_transaction;验证没有部分写入事务 Producer 脚本异常后 kafka-console-consumer 看不到部分消息;transactional.idacks=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 5Watermark:写一份乱序事件流(用 assignTimestampsAndWatermarks + forBoundedOutOfOrderness(Duration.ofSeconds(5)));用 TumblingEventTimeWindows.of(Time.minutes(1)) 聚合;故意晚到 30s 的事件走 Side OutputWatermark Job + Side Output 文件乱序 ≤5s 的事件全部进窗口;晚 30s 的事件落 Side Output 目录,验证 OutputTag 配置正确
Day 6State 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 分钟)

  1. kafka-topics --create --topic orders
  2. Producer 写 1 万条订单到 Kafka;
  3. Flink Job:Kafka Source → keyBy(user_id) → TumblingEventTimeWindow(5min) → reduce(SumAmount) → IcebergSink;
  4. SELECT * FROM iceberg.dws_user_order_5min 查汇总结果;
  5. 端到端 P95 < 30s。

步骤 B —— 4 类边界用例(30~45 分钟)

#用例期望行为验证命令
B1Kafka broker 宕机 1 个Producer 自动重试;Consumer 不丢数据docker stop kafka-1 后观察
B2TaskManager OOMJobManager 重启 Task;state 从 checkpoint 恢复docker kill taskmanager-1
B3迟到 1 小时事件落 Side Output;不进窗口结果pytest -k test_late_event
B4Schema 漂移(新字段)Flink Job 继续;警告日志改 JSON 加新字段再发

Day 7 当天必完成步骤 A;步骤 B 至少完成 B1、B3。

5. 阶段通用验收

  1. 不看答案独立重写一条 Kafka → Flink → Iceberg 的端到端管道 + 一次 EOS 验证;
  2. 用自己的话解释”为什么流处理需要 Watermark""为什么 Flink 与 Kafka 各自管一份 EOS""为什么 State 必须可序列化”;
  3. 画一张图:流处理架构图 + Watermark 时间线 + State 后端选型决策树(三件套之一);
  4. 测试乱序、迟到、节点宕机、Schema 漂移、状态膨胀 5 类边界;
  5. 准备至少 3 组自定义事件流(正常 / 乱序 / 迟到)并贴出实际输出与端到端 P95;
  6. 记录端到端 P50 / P99 延迟、吞吐(条/秒)、checkpoint 时长、state 大小 4 个核心指标;
  7. 能修改已有 Flink Job(加 Window、改 Watermark 策略、换 State 后端)并用同样指标验证。

交付存放:第 3 项的图、第 5 项的输出、第 6 项的指标,统一存到 week1/notes/week1/<job>/README.md

6. 最终验收

7. 综合项目

首选:CDC + 实时 dashboard(必做:MySQL → Kafka → Flink → Iceberg → Superset)。
备选:点击流实时管道(Nginx → Kafka → Flink → Iceberg → 实时 UV/PV 看板)。
备选:IoT 设备告警(MQTT → Kafka → Flink → Iceberg + 规则告警)。

CDC + 实时 dashboard 必做要求:

任何综合项目都必须包含:

  1. 需求说明与数据契约(源 schema、目标 schema、SLA、迟到容忍);
  2. 流处理架构图与命名规范;
  3. 核心代码(Debezium 配置 + Flink Job + Iceberg DDL + Superset 看板配置);
  4. 边界测试(broker 宕机、TaskManager OOM、迟到 1 小时、Schema 漂移、状态膨胀);
  5. 可观测(Prometheus 抓 Flink metrics + Grafana 看板);
  6. README(设计取舍、已知限制、下一步);
  7. 复盘记录(模拟 1 次 P0 故障,写 retrospective.md);
  8. 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 职责

  1. 用 Kafka 事务 Producer(transactional.id + enable.idempotence=true + acks=all)保证”发出去的消息要么全成功要么全失败”;Flink Kafka Source 用 setStartOffsets(OffsetsInitializer.committedOffsets()) + isolation.level=read_committed 只消费已提交消息。
  2. 用 Flink Watermark(WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)))+ TumblingEventTimeWindow + Side Output 处理乱序与迟到;两阶段提交 Sink(TwoPhaseCommitSinkFunction + Iceberg IcebergSink)保证端到端 exactly-once。
  3. 用 Flink Checkpoint(60 秒,CheckpointingMode.EXACTLY_ONCE)+ EmbeddedRocksDBStateBackend 把状态做可恢复;kill TaskManager 后从最近 checkpoint 续跑,Kafka offset 与 RocksDB state 双一致。

4 交付物

  1. 一条 Kafka → Flink → Iceberg 的端到端管道(FlinkKafkaConsumer + keyBy(user_id) + TumblingEventTimeWindow(5min) + IcebergSink),含 Watermark 策略与迟到容忍 5 秒。
  2. 一份 exactly-once 验证脚本(事务 Producer 抛异常 → abort_transaction → 消费侧看不到部分消息;kill TaskManager → 自动重启 → offset 与 state 完全一致),含 checkpoint 恢复证据目录。
  3. 一份 RocksDB State Backend 配置(state.backend=rocksdb + state.checkpoints.dir=s3://flink/checkpoints/),跑 1GB 状态不 OOM,checkpoint 时长 < 30 秒。
  4. 一份 Flink metrics + Prometheus 抓取配置(latencyMarker / numRecordsIn / numRecordsOut / checkpointDuration),4 个核心指标基线:端到端 P50/P99 延迟 / 吞吐(条/秒) / checkpoint 时长 / state 大小。

3 指标

  1. exactly-once 验证 0 重复 0 丢失(事务 abort 后消费侧消息数为 0;checkpoint 恢复后 Kafka offset 与 RocksDB state 完全一致)。
  2. 端到端 P99 延迟 < 30 s(含 Kafka produce + Flink 处理 + Iceberg commit)。
  3. 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.pdfWatermark / 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 案例
5WatermarkConfluent Bloghttps://www.confluent.io/blog/看 ksqlDB 与 Schema Registry
6StateVerverica Bloghttps://www.ververica.com/blogState 后端调优案例
7SinkIceberg 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.11https://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 BlogKafka 生态看 Schema Registry 与 ksqlDB
Flink Forward大会视频看 Best Practices
Ververica BlogFlink 商业看 Exactly-Once 案例

9.3.4 核心人物

人物影响材料
Jay KrepsKafka / KappaConfluent 博客
Tyler AkidauDataflow / BeamApache Beam 博客
Kostas TzoumasFlinkVerverica 博客
Stephan EwenFlink 起源Flink Forward

9.3.5 开发方法

方法动作何时用
Event-time first用事件时间,不用处理时间任何实时分析
Watermark with bounded out-of-orderness显式设乱序容忍任何流处理
State is checkpointedState 必须可序列化 + 增量任何有状态 Job
Idempotent Sink或用两阶段提交任何写入外部存储
Savepoint for migration升级前 Savepoint任何状态变更

9.4 经典问题与经典案例(≥5 道)

#问题重要性最简答案
1Kafka Exactly-Once 怎么保证一致性幂等 Producer + 事务 + 幂等 Consumer
2Flink 与 Kafka EOS 区别概念混淆Flink 用 barrier + 两阶段提交
3Watermark 怎么设乱序与迟到用 bounded out-of-orderness
4Window 三类怎么选业务决定Tumbling 切片、Sliding 滑窗、Session 会话
5State Backend 怎么选性能小状态 HashMap、大状态 RocksDB
6迟到数据怎么处理业务常见用 Allowed Lateness + Side Output
7Checkpoint 与 Savepoint 区别运维Checkpoint 自动、Savepoint 手动
8双流 Join 怎么写关联场景用 KeyedCoProcessFunction + 状态 TTL
9流处理重启怎么恢复容错配 checkpoint + 状态后端
10Kafka 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 Kafka3.xApacheGAApache-2.0
Apache Flink1.18+ApacheGAApache-2.0
Apache Spark Streaming3.5+ApacheGAApache-2.0
Apache Pulsar3.xApacheGAApache-2.0
Apache Beam2.xApacheGAApache-2.0

Scope

Kafka 是消息层,Flink/Spark 是流处理引擎;Pulsar 是消息+存储一体;Beam 是多引擎抽象。

Structure

Ecosystem

Depth Tiers

层级能力标准
L0知道存在知道 Kafka/Flink 基本概念
L1看得懂示例能读懂简单 Flink Job
L2能正确调用能用 Kafka/Flink 跑通简单管道
L3能解释与排错能定位乱序、迟到、状态丢失
L4能设计与扩展能设计 EOS 管道与状态管理

本子主题目标:L3

Source

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 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: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. 常见误区

11. 所有知识点分类

  1. 编程语言
  2. 数据结构与算法
  3. 计算机基础
  4. 工程技术
  5. Web 与后端
  6. 前端与客户端
  7. 数据与人工智能
  8. 项目与职业能力
  9. 安全与可靠性

本计划归属:数据与人工智能 主 + 工程技术 辅。


直接依赖(2)

查看知识图谱