Skip to content

从 0 到 1 学习 Kafka · 总学习计划

本计划由 task.md 第六节大纲展开而来,是后续 21 章 + 3 个附录的「施工图」。 配套环境:根目录 docker-compose.yml(3 节点 KRaft + Kafka UI + Schema Registry)+ requirements.txt(Python confluent-kafka 等依赖)。

📌 教程总目标

完成本教程后,你应当具备以下能力:

  1. 会用:写得出生产 / 消费代码,会用 kafka-*.sh 系列命令调试集群、查 Topic、看 Lag、重置 Offset。
  2. 懂原理:能用生活类比把「顺序写 / PageCache / 零拷贝 / ISR / HW / LEO / Leader Epoch / KRaft / Rebalance / Exactly Once / Compaction」讲清楚。
  3. 能调优:知道分区数 / 副本数 / acks / linger.ms / batch.size / compression.type / min.insync.replicas 怎么选;能从 JMX 指标定位瓶颈。
  4. 能答题:覆盖大厂后端 / 中间件 / 大数据方向 Kafka 高频面试题(与 RabbitMQ / RocketMQ / Pulsar 横向对比、消息语义、事务、KRaft 演进、踩坑案例)。

一、章节总表(21 章 + 3 个附录)

字数为文档主体的预估值,不含代码 / demo.html。「demo 要点」是同名子目录 XX_xxx/demo.html 必须实现的可视化交互。

#章节核心知识点实战案例 / 代码demo 要点面试题数预计字数
1Kafka 是什么 & 为什么快MQ 价值(解耦/异步/削峰/日志总线)、Kafka 诞生(LinkedIn)、设计哲学(以日志为中心)、五大性能基石(顺序写 / PageCache / 零拷贝 / 批量 / 压缩)、与 RabbitMQ / RocketMQ / Pulsar / Redis Stream / NATS 对比hello_kafka.py:建 Topic + 生产 5 条 + 消费打印① 顺序写 vs 随机写吞吐柱状图;② 5 大 MQ 存储/消费模型卡片;③ Kafka 生态全景拓扑64500
2核心概念与架构总览Producer / Broker / Consumer / Topic / Partition / Replica / Offset / Consumer Group / Coordinator / Controller / KRaft 全部术语;三横三纵分层;内部 Topic(__consumer_offsets / __transaction_state / __cluster_metadatacluster_probe.py:用 AdminClient 打印 Broker / Topic / Leader / ISR / Cluster ID交互式集群拓扑:可调 Broker / Partition / Replica,模拟 Broker 宕机看 Leader 切换54500
3命令行与客户端基础kafka-topics.sh / kafka-console-producer.sh / kafka-console-consumer.sh / kafka-consumer-groups.sh / kafka-configs.sh / kafka-dump-log.sh;Python confluent-kafka Hello World;Kafka UI 界面导览topic_admin.py:增删改查 Topic;producer_consumer.py:最小可运行示例命令速查交互卡片:点命令显示真实输入 / 输出对照54000
4Producer 深入同步 / 异步 send;acks=0/1/-1 语义;retries / delivery.timeout.ms / max.in.flight.requests.per.connectionlinger.ms / batch.size / compression.type;Partitioner(Murmur2 / Sticky / 自定义);幂等 Producer(PID + 序列号 + Broker 去重)sync_producer.py / async_producer.py / benchmark.py(不同 acks × linger × 压缩算法的吞吐对比)linger.ms / batch.size / 压缩算法滑块 → 实时算吞吐 / 延迟曲线76500
5Consumer 深入subscribe vs assignpoll 循环模型;Offset 提交(自动 / 同步 / 异步 / 手动);__consumer_offsetsauto.offset.reset;心跳 / session.timeout.ms / max.poll.interval.ms;Consumer Group ↔ Coordinator 6 步协议;seek / pause / resumemanual_commit.py:业务主键 + 数据库去重;seek_demo.py:按时间戳回溯消费poll 循环时序图 + Offset 提交时机决策树66500
6Topic 设计与分区策略分区数选型(吞吐 / 并行度 / 有序性 / 扩容代价);副本数选型;Key 与顺序保证 / 热点;min.insync.replicas;何时「重建 Topic」而非「加分区」partition_calc.py:根据目标吞吐推荐分区数;hot_key_demo.py:复现热点分区多车道高速公路动画:单分区 vs 多分区的顺序差异 + Key Hash 路由可视化55000
7存储与日志格式日志目录 log.dirs;Topic-Partition-Segment 三层;.log / .index / .timeindex / .snapshot / leader-epoch-checkpointindex.interval.bytes 稀疏索引;V0/V1/V2 消息格式与 RecordBatch;kafka-dump-log.sh 解码;segment 滚动条件dump_segment.sh 封装;造数 + 滚动 + 解码完整流水线.log / .index / .timeindex 二分查找定位 Offset 动画66500
8高吞吐的底层原理顺序写 vs 随机写量级差;PageCache(为什么不在 JVM 堆里缓存);零拷贝 sendfile / transferTo(4 拷贝 → 2 拷贝);批量 + 压缩三端协同;mmap 在索引文件上的应用bench_seq_vs_rand.py:用 fio / Python 复现顺序 vs 随机写差距4 拷贝 vs 2 拷贝动画 + PageCache 命中率仪表盘66500
9副本机制与 ISRLeader / Follower / Observer;ISR 定义与 replica.lag.time.max.ms;HW vs LEO;Leader Epoch 解决 HW 截断;unclean.leader.election.enable 代价;min.insync.replicas + acks=all 组合;故障场景三连(Follower 掉队 / Leader 宕机 / 脑裂)kill_leader_demo.py:docker stop broker 观察 HA;isr_lag_inject.py:注入网络延迟看 ISR 收缩LEO / HW 推进时序动画 + Leader Epoch 截断防护对比77000
10Controller 与 KRaftZK 时代 Controller 选举与元数据广播 → 历史包袱;KRaft 诞生(Raft 化的 __cluster_metadata Topic、Controller Quorum、事件溯源);KRaft 选举 / 快照 / 恢复;ZK → KRaft 迁移路径;4.0 全面移除 ZKkraft_election_probe.py:从 metadata.log 读 Raft 选举记录;metadata_quorum_status.sh 解读KRaft 选举动画:3 controller 投票 / 日志复制 / 快照66500
11消费者组与 RebalanceGroup Coordinator 角色;JoinGroup / SyncGroup / Heartbeat / LeaveGroup;分配策略(Range / RoundRobin / Sticky / CooperativeSticky);Eager Rebalance 的 STW 之痛;Static Membership;Rebalance 监听器;Lag 排查rebalance_observer.py:监听 ConsumerRebalanceListener;static_membership_demo.py4 种策略动画对比:增减 Consumer / Topic 看分区如何重分配77000
12Offset 与「消息语义」At Most Once / At Least Once / Exactly Once 工程含义;自动提交陷阱;手动提交正确姿势;消费者端幂等模式(业务主键 / Redis SETNX / DB 唯一约束);seek 到 Offset / 时间戳;__consumer_offsets 存储结构idempotent_consumer.py:业务表去重;seek_to_timestamp.py三种语义对比时序图 + 自动提交 vs 手动提交风险点55500
13幂等与事务:Exactly Once 的真相幂等 Producer(PID + Epoch + 序列号);事务 Producer 全 API;__transaction_state;Transaction Coordinator 两阶段提交;isolation.level=read_committed;Streams EOS v2;EOS 性能代价idempotent_producer.py / transactional_producer.py / read_committed_consumer.py;扣款转账演示EOS 两阶段提交时序动画:Producer ↔ TC ↔ Partition Leader77500
14Log Compaction 与 Tombstonedelete / compact / compact,deletemin.cleanable.dirty.ratio / segment.ms;null value 作 Tombstone;Compaction 在状态快照 / CDC / 内部 Topic 的应用;线程模型与性能影响compaction_demo.py:构造 100 万 key + 反复 update 看压缩;tombstone_delete.py同 Key 多版本压缩动画 + Tombstone 删除流程55000
15安全与多租户SSL / SASL(PLAIN / SCRAM / GSSAPI / OAUTHBEARER);ACL(Topic / Group / Cluster / TransactionalId);Quota(Producer / Consumer / Request);SuperUser 与审计;最小权限模型setup_scram_user.shacl_demo.sh;带 SASL 的 Python Producer权限矩阵交互查询 + Quota 限速效果曲线55000
16Kafka ConnectSource / Sink Connector;Standalone vs Distributed;常见 Connector(JDBC / Debezium / Elasticsearch / S3 / MongoDB / HDFS);Converter(JSON / Avro / Protobuf);SMT;DLQ;Exactly Once SourceDebezium MySQL → Kafka → JDBC Sink → PG 完整链路;connect_rest_calls.pyConnect 集群拓扑 + Connector 状态可视化 + DLQ 流转动画55500
17Kafka Streams 与 ksqlDB流-表二元性;KStream / KTable / GlobalKTable;无状态 / 有状态算子;State Store(RocksDB) + Changelog;时间语义(Event / Processing / Ingestion);窗口(Tumbling / Hopping / Sliding / Session);ksqlDB SQL 流处理wordcount_streams.java(Streams DSL);realtime_topn.py(Faust 版);ksqlDB 实时 GMV 大屏KStream / KTable 二元性可视化 + 4 种窗口聚合演示67000
18Schema Registry 与数据治理为什么要 Schema;Avro / Protobuf / JSON Schema 三家对比;兼容性策略(BACKWARD / FORWARD / FULL / NONE 及 _TRANSITIVE);_schemas Topic;CDC 表结构变更应对avro_producer.py + avro_consumer.pyschema_evolution_demo.py6 种兼容性策略测试矩阵:加字段 / 删字段 / 改类型 → 自动判定55500
19可观测性与运维核心指标(UnderReplicatedPartitions / OfflinePartitionsCount / ActiveControllerCount / RequestHandlerAvgIdlePercent / Network IO / PageCache 命中率 / record-error-rate / records-lag-max);JMX + Prometheus JMX Exporter + Grafana;日志定位(server.log / controller.log / state-change.log / log-cleaner.log);运维命令(kafka-reassign-partitions.sh / preferred replica election / 限流);调优清单jmx_exporter_setup.shreassign_partitions.py;模拟 Grafana 面板的 Python 仪表盘模拟 Grafana 看板 + 异常注入触发指标变化66500
20常见踩坑与排障案例集分区数选少 → 扩容难;Key 设计不当 → 热点;acks=1 丢消息;auto.offset.reset=latest 丢历史;消费慢 + max.poll.interval.ms 超时 → 死循环 Rebalance;unclean.leader.election=true → 已 ack 消息丢失;大消息撑爆 message.max.bytes;Compaction 不生效;Connect Offset Topic 污染;跨机房 ISR 抖动每个踩坑都有可复现脚本 + 修复后对照「踩坑还原器」交互菜单:选案例 → 看症状 → 一键复现 → 看根因87500
21综合实战项目任选其一贯穿:①实时订单总线 ②秒杀削峰 ③CDC 实时数仓(MySQL → Debezium → Kafka → Streams → ClickHouse)④用户行为埋点分析 ⑤IoT 上报平台完整项目源码 + docker-compose + 压测报告 + 监控接入综合项目全链路看板:QPS / Lag / ISR / 端到端延迟实时刷新58000
A附录 A · 命令 / 参数 / JMX 速查表kafka-*.sh 全集;Broker / Producer / Consumer 关键参数;重要 JMX 指标;内部 Topic 速查3500
B附录 B · Kafka vs RabbitMQ / RocketMQ / Pulsar / Redis Stream / NATS存储模型 / 消费模型 / 顺序保证 / 事务 / 延迟 / 吞吐 / 生态 7 个维度对比表3500
C附录 C · 踩坑案例集第 20 章的精简索引 + 常见生产事故复盘3000

合计:21 章正文 ≈ 12 万字,3 个附录 ≈ 1 万字,总规模 ≈ 13 万字 + 50+ 个 Python / Java 脚本 + 21 个交互式 demo.html。


二、章节深度分级

参考 task.md 第八节「写作执行顺序」,本教程将章节分为三档投入:

🔥 灵魂章节(重点打磨,篇幅 / 动画 / 案例都要顶配)

1、4、5、7、9、10、11、13 章 —— 它们对应 Kafka 的 8 个核心命题:「为什么快 / Producer / Consumer / 存储 / 副本 / KRaft / Rebalance / EOS」。这些章节即使读者其它章节略读,也必须能从这 8 章拼出 Kafka 的全貌。

🟢 主线章节(按节奏完整产出)

2、3、6、8、12、14、19、20 章 —— 提供完整知识 + 至少 1 个实战 + 1 个 demo + 5 题面试题,但篇幅不会刻意拉长。

🔵 生态章节(讲清「能干什么」+ 1 个端到端 demo)

15、16、17、18、21 章 —— 重点是让读者建立全景认知,不深挖每个 Connector / Stream 算子的源码。


三、学习路径建议

教程对不同读者给出 3 条推荐路线,按需挑选。

🚀 3 周入门路线(学完会用,覆盖 80% 日常场景)

阶段章节每天投入学完达成
第 1 周入门 + Producer/Consumer第 1、2、3、4、5 章1.5 h能写得出生产 / 消费代码、看得懂 Topic / Partition / Offset / Consumer Group
第 2 周设计 + 存储 + 高可用第 6、7、9 章 + 第 8 章选读1.5 h会算分区数 / 副本数 / min.insync.replicas,能解释「acks=all 为什么不丢」
第 3 周消费语义 + 综合第 12、20 章 + 第 21 章简化版1.5 h能避开「自动提交丢消息」「Rebalance 风暴」两大坑,能搭一个最小可用的事件总线

跳过:KRaft 内核(第 10)、EOS(第 13)、Streams / Connect / Schema Registry / 安全(第 15-18)。 后续遇到具体场景再回头补即可。

⚡ 1 个月进阶路线(学完懂原理,覆盖 95% 生产问题)

章节
第 1 周第 1-5 章(同入门路线)
第 2 周第 6、7、8、9 章(设计 + 存储 + 性能基石 + 副本)
第 3 周第 10、11、12、13 章(KRaft + Rebalance + 语义 + EOS)
第 4 周第 14、19、20、21 章(Compaction + 可观测性 + 踩坑 + 综合实战)

此路线产出的水平:能独立 Code Review Kafka 相关 PR、能在生产事故中定位 Lag / ISR / Rebalance 风暴根因、能为业务方制定 Topic 设计规范。

🎯 面试冲刺路线(2 周,覆盖大厂高频考点)

维度必看章节高频题数
架构总览第 1、2 章11 题
存储与吞吐第 7、8 章12 题
副本与 HA第 9 章7 题
KRaft / 元数据第 10 章6 题
Producer / Consumer / Rebalance第 4、5、11 章20 题
消息语义 / EOS第 12、13 章12 题
横向对比 MQ附录 B + 各章对比小框8 题
运维与踩坑第 19、20 章14 题

方法:先用一周通读这些章节末尾的「面试高频题」,第二周对照 interview.md 总索引,按主题分类做题、复述给别人听。


四、写作 / 阅读约定

为保持全教程上下一致:

项目约定说明
服务端版本Kafka 3.8+(KRaft 模式)本教程默认 KRaft,ZK 模式仅在第 10 章对比讲解
Python 客户端confluent-kafka ≥ 2.5librdkafka 实现,事务 / EOS / SASL 支持最全
Java 客户端kafka-clients 3.8.0关键章节给对照实现
配置项命名官方点分小写(如 bootstrap.serversenable.idempotenceacksPython 与 Java 配置键对照见第 1 章末附表
Topic 命名learn.<chapter>.<scene>(如 learn.01.hello不污染业务命名空间
默认 bootstraplocalhost:9092对应根目录 docker-compose.yml 暴露的 broker1
演示语言中文,技术名词保留英文Producer / Partition / Offset 等专有名词不翻译
代码块必标语言(```bash / ```python / ```java / ```yaml / ```mermaid / ```properties方便高亮与复制

五、配套文件总览

learnNote/kafka/
├── 0_learn_plan.md                # 本文档
├── docker-compose.yml             # 3 broker KRaft + Kafka UI + Schema Registry
├── requirements.txt               # Python 依赖
├── interview.md                   # 各章面试题总索引(边写边汇总)

├── 01_intro.md                    # ✅ 本批次产出
├── 01_intro/
│   ├── demo.html                  # ✅ 顺序写 vs 随机写 / MQ 对比 / Kafka 生态
│   └── code/hello_kafka.py        # ✅ 第一段可运行代码

├── 02_architecture.md             # ✅ 本批次产出
├── 02_architecture/
│   ├── demo.html                  # ✅ 集群拓扑可调 + 故障切换动画
│   └── code/cluster_probe.py      # ✅ AdminClient 探针

├── 03_cli_basics.md               # 由其它 Agent 产出
├── ...
├── 21_project.md

├── appendix_cheatsheet.md         # 附录 A
├── appendix_kafka_vs_others.md    # 附录 B
└── appendix_pitfalls.md           # 附录 C

本批次(基础设施 + 第 1-2 章)由本 Agent 完成,第 3-21 章及附录由并行的另外 7 个 Agent 同步推进。所有章节统一遵守上述约定。


六、给读者的建议

  1. 先跑环境再读文档:进 learnNote/kafka/ 目录执行 docker compose up -d,浏览器开 http://localhost:8080 确认 Kafka UI 能列出 3 个 broker,再开始第 1 章。
  2. demo.html 不要跳过:每个章节配套的可视化演示是为「打通直觉」设计的,比文字快 10 倍。
  3. 手敲一遍命令kafka-topics.sh / kafka-console-consumer.sh 这类命令亲手敲一次,比看 100 行文档管用。
  4. 横向对比记忆:每读完一章问自己「RabbitMQ / RocketMQ / Pulsar 这里是怎么做的?」—— Kafka 的设计选择只有放在对比里才能被真正理解。
  5. 遇到问题查附录 C:90% 的「Kafka 怎么这样」都能在踩坑案例集找到答案。

祝你早日成为团队里说得清「Kafka 为什么快、ISR 怎么算、EOS 怎么实现」的那个人。