主题
从 0 到 1 学习 Kafka · 总学习计划
本计划由
task.md第六节大纲展开而来,是后续 21 章 + 3 个附录的「施工图」。 配套环境:根目录docker-compose.yml(3 节点 KRaft + Kafka UI + Schema Registry)+requirements.txt(Pythonconfluent-kafka等依赖)。
📌 教程总目标
完成本教程后,你应当具备以下能力:
- 会用:写得出生产 / 消费代码,会用
kafka-*.sh系列命令调试集群、查 Topic、看 Lag、重置 Offset。 - 懂原理:能用生活类比把「顺序写 / PageCache / 零拷贝 / ISR / HW / LEO / Leader Epoch / KRaft / Rebalance / Exactly Once / Compaction」讲清楚。
- 能调优:知道分区数 / 副本数 /
acks/linger.ms/batch.size/compression.type/min.insync.replicas怎么选;能从 JMX 指标定位瓶颈。 - 能答题:覆盖大厂后端 / 中间件 / 大数据方向 Kafka 高频面试题(与 RabbitMQ / RocketMQ / Pulsar 横向对比、消息语义、事务、KRaft 演进、踩坑案例)。
一、章节总表(21 章 + 3 个附录)
字数为文档主体的预估值,不含代码 / demo.html。「demo 要点」是同名子目录
XX_xxx/demo.html必须实现的可视化交互。
| # | 章节 | 核心知识点 | 实战案例 / 代码 | demo 要点 | 面试题数 | 预计字数 |
|---|---|---|---|---|---|---|
| 1 | Kafka 是什么 & 为什么快 | MQ 价值(解耦/异步/削峰/日志总线)、Kafka 诞生(LinkedIn)、设计哲学(以日志为中心)、五大性能基石(顺序写 / PageCache / 零拷贝 / 批量 / 压缩)、与 RabbitMQ / RocketMQ / Pulsar / Redis Stream / NATS 对比 | hello_kafka.py:建 Topic + 生产 5 条 + 消费打印 | ① 顺序写 vs 随机写吞吐柱状图;② 5 大 MQ 存储/消费模型卡片;③ Kafka 生态全景拓扑 | 6 | 4500 |
| 2 | 核心概念与架构总览 | Producer / Broker / Consumer / Topic / Partition / Replica / Offset / Consumer Group / Coordinator / Controller / KRaft 全部术语;三横三纵分层;内部 Topic(__consumer_offsets / __transaction_state / __cluster_metadata) | cluster_probe.py:用 AdminClient 打印 Broker / Topic / Leader / ISR / Cluster ID | 交互式集群拓扑:可调 Broker / Partition / Replica,模拟 Broker 宕机看 Leader 切换 | 5 | 4500 |
| 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:最小可运行示例 | 命令速查交互卡片:点命令显示真实输入 / 输出对照 | 5 | 4000 |
| 4 | Producer 深入 | 同步 / 异步 send;acks=0/1/-1 语义;retries / delivery.timeout.ms / max.in.flight.requests.per.connection;linger.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 / 压缩算法滑块 → 实时算吞吐 / 延迟曲线 | 7 | 6500 |
| 5 | Consumer 深入 | subscribe vs assign;poll 循环模型;Offset 提交(自动 / 同步 / 异步 / 手动);__consumer_offsets;auto.offset.reset;心跳 / session.timeout.ms / max.poll.interval.ms;Consumer Group ↔ Coordinator 6 步协议;seek / pause / resume | manual_commit.py:业务主键 + 数据库去重;seek_demo.py:按时间戳回溯消费 | poll 循环时序图 + Offset 提交时机决策树 | 6 | 6500 |
| 6 | Topic 设计与分区策略 | 分区数选型(吞吐 / 并行度 / 有序性 / 扩容代价);副本数选型;Key 与顺序保证 / 热点;min.insync.replicas;何时「重建 Topic」而非「加分区」 | partition_calc.py:根据目标吞吐推荐分区数;hot_key_demo.py:复现热点分区 | 多车道高速公路动画:单分区 vs 多分区的顺序差异 + Key Hash 路由可视化 | 5 | 5000 |
| 7 | 存储与日志格式 | 日志目录 log.dirs;Topic-Partition-Segment 三层;.log / .index / .timeindex / .snapshot / leader-epoch-checkpoint;index.interval.bytes 稀疏索引;V0/V1/V2 消息格式与 RecordBatch;kafka-dump-log.sh 解码;segment 滚动条件 | dump_segment.sh 封装;造数 + 滚动 + 解码完整流水线 | .log / .index / .timeindex 二分查找定位 Offset 动画 | 6 | 6500 |
| 8 | 高吞吐的底层原理 | 顺序写 vs 随机写量级差;PageCache(为什么不在 JVM 堆里缓存);零拷贝 sendfile / transferTo(4 拷贝 → 2 拷贝);批量 + 压缩三端协同;mmap 在索引文件上的应用 | bench_seq_vs_rand.py:用 fio / Python 复现顺序 vs 随机写差距 | 4 拷贝 vs 2 拷贝动画 + PageCache 命中率仪表盘 | 6 | 6500 |
| 9 | 副本机制与 ISR | Leader / 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 截断防护对比 | 7 | 7000 |
| 10 | Controller 与 KRaft | ZK 时代 Controller 选举与元数据广播 → 历史包袱;KRaft 诞生(Raft 化的 __cluster_metadata Topic、Controller Quorum、事件溯源);KRaft 选举 / 快照 / 恢复;ZK → KRaft 迁移路径;4.0 全面移除 ZK | kraft_election_probe.py:从 metadata.log 读 Raft 选举记录;metadata_quorum_status.sh 解读 | KRaft 选举动画:3 controller 投票 / 日志复制 / 快照 | 6 | 6500 |
| 11 | 消费者组与 Rebalance | Group Coordinator 角色;JoinGroup / SyncGroup / Heartbeat / LeaveGroup;分配策略(Range / RoundRobin / Sticky / CooperativeSticky);Eager Rebalance 的 STW 之痛;Static Membership;Rebalance 监听器;Lag 排查 | rebalance_observer.py:监听 ConsumerRebalanceListener;static_membership_demo.py | 4 种策略动画对比:增减 Consumer / Topic 看分区如何重分配 | 7 | 7000 |
| 12 | Offset 与「消息语义」 | At Most Once / At Least Once / Exactly Once 工程含义;自动提交陷阱;手动提交正确姿势;消费者端幂等模式(业务主键 / Redis SETNX / DB 唯一约束);seek 到 Offset / 时间戳;__consumer_offsets 存储结构 | idempotent_consumer.py:业务表去重;seek_to_timestamp.py | 三种语义对比时序图 + 自动提交 vs 手动提交风险点 | 5 | 5500 |
| 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 Leader | 7 | 7500 |
| 14 | Log Compaction 与 Tombstone | delete / compact / compact,delete;min.cleanable.dirty.ratio / segment.ms;null value 作 Tombstone;Compaction 在状态快照 / CDC / 内部 Topic 的应用;线程模型与性能影响 | compaction_demo.py:构造 100 万 key + 反复 update 看压缩;tombstone_delete.py | 同 Key 多版本压缩动画 + Tombstone 删除流程 | 5 | 5000 |
| 15 | 安全与多租户 | SSL / SASL(PLAIN / SCRAM / GSSAPI / OAUTHBEARER);ACL(Topic / Group / Cluster / TransactionalId);Quota(Producer / Consumer / Request);SuperUser 与审计;最小权限模型 | setup_scram_user.sh;acl_demo.sh;带 SASL 的 Python Producer | 权限矩阵交互查询 + Quota 限速效果曲线 | 5 | 5000 |
| 16 | Kafka Connect | Source / Sink Connector;Standalone vs Distributed;常见 Connector(JDBC / Debezium / Elasticsearch / S3 / MongoDB / HDFS);Converter(JSON / Avro / Protobuf);SMT;DLQ;Exactly Once Source | Debezium MySQL → Kafka → JDBC Sink → PG 完整链路;connect_rest_calls.py | Connect 集群拓扑 + Connector 状态可视化 + DLQ 流转动画 | 5 | 5500 |
| 17 | Kafka 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 种窗口聚合演示 | 6 | 7000 |
| 18 | Schema Registry 与数据治理 | 为什么要 Schema;Avro / Protobuf / JSON Schema 三家对比;兼容性策略(BACKWARD / FORWARD / FULL / NONE 及 _TRANSITIVE);_schemas Topic;CDC 表结构变更应对 | avro_producer.py + avro_consumer.py;schema_evolution_demo.py | 6 种兼容性策略测试矩阵:加字段 / 删字段 / 改类型 → 自动判定 | 5 | 5500 |
| 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.sh;reassign_partitions.py;模拟 Grafana 面板的 Python 仪表盘 | 模拟 Grafana 看板 + 异常注入触发指标变化 | 6 | 6500 |
| 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 抖动 | 每个踩坑都有可复现脚本 + 修复后对照 | 「踩坑还原器」交互菜单:选案例 → 看症状 → 一键复现 → 看根因 | 8 | 7500 |
| 21 | 综合实战项目 | 任选其一贯穿:①实时订单总线 ②秒杀削峰 ③CDC 实时数仓(MySQL → Debezium → Kafka → Streams → ClickHouse)④用户行为埋点分析 ⑤IoT 上报平台 | 完整项目源码 + docker-compose + 压测报告 + 监控接入 | 综合项目全链路看板:QPS / Lag / ISR / 端到端延迟实时刷新 | 5 | 8000 |
| 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.5 | librdkafka 实现,事务 / EOS / SASL 支持最全 |
| Java 客户端 | kafka-clients 3.8.0 | 关键章节给对照实现 |
| 配置项命名 | 官方点分小写(如 bootstrap.servers、enable.idempotence、acks) | Python 与 Java 配置键对照见第 1 章末附表 |
| Topic 命名 | learn.<chapter>.<scene>(如 learn.01.hello) | 不污染业务命名空间 |
| 默认 bootstrap | localhost: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 同步推进。所有章节统一遵守上述约定。
六、给读者的建议
- 先跑环境再读文档:进
learnNote/kafka/目录执行docker compose up -d,浏览器开 http://localhost:8080 确认 Kafka UI 能列出 3 个 broker,再开始第 1 章。 - demo.html 不要跳过:每个章节配套的可视化演示是为「打通直觉」设计的,比文字快 10 倍。
- 手敲一遍命令:
kafka-topics.sh/kafka-console-consumer.sh这类命令亲手敲一次,比看 100 行文档管用。 - 横向对比记忆:每读完一章问自己「RabbitMQ / RocketMQ / Pulsar 这里是怎么做的?」—— Kafka 的设计选择只有放在对比里才能被真正理解。
- 遇到问题查附录 C:90% 的「Kafka 怎么这样」都能在踩坑案例集找到答案。
祝你早日成为团队里说得清「Kafka 为什么快、ISR 怎么算、EOS 怎么实现」的那个人。