主题
第 1 章 Kafka 是什么 & 为什么这么快
学习目标:能用一句话向产品经理解释「Kafka 到底是个什么玩意儿」;能向面试官讲清楚「Kafka 凭什么吞吐能比 RabbitMQ 高一个数量级」;能跑通第一段 Python 代码,连上一个 KRaft 模式的 Kafka 集群,把 5 条消息从生产者端送到消费者端。
0. 全书约定(先把规矩定下来)
为了让全教程上下一致,从本章开始遵守如下约定:
| 项目 | 约定 | 说明 |
|---|---|---|
| 服务端版本 | Apache Kafka 3.8 + KRaft 模式 | 后续所有命令、参数、行为都以 3.8 / KRaft 为准。ZK 模式只在第 10 章对比讲解。 |
| 命令行 | kafka-*.sh(官方自带) | kafka-topics.sh / kafka-console-producer.sh / kafka-consumer-groups.sh / kafka-dump-log.sh 等。 |
| 编程语言客户端 | Python confluent-kafka ≥ 2.5 | 基于 librdkafka,性能与协议完整度(事务 / EOS / SASL)最好。关键章节会补 Java kafka-clients 对照。 |
| Topic 命名 | learn.<chapter>.<scene> | 例:learn.01.hello、learn.04.orders。不污染业务命名空间。 |
| 默认连接 | bootstrap.servers=localhost:9092 | 对应根目录 docker-compose.yml 暴露的 broker1。 |
| 配置项写法 | 官方点分小写 | bootstrap.servers、enable.idempotence、acks=all,不缩写不大小写混写。Python 与 Java 的配置键对照见 §1.10。 |
| 可视化 | Kafka UI(http://localhost:8080) | 浏览器看 Topic / 消费组 / Lag 比命令行直观。 |
📌 与 MySQL / PostgreSQL 的最大区别:你写完 SQL 后通常立刻就能查到结果;而 Kafka 是异步、流式、可重放的 —— 生产者把消息塞进 Kafka 就走,消费者按自己的节奏来取。这不是 bug,是 Kafka 的灵魂。
1.1 一段「家谱」:Kafka 从哪来
1.1.1 消息中间件简史(30 秒版本)
1983 ─── IBM MQ Series(消息队列鼻祖,企业级)
1995 ─── 开源 ActiveMQ 的前身 OpenJMS
2001 ─── JMS 规范发布(Java Message Service)
2007 ─── RabbitMQ 1.0(基于 Erlang,AMQP 协议)
2010 ─── ★ Kafka 0.6 在 LinkedIn 内部诞生(Jay Kreps、Neha Narkhede、Jun Rao)
2011 ─── Kafka 开源、捐给 Apache
2012 ─── Apache RocketMQ(阿里巴巴)
2016 ─── Apache Pulsar(雅虎,存算分离)
2017 ─── Kafka 1.0
2021 ─── Kafka 2.8 引入 KRaft 预览(开始去 ZooKeeper)
2022 ─── Kafka 3.3 KRaft 生产可用
2024 ─── Kafka 3.7 / 3.8(KRaft 全面成熟)
2025 ─── Kafka 4.0(彻底移除 ZooKeeper 支持)Kafka 的诞生有一个非常具体的痛点:LinkedIn 的活动日志聚合。当时 LinkedIn 的服务每秒产生海量用户行为数据(搜索、点击、推荐曝光),这些数据要分发给:
- 离线数仓(Hadoop) —— 跑日报、训练推荐模型
- 实时监控(Storm 等) —— 安全风控、运维告警
- 搜索索引(Elasticsearch 等) —— 准实时索引更新
- 业务中台 —— 用户画像、广告效果分析
如果用传统的「点对点 RPC」或「JMS 队列」做法:
┌──────────┐
┌───────────►│ 数仓 │
│ └──────────┘
│ ┌──────────┐
┌────────┐ HTTP / RPC ├───────────►│ 风控监控 │
│ 上游 │ ────────────┤ └──────────┘
│ 服务 │ │ ┌──────────┐
└────────┘ ├───────────►│ 搜索索引 │
│ └──────────┘
│ ┌──────────┐
└───────────►│ 画像中台 │
└──────────┘
M 个上游 × N 个下游 = M × N 条点对点连接,加一个新下游就要改所有上游。这就是著名的 N×M 网状耦合问题。
LinkedIn 的工程师想要的是这样:
┌──────┐ ┌──────────┐
│ 上游1│─┐ ┌──►│ 数仓 │
├──────┤ │ ╔══════════════╗ │ ├──────────┤
│ 上游2│─┼─────►║ Kafka 总线 ║───┼──►│ 风控监控 │
├──────┤ │ ║(可重放日志) ║ │ ├──────────┤
│ 上游3│─┘ ╚══════════════╝ ├──►│ 搜索索引 │
└──────┘ └──►│ 画像中台 │
└──────────┘
上游只需「写一份到 Kafka」;新增下游只是再起一个消费者,
不需要改上游、不需要扩容上游、不需要协调发布窗口。也就是「一份数据多次消费」。但已有的消息队列(ActiveMQ、JMS)有两个致命缺陷:
- 吞吐撑不住:消息一旦被消费就丢,而 LinkedIn 需要保留 7 天供任意下游回溯。
- 不能多订阅:JMS 的 Topic 不能让消费者「重置位置回看历史」,活动日志这种场景没法用。
所以他们决定自己写一个 —— 「底座是磁盘上的不可变追加日志(commit log),上面才有消息队列的语义」。这就是 Kafka 的灵魂,也是 Jay Kreps 那句名言:
"Kafka is a distributed, partitioned, replicated commit log service." —— Kafka 不是消息队列,它是一个分布式的提交日志服务。
1.1.2 名字的故事
Jay Kreps 在博客里解释:「我喜欢 Franz Kafka(卡夫卡)的小说,而且这个系统主要负责处理写入数据,所以用一个作家的名字挺合适。」就这么随性。
1.2 消息中间件的价值:解耦 / 异步 / 削峰 / 日志总线
在讲「Kafka 为什么这么牛」之前,先讲「为什么我们需要消息中间件」。这是面试官经常用来探底的开场题,请用以下 4 个生活化场景回答。
1.2.1 价值一:解耦(用「邮局」类比)
没有邮局:
小张 ───亲手送信───► 小李 (小李搬家就找不到了)
小张 ───亲手送信───► 小王 (小王换电话就送不出去)
小张 ───亲手送信───► 小赵 (加 1 个收件人,小张要多跑一趟)
有了邮局:
小张 ──寄信──► ╔═══ 邮局信箱柜 ═══╗ ──取信──► 小李 / 小王 / 小赵
║ 邮筒 + 流水号 ║
╚═══════════════╝
小张只管把信投到对应的「信箱柜」(Topic),收件人变更不影响他。Topic 就是邮局的「信件分类柜」 —— 订单事件柜 / 用户行为柜 / 支付通知柜。生产者只关心把信投对柜子,消费者只关心从对应柜子取信,双方互不知道对方存在。这就是「解耦」。
1.2.2 价值二:异步(用「快递寄送」类比)
同步调用 异步消息
───────── ─────────
你 → 直接送货上门 → 收件人 你 → 快递柜 → 走人
↑ 必须等到对方签收才能离开 ↑ 一秒就走,快递员后续派送
↑ 对方不在你只能干等 ↑ 对方有空再去取下单后要做的事可能有:扣库存、发短信、发优惠券、推荐召回、日志归档…… 同步串行做完用户可能等 2 秒。异步把这些任务塞进 Kafka,主流程 50ms 返回,下游慢慢消费。
1.2.3 价值三:削峰(用「水库」类比)
直连模式:流量洪峰直接砸到下游
▲ QPS
│ ╱╲╱╲ 下游 DB 扛不住
│ ╱╲╱ ╲ ↓ 直接挂
│ ╱╲╱ ╲___ ✗
└──────────────────────►时间
水库模式(Kafka 削峰):
▲ QPS
│ ╱╲╱╲
│ ╱╲╱ ╲ Kafka 把洪峰存下来
│ ╱╲╱ ╲___ ↓ 下游按自己速率消费
└─[Kafka]──────────────► ──── ✓ 平稳
│
▼ 下游恒定 5K QPS 慢慢拉经典场景:秒杀。前端 100 万人同时点抢购,库存服务只能扛 1 万 QPS。把请求先写到 Kafka,后端按 1 万 QPS 拉队列扣库存 —— 用户看到的是「排队中」,而不是「服务挂了」。
1.2.4 价值四:日志总线(用「广播电台回放」类比)
传统点对点队列:消息一旦消费就消失(像电话)
生产 ───► [消息] ───► 消费A 取走 → 消息没了
消费B 也想看?对不起,没了。
Kafka 日志总线:消息按时间归档保留,谁都可以反复来听
生产 ───► ░░░░░░░░░░░░░░░░░░░░░ (磁盘日志,按 Offset 编号)
▲ ▲ ▲
│ │ │
消费A 消费B 消费C (各自记着自己读到哪了,互不干扰)
新消费可以「从头开始重放」整段历史这就是 Kafka 与 RabbitMQ 最本质的区别 —— 消息不会因为被消费而删除,只会因为「过期」(默认 7 天)或「Compaction」而清理。
📌 生活化记忆:RabbitMQ 像电话 —— 接听完就挂了;Kafka 像广播电台 + 录音存档 —— 现场听完磁带还留着,明天来还能重放。
1.3 Kafka 的设计哲学:以日志为中心
1.3.1 核心命题:「日志是分布式系统最简单也最强大的抽象」
Jay Kreps 在 2013 年写过一篇神文 《The Log: What every software engineer should know about real-time data's unifying abstraction》,核心观点只有一句:
「数据库的 redo log、复制协议的 WAL、消息队列、事件溯源、流处理 —— 它们的本质都是一份『不可变、有序、可追加』的日志。」
那么直接把这个「日志」做大、做分布式、做高吞吐,不就是一个万能基础设施了吗?
┌──────────────────────┐
│ 不可变追加日志 │ ← Kafka 的本体
│ Append-only Log │
└──────────┬───────────┘
│
┌─────────────┬────────────┼────────────┬─────────────┐
▼ ▼ ▼ ▼ ▼
消息队列 事件总线 CDC 同步 流处理状态 审计追溯
(MQ) (Event Bus) (Debezium) (Streams) (Audit)这就是 Kafka 「不是 MQ、是 Log Platform」 的核心。
1.3.2 Kafka 的「五大性能基石」(本章只做概览,第 7、8 章详讲)
读者第一次见到 Kafka 都会问:「为什么吞吐这么高?单机上百万 QPS?」答案是这 5 个底层选择协同作用:
| # | 基石 | 一句话原理 | 详讲章节 |
|---|---|---|---|
| 1 | 顺序写磁盘 | Producer 把消息只追加到日志末尾,机械盘顺序写就能跑到 600 MB/s,比内存随机访问还快 | 第 8 章 |
| 2 | PageCache | Broker 不在 JVM 堆里缓存消息,全部交给操作系统的 PageCache,避免 GC 同时享受 OS 预读 / 回写 | 第 8 章 |
| 3 | 零拷贝(sendfile) | Consumer 拉取消息时,数据从磁盘 → 内核缓冲 → 网卡,不经过用户态 JVM,省掉 2 次拷贝和 1 次上下文切换 | 第 8 章 |
| 4 | 批量化(Batch) | Producer 攒一波再发、Broker 攒一波再写、Consumer 一次拉一批,把单条消息开销摊到几百上千条上 | 第 4 章 |
| 5 | 压缩 | 在 Producer 端按批压缩(gzip / snappy / lz4 / zstd),网络与磁盘的有效吞吐再翻 3-10 倍 | 第 4 章 |
Producer 攒批 + 压缩 ──► 网络(少帧)──► Broker 顺序追加(PageCache)
│
▼
Consumer 批量拉取 ◄── 网络(少帧)◄── sendfile 零拷贝出页5 个机制不是孤立的,是一整套「面向吞吐」的工程哲学:从客户端到服务端、从内存到磁盘、从 CPU 到网卡,全链路都按「批量 + 顺序 + 零拷贝」来打磨。
1.3.3 第二个核心命题:「分区是一切的原子单位」
Topic = "learn.01.orders" (一个逻辑主题)
┌──────────┬──────────┬──────────┐
│ P0 │ P1 │ P2 │ ← 3 个分区(Partition)
│ ░░░░░░░░ │ ░░░░░░░░ │ ░░░░░░░░ │
│ ↑LEO │ ↑LEO │ ↑LEO │
└────┬─────┴────┬─────┴────┬─────┘
│ │ │
Leader Leader Leader
↓ 复制 ↓ ↓
Follower Follower Follower ← 副本(Replica)**分区(Partition)**几乎是 Kafka 所有重要概念的边界:
- 并行度的边界:N 个分区最多被 N 个消费者并行消费。
- 顺序的边界:单分区内严格有序,跨分区不保证。
- 副本的边界:副本以分区为粒度复制,不是以 Topic。
- Rebalance 的最小单元:消费者组重平衡也是按分区分配。
- 存储的物理单元:每个分区对应磁盘上一个目录,里面装一堆
.logSegment。
📌 记住一句话:「Kafka = Topic + Partition + 不可变追加日志 + 多副本 + 消费组」。后面所有章节都是在这五个词基础上展开。
1.4 与其它消息中间件的全面对比
这是面试高频题,必须能脱口而出。
1.4.1 五大消息中间件全景对比表
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar | Redis Stream / NATS |
|---|---|---|---|---|---|
| 诞生年份 | 2010(LinkedIn) | 2007(Rabbit Tech) | 2012(阿里) | 2016(Yahoo) | Redis 5.0 / 2010 |
| 协议 | 自研二进制 | AMQP 0.9.1 / STOMP / MQTT | 自研 | 自研(兼容 Kafka) | RESP / NATS |
| 语言 | Scala / Java | Erlang | Java | Java(Broker)+ C++ | C / Go |
| 存储模型 | 磁盘顺序日志 + 多副本 | 内存队列 + 持久化(lazy queue) | 磁盘日志(CommitLog 单文件) | 存算分离:Broker 计算 + BookKeeper 存储 | 内存(Stream 持久化为 RDB / AOF) |
| 消费模型 | 拉(pull)+ Offset 提交 | 推(push)+ ACK | 拉为主,支持长轮询 | 推 + ACK / 拉皆可 | 拉 |
| 顺序保证 | 单分区严格有序 | 单队列有序 | 单队列有序(顺序消费需指定 Queue) | 单分区有序 | Stream 内有序 |
| 路由 | 仅按 Key Hash 到 Partition | Exchange + Binding + Routing Key(最灵活) | Topic + Tag | Topic + Subscription | Stream Key |
| 吞吐(单机 4C8G 参考) | 百万级 QPS | 1~10 万 QPS | 10~50 万 QPS | 50~100 万 QPS | 10~50 万 QPS |
| 延迟(p99) | 5~50 ms | < 5 ms | < 10 ms | < 10 ms | < 1 ms(内存) |
| 是否可重放 | 是(按 Offset / 时间戳) | 否(消息消费即删) | 是(按 Offset) | 是 | 是(按 ID) |
| 事务 / EOS | 幂等 Producer + 事务 | 简单事务(性能差) | 事务消息(半消息 + 二次确认) | 事务(接近 Kafka) | 不支持 |
| 延迟 / 定时消息 | 不原生支持(需外部调度) | 原生(TTL + DLX) | 原生支持(多个延迟级别) | 原生 | 不支持 |
| 死信队列 DLQ | 第三方 / 自建 | 原生 | 原生 | 原生 | 自建 |
| 多租户 / 命名空间 | ACL + Quota | vhost | Topic Group | Tenant / Namespace 三层 | DB Index |
| 生态 | 最丰富(Connect / Streams / ksqlDB / Schema Registry) | 中等 | 中等(中文文档好) | 增长中 | Redis 生态 / 轻量 |
| 典型场景 | 大数据日志总线 / 流处理 / CDC / 事件溯源 | 业务异步 / RPC 风格队列 | 电商交易(订单 / 支付) | 多租户云原生 / 跨地域 | 缓存里的轻量队列 / 微服务 |
| 架构特点 | Broker = 存储 + 计算 耦合 | 内存为主,集群方案弱 | NameServer + Broker(去 ZK) | Broker 无状态 + BookKeeper 存储 | 单进程内存 |
1.4.2 5 个最容易踩坑的差异
- Kafka 是拉模型,RabbitMQ 是推模型:拉模型让消费者按自己节奏(防止打挂自己),推模型反应快但消费者一旦慢下来 Broker 就要做「流控」很复杂。
- Kafka 路由能力弱:只能按 Key Hash,做不了「订单金额 > 1000 走 VIP 队列」这种规则。要么客户端自己路由,要么用 Streams 做拆流。这不是 bug,是为了换吞吐做的取舍。
- Kafka 没有原生延迟消息:要做「30 分钟未支付自动取消」这种业务,得借助外部调度器或 Confluent 的 KIP-815。RocketMQ 一行配置搞定。
- Pulsar 不是 Kafka 的替代品:Pulsar 强在「存算分离 + 多租户」,适合云厂商;Kafka 强在「生态 + 简单」,适合自建。
- Redis Stream 不是「Kafka 替代品」:Redis 内存够,单机 50 万 QPS 没问题,但没有副本同步语义(Cluster 模式下跨槽问题多),扛不住 Kafka 那种 PB 级日志。
📌 一句金句送给面试官:「RabbitMQ 是消息队列,Kafka 是分布式日志。前者关注消息能不能可靠送达,后者关注海量数据能不能高吞吐持久化并被多次消费。它们解决的是不同问题。」
1.4.3 选型决策树
1.5 Kafka 的典型场景
1.5.1 LinkedIn 同款 —— 活动日志聚合
┌───────────┐ ┌───────────┐ ┌───────────┐
│ Web 服务 │ │ App 服务 │ │ 推荐服务 │
└─────┬─────┘ └─────┬─────┘ └─────┬─────┘
│ user_click │ search │ recommend
▼ ▼ ▼
╔═══════════════════════════════════════════╗
║ Kafka Topic: learn.01.user_events ║
║ Partition × 32 / Replica × 3 ║
╚════════════════╤══════════════════════════╝
│
┌────────────┼────────────┬─────────────┐
▼ ▼ ▼ ▼
实时风控 实时推荐召回 数仓 ETL A/B 实验平台
(Flink) (Streams) (Spark) (ksqlDB)1.5.2 微服务事件总线 / 事件溯源(Event Sourcing)
订单服务发出 OrderCreated / OrderPaid / OrderShipped 事件,所有下游(库存 / 物流 / 通知 / 推荐 / 财务)订阅 Kafka,任何下游崩了重启后都能从上次 Offset 继续,业务流不丢。
1.5.3 CDC 实时数据同步
MySQL ──Debezium──► Kafka ──Connect/Streams──► ClickHouse / ES / Doris
(binlog) (Topic per Table) (实时数仓)MySQL → Kafka 这一步几乎成了所有「实时数仓 / 数据湖」的事实标准。
1.5.4 IoT / 监控数据采集
百万级设备每秒上报心跳 + 指标 → 写到 Kafka 的 100+ Partition Topic → 规则引擎实时告警 + 时序库存档。Kafka 的高吞吐 + 顺序写在这里发挥到极致。
1.5.5 流处理平台底座
Kafka Streams / Flink / Spark Streaming 全部把 Kafka 当作输入 + 输出 + 状态后端的 Changelog Topic。没有 Kafka 就没有现代实时流处理生态。
1.6 实操:Hello Kafka
1.6.1 启动集群(30 秒)
进入 learnNote/kafka/ 目录:
bash
$ docker compose up -d
[+] Running 6/6
✔ Container kafka1 Healthy 28.4s
✔ Container kafka2 Healthy 28.6s
✔ Container kafka3 Healthy 28.8s
✔ Container schema-registry Healthy 42.1s
✔ Container kafka-ui Healthy 45.3s打开浏览器 http://localhost:8080,应能看到:
┌─────────────────────────────────────────────────────────────┐
│ Kafka UI · learn-kafka │
│ ── Brokers ───────────────────────────────────────────── │
│ ID HOST PORT STATUS │
│ 1 kafka1 19092 ✓ online │
│ 2 kafka2 19094 ✓ online │
│ 3 kafka3 19096 ✓ online │
└─────────────────────────────────────────────────────────────┘1.6.2 用命令行收发第一条消息(60 秒)
bash
# 1) 进入 broker1 容器
$ docker exec -it kafka1 bash
# 2) 创建 Topic:3 分区 / 3 副本
$ /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:19092 \
--create --topic learn.01.hello \
--partitions 3 --replication-factor 3
Created topic learn.01.hello.
# 3) 看一眼
$ /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:19092 \
--describe --topic learn.01.hello
Topic: learn.01.hello TopicId: aB3xY... PartitionCount: 3 ReplicationFactor: 3
Topic: learn.01.hello Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
Topic: learn.01.hello Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
Topic: learn.01.hello Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2
# 4) 启动一个 console-producer(保持窗口开着)
$ /opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server localhost:19092 --topic learn.01.hello
> hello kafka
> 这是第二条
> ^C
# 5) 另开一个终端,启动 console-consumer(从头消费)
$ docker exec -it kafka2 bash
$ /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:19094 --topic learn.01.hello \
--from-beginning
hello kafka
这是第二条恭喜,你刚刚跑通了 Kafka 的最小可用闭环:生产 → 持久化 → 多副本 → 消费。
1.6.3 用 Python 客户端
完整脚本见 01_intro/code/hello_kafka.py。核心片段:
python
from confluent_kafka import Producer, Consumer
from confluent_kafka.admin import AdminClient, NewTopic
BOOTSTRAP = "localhost:9092"
TOPIC = "learn.01.hello"
# 1) 建 Topic(幂等:如果已存在直接跳过)
admin = AdminClient({"bootstrap.servers": BOOTSTRAP})
admin.create_topics([NewTopic(TOPIC, num_partitions=3, replication_factor=3)])
# 2) 生产 5 条
producer = Producer({"bootstrap.servers": BOOTSTRAP})
for i in range(5):
producer.produce(TOPIC, key=str(i), value=f"hello-{i}")
producer.flush()
# 3) 消费打印
consumer = Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": "learn-01",
"auto.offset.reset": "earliest",
})
consumer.subscribe([TOPIC])
for _ in range(5):
msg = consumer.poll(5.0)
print(msg.partition(), msg.offset(), msg.value().decode())输出:
0 0 hello-1
0 1 hello-4
1 0 hello-2
2 0 hello-0
2 1 hello-3📌 观察:5 条消息按 Key Hash 散到 3 个分区,每个分区内严格按 Offset 递增。这就是 Kafka 顺序保证的边界 —— 单分区。
1.6.4 浏览器可视化演示
打开 01_intro/demo.html,可以点击按钮一步步看:
- ① 「为什么 Kafka 快」可视化 —— 顺序写 vs 随机写吞吐柱状图、零拷贝开 / 关对比、PageCache 命中率仪表
- ② 5 大消息队列对比 —— 存储模型 / 消费模型 / 顺序保证 / 路由能力 / 典型场景卡片切换
- ③ Kafka 生态全景拓扑 —— 鼠标悬停看每个组件(Connect / Streams / Schema Registry / Mirror Maker 等)的作用与替代品
1.7 五大性能基石的初步直觉(深入留给第 7、8 章)
1.7.1 顺序写 vs 随机写:一个「反直觉」的事实
HDD 机械盘:
随机写 4KB : ~100 ops/s ⇒ 约 400 KB/s
顺序写 1MB : ~600 ops/s ⇒ 约 600 MB/s ← 1500 倍差距!
SSD:
随机写 4KB : 约 100 MB/s
顺序写 1MB : 约 3 GB/s ← 仍然 30 倍差距为什么差这么多? 机械盘要寻道(10ms 级),SSD 要做闪存擦除 + GC。而顺序写完全规避这两件事。Kafka 强制所有消息只能 append 到日志末尾 —— 这一个设计决定,让磁盘吞吐可以接近网络带宽。
1.7.2 零拷贝(sendfile):从 4 次拷贝变 2 次
传统 read + write(消费消息 = 读盘 + 发网卡):
磁盘 ───DMA───► 内核 PageCache ───CPU 拷贝───► 用户态 buffer
│
│ CPU 拷贝
▼
网卡 ◄───DMA──── Socket buffer ◄─────────────── 用户态 buffer
总计:4 次拷贝 + 2 次用户态/内核态切换
零拷贝 sendfile(Kafka 消费走这条):
磁盘 ───DMA───► 内核 PageCache ───CPU/sg-DMA───► 网卡
(只是把元数据传给网卡)
总计:2 次拷贝 + 0~1 次切换「零拷贝」不是真的 0 次拷贝,是 0 次「用户态拷贝」。意义在于:
- 内存带宽节省 50%+
- CPU 几乎不参与 → 消费侧 broker CPU 利用率极低
- 这就是 Kafka 一台服务器能服务 10 万消费者拉取的关键
1.7.3 PageCache:把缓存交给 OS
Kafka 故意不在 JVM 堆里缓存消息。原因:
| 自己缓存(如 RabbitMQ) | 交给 PageCache(Kafka) |
|---|---|
| GC 压力大 | JVM 堆只放元数据,几百 MB 够用 |
| 进程崩溃缓存丢失 | OS 级,进程重启缓存还在 |
| 没有预读 / 回写优化 | OS 自带智能预读 / 后写 |
| 内存利用率低(重复缓存) | 与读写共用一份 |
📌 记住:Kafka Broker 给 JVM 设 6~8 GB 已经富裕,剩下的内存全留给 PageCache。生产经验值:机器内存 = JVM 堆 × 4~8。
1.7.4 批量化:把每条消息开销摊薄
Producer 的 linger.ms + batch.size 让一组消息攒在一起以一个 TCP 包发出。代价是延迟稍增(10ms 级),收益是吞吐提升 5~50 倍。Broker 接收、Consumer 拉取也都批量进行 —— 整条链路都按批走。
1.7.5 压缩:一次开销,多次受益
compression.type=zstd 在 Producer 端按批压缩一次,Broker 直接存压缩后的字节,Consumer 拉取时再解压。网络 + 磁盘吞吐再翻 3~10 倍,CPU 代价集中在 Producer / Consumer 两端。
1.8 Kafka 的整体架构(先看全景,第 2 章详讲)
ASCII 全景图(更直观):
┌─────────────────────────────────────────┐
│ 生产者集群 │
└─────┬──────────────┬──────────────┬──────┘
│ │ │
┌─────────────────▼──────────────▼──────────────▼─────────────────┐
│ │
│ Broker1 Broker2 Broker3 ← KRaft Controller Quorum │
│ ┌─────┐ ┌─────┐ ┌─────┐ │
│ │ P0L │ │ P0F │ │ P0F │ ← Topic A · Partition 0 │
│ │ P1F │ │ P1L │ │ P1F │ ← Topic A · Partition 1 │
│ │ P2F │ │ P2F │ │ P2L │ ← Topic A · Partition 2 │
│ └──┬──┘ └──┬──┘ └──┬──┘ │
│ │ │ │ L=Leader, F=Follower │
│ └──────────┴──────────┘ │
│ 副本同步(ISR) │
└────────────────────┬───────────────────────────────────────────────┘
│ fetch (拉模型)
┌────────────────┼────────────────┐
▼ ▼ ▼
Consumer Group A Consumer Group B Connect / Streams
(订单服务) (数仓 ETL) (实时同步)1.9 本章小结
┌────────────────────────────────────────────────────────────────┐
│ 第 1 章核心要点 │
├────────────────────────────────────────────────────────────────┤
│ │
│ ① Kafka 不是消息队列,是「分布式提交日志服务」 │
│ 生于 LinkedIn 2010,为解决 N×M 网状耦合 + 大数据日志总线 │
│ │
│ ② 消息中间件四大价值:解耦 / 异步 / 削峰 / 日志总线 │
│ 生活类比:邮局 / 快递 / 水库 / 广播电台 │
│ │
│ ③ 与其它 MQ 的关键差异: │
│ • 拉模型 + Offset 提交(不是推 + ACK) │
│ • 消息消费完不删(按时间 / Compaction 清理) │
│ • 路由弱(仅 Key Hash),换来吞吐 │
│ • 单分区有序 ≠ 全 Topic 有序 │
│ │
│ ④ 五大性能基石(深入第 7、8 章): │
│ 1. 顺序写磁盘(vs 随机写 1500×) │
│ 2. PageCache(不在 JVM 堆里缓存) │
│ 3. 零拷贝 sendfile(4 拷贝 → 2 拷贝) │
│ 4. 批量化(Producer / Broker / Consumer 三端) │
│ 5. 压缩(zstd / lz4,Producer 端一次) │
│ │
│ ⑤ 分区是一切的原子单位 —— 并行度 / 顺序 / 副本 / Rebalance 边界 │
│ │
│ ⑥ 默认环境:3 broker KRaft + Kafka UI + Schema Registry │
│ Topic 命名 learn.<chapter>.<scene>,bootstrap localhost:9092 │
│ │
└────────────────────────────────────────────────────────────────┘1.10 附:Python confluent-kafka 与 Java kafka-clients 配置键对照(高频版)
| 含义 | Python confluent-kafka | Java kafka-clients |
|---|---|---|
| 集群地址 | bootstrap.servers | bootstrap.servers |
| Producer ack 策略 | acks | acks |
| 启用幂等 Producer | enable.idempotence | enable.idempotence |
| 攒批等待 | linger.ms | linger.ms |
| 单批最大字节 | batch.size | batch.size |
| 压缩算法 | compression.type | compression.type |
| Consumer 组 | group.id | group.id |
| 起始位置策略 | auto.offset.reset | auto.offset.reset |
| 是否自动提交 | enable.auto.commit | enable.auto.commit |
| 心跳间隔 | heartbeat.interval.ms | heartbeat.interval.ms |
| 单次 poll 最大数 | max.poll.records (注: librdkafka 中通过 queued.max.messages.kbytes/fetch.message.max.bytes 控制) | max.poll.records |
| 客户端 ID | client.id | client.id |
📌 核心一致:
confluent-kafka是对 librdkafka 的 Python 封装,librdkafka 又是 Kafka 协议的 C 重写,配置键基本与 Java 一一对应。少数 librdkafka 特有的(如queued.max.messages.kbytes)会在用到的章节单独说明。
1.11 面试高频题
Q1:Kafka 和 RabbitMQ 的本质区别是什么?
考察点:是否真的理解两者解决的不是同一个问题,而不是只会背「Kafka 吞吐高 / RabbitMQ 延迟低」。
标准答案(按重要性分点):
- 设计定位完全不同:
- RabbitMQ 是消息队列(Queue),核心抽象是「队列里的消息被谁取走就消失了」,关注「可靠投递」。
- Kafka 是分布式提交日志(Distributed Commit Log),核心抽象是「磁盘上一份不可变的追加日志」,关注「海量数据持久化 + 多次重放」。
- 存储模型不同:RabbitMQ 内存为主 + 持久化补充;Kafka 磁盘顺序写 + PageCache 为核心。
- 消费模型不同:RabbitMQ 是 Broker 推(push) + 客户端 ACK;Kafka 是客户端拉(pull) + 自己提交 Offset。Kafka 拉模型让消费者按自己速率消费,永远不会被压垮。
- 路由能力不同:RabbitMQ 有 Exchange + Binding + Routing Key 三层灵活路由(topic / direct / fanout / headers);Kafka 只能按 Key Hash 到 Partition,路由极简,换吞吐。
- 顺序保证不同:RabbitMQ 单队列内有序;Kafka 单分区内有序,跨分区不保证。
- 吞吐量级不同:单机 RabbitMQ 1~10 万 QPS,Kafka 百万 QPS。
- 消息生命周期不同:RabbitMQ 消费完即删;Kafka 按时间(默认 7 天)或 Compaction 策略清理,消费完不删,可被任意消费组重复读。
加分项:能讲一句「RabbitMQ 是消息队列,Kafka 是日志平台」,并举例:要做 RPC 风格异步任务(订单 → 短信通知)选 RabbitMQ;要做大数据日志总线 / CDC / 流处理选 Kafka。
易错点:不要说「Kafka 性能更好」 —— 在「单条消息 + 复杂路由 + 低延迟」场景,RabbitMQ p99 延迟(< 5ms)甚至比 Kafka(5-50ms)更低。
Q2:为什么 Kafka 这么快?吞吐量为什么能比传统 MQ 高一个数量级?
考察点:「五大性能基石」必须能脱口而出,并能讲出每一条的「为什么」。
标准答案:5 个核心机制协同作用:
- 顺序写磁盘:Producer 把消息只追加到 Partition 日志末尾,机械盘顺序写吞吐能到 600 MB/s(vs 随机写 0.4 MB/s,差 1500 倍),SSD 上仍有 30 倍差距。Kafka 强制 append-only,从根本上规避了寻道 / GC 开销。
- PageCache(操作系统页缓存):Broker 不在 JVM 堆里缓存消息,全部交给 OS 的 PageCache 管理。好处:① 避免 JVM GC 压力;② 进程重启缓存不丢;③ 享受 OS 自带的智能预读 / 后写;④ 与读写共用一份内存,利用率极高。
- 零拷贝 sendfile:消费者拉消息时,Broker 用
sendfile系统调用直接把 PageCache 的数据 → 网卡,绕过用户态 JVM 堆。从传统的「4 次拷贝 + 2 次切换」减到「2 次拷贝 + 0~1 次切换」,CPU 几乎不参与,单台 Broker 能服务上万个 Consumer。 - 批量化:Producer 端
linger.ms+batch.size攒一波再发;Broker 端按批写盘;Consumer 端按批拉取。单条消息开销摊薄到几百上千条上,吞吐提升 5-50 倍。 - 压缩:Producer 按批用 lz4 / snappy / zstd 压缩一次,Broker 直接存压缩字节、Consumer 拉到再解压。网络 + 磁盘有效吞吐再翻 3-10 倍。
加分项:
- 能补充「Sticky Partitioner」(KIP-480):Producer 默认会把同一批消息贴到同一个分区,让批次更密,实际生产中比 Murmur2 Hash 吞吐高 30-50%。
- 能讲「MMap 应用在索引文件上」:
.index/.timeindex用 mmap 让查找像访问内存一样。 - 能对比「为什么 RabbitMQ 不能这样」:RabbitMQ 的 Erlang 进程内消息队列 + 复杂路由表 + Push 模型,注定无法把整链路按批 + 顺序 + 零拷贝来打磨。
易错点:
- 不要只说「Kafka 用了 PageCache 所以快」 —— 5 个机制是协同工作的,单拎一条都不能解释为什么快一个数量级。
- 不要把「零拷贝」说成「0 次拷贝」 —— 是0 次用户态拷贝,DMA 拷贝还是有的。
Q3:消息中间件解决了什么问题?为什么不直接用 RPC?
考察点:基础概念是否扎实,能否从架构角度解释「中间件存在的意义」。
标准答案:
消息中间件解决 4 类问题,都是 RPC(同步调用)解决不了的:
- 解耦:上游不需要知道下游有谁、在哪、活没活。新增一个下游消费者只需要订阅 Topic,上游不改一行代码。RPC 模式下 N 个上游 × M 个下游 = N×M 条耦合关系。
- 异步:上游写完消息立刻返回(毫秒级),不需要等下游处理。用户体验大幅提升(下单 → 立即提示成功,发短信 / 加积分 / 推荐召回慢慢做)。RPC 必须串行等所有下游响应。
- 削峰:流量洪峰先进 Kafka 这个「水库」,下游按自己稳定速率消费。保护下游不被瞬时高 QPS 打挂(秒杀场景必备)。
- 日志总线 / 数据集成:一份数据写一次到 Kafka,N 个下游(数仓 / 风控 / 推荐 / 监控)各自独立消费,避免上游被多次重复调用。这是 RPC 完全做不到的。
加分项:
- 能讲清「同步 RPC 的三个核心痛点」:① 强耦合;② 强一致依赖(一个下游挂了上游失败);③ 流量峰值雪崩。
- 能补一句「消息中间件的代价」:① 异步 → 链路追踪难;② 多了一个组件 → 多一份运维;③ 消费幂等设计成本(消费侧要扛 At-Least-Once 重复)。
- 能对比「RPC 框架(如 gRPC、Dubbo) + MQ 的搭配」:RPC 用于强一致同步链路,MQ 用于异步事件链路。
易错点:不要说「MQ 完全替代 RPC」 —— 强一致、低延迟、需要立即拿到响应的场景(如查余额、扣库存)必须 RPC。
Q4:Kafka 为什么用「拉模型」而不是「推模型」?相比 RabbitMQ 的推有什么优势?
考察点:拉 / 推模型的本质权衡,以及 Kafka 选择拉的工程原因。
标准答案:
Kafka 选择 Pull(拉)模型,每个 Consumer 主动向 Broker fetch 消息。RabbitMQ 是 Push(推)模型,Broker 主动把消息推给消费者。
Kafka 拉模型的 4 个优势:
- 消费者按自己速率消费,永远不会被压垮:消费慢就少拉一点,消费快就多拉一点。Push 模型下 Broker 不知道消费者能扛多少,必须做复杂的「流控 / 反压」(RabbitMQ 的
prefetch配置就是干这个的)。 - 天然支持批量拉取:一次
fetch拿一批,吞吐高。Push 是 Broker 一条一条推,批量化困难。 - 消费者宕机 / 恢复对 Broker 透明:Broker 不维护「消息发给谁了」状态,Consumer 自己维护 Offset。Push 模型 Broker 必须记 ACK 状态。
- 支持「重置消费位置 / 重放历史」:Consumer 改一下 Offset 就能从任意位置重新消费。Push 模型消息推完就走了,回放成本高。
Pull 模型的代价:
- 空轮询问题:没消息时 Consumer 还在不停问 Broker。Kafka 用
fetch.min.bytes+fetch.max.wait.ms做「长轮询」缓解 —— 没攒够数据就让 Broker 把请求 hold 住到超时。
加分项:
- 能讲「Kafka 实际上是 long-polling」:
fetch.max.wait.ms=500ms时,Broker 没数据会等到 500ms 或攒够fetch.min.bytes再返回。这样既不空转又不延迟太大。 - 能对比 Pulsar:Pulsar 同时支持 Push 和 Pull,Consumer 端可选;底层走 BookKeeper 拉,Broker 缓存后推给客户端。
易错点:不要说「Push 一定比 Pull 延迟低」 —— Kafka 在配置得当下端到端延迟可以做到 5ms 以内,与 Push 模型没有本质差距。
Q5:什么是 Topic / Partition / Offset?为什么 Kafka 单分区才有序?
考察点:核心概念是否真的理解,而不是只会背术语。
标准答案:
Topic = 消息的逻辑分类(邮局信件分类柜)。一个业务事件流对应一个 Topic,比如 learn.01.orders 装订单事件、learn.01.user_clicks 装用户点击事件。
Partition(分区) = Topic 内部的物理切分单元(同一类柜子下的多条流水线)。每个 Partition 在磁盘上是一个独立的日志目录,里面装一组 Segment 文件(.log / .index / .timeindex)。一个 Topic 通常有 N 个 Partition(N = 期望的并行度)。
Offset(偏移量) = 单个 Partition 内部的消息编号(每条流水线上的流水号)。从 0 开始单调递增,永远不重复。Consumer 通过提交 Offset 标记「我消费到哪儿了」,宕机重启后从这个 Offset 继续。
为什么单分区才有序?
Topic "orders"
├── Partition 0: [P0-0] [P0-1] [P0-2] [P0-3] ← 单分区内严格顺序写
├── Partition 1: [P1-0] [P1-1] [P1-2]
└── Partition 2: [P2-0] [P2-1] [P2-2] [P2-3] [P2-4]- 单分区内有序:Producer 同步追加 + Broker append-only 保证写入顺序,Consumer 顺序拉取保证消费顺序。
- 跨分区不保证:3 个分区 Leader 在 3 台机器,网络 / 磁盘 / GC 节奏不同,消息写入到达 Broker 的真实时刻不同步;Consumer 也按 fetch 节奏拉取,不可能整体有序。
实战建议:
- 想保证某个业务实体的事件有序(如同一笔订单的多个状态变更),就让它们用同一个 Key(如
order_id),Kafka 默认按Murmur2(key) % NHash 到固定分区,自然顺序。 - 千万不要把整个 Topic 设为 1 分区来求全局有序 —— 直接放弃了 Kafka 最大优势(并行度),单 broker 单线程顶天 1-2 万 QPS。
加分项:
- 能讲「Sticky Partitioner(KIP-480)」:从 Kafka 2.4 起,Producer 没指定 Key 时不再按轮询,而是「粘住」一个分区直到这一批攒满 —— 让批次更密,吞吐高 30-50%。
- 能讲「Compaction Topic 的 Offset 不连续」:被压缩掉的 key 留下空洞,Consumer 看到的 Offset 会跳。
易错点:
- 不要说「Offset 是全局唯一」 —— Offset 只在单 Partition 内唯一。两个不同 Partition 的 Offset 可以是同一个数字。
- 不要混淆 Topic 的副本数(Replica) 与 分区数(Partition):Replica 是同一份数据的多份拷贝(HA),Partition 是把数据切成多片(并行度)。
Q6:什么是消费者组(Consumer Group)?跟广播 / 订阅的区别是什么?
考察点:能否用「家庭信箱」类比讲清楚 Consumer Group 的语义。
标准答案:
Consumer Group(消费者组) 是 Kafka 实现「广播 / 单播」的核心机制:
Topic "orders" (3 个 Partition: P0 P1 P2)
│
├──► Group A (订单服务,3 个实例)
│ ├── Consumer A1 ── 消费 P0
│ ├── Consumer A2 ── 消费 P1
│ └── Consumer A3 ── 消费 P2 ← Group A 内部分摊
│
├──► Group B (数仓 ETL,2 个实例)
│ ├── Consumer B1 ── 消费 P0, P2
│ └── Consumer B2 ── 消费 P1 ← Group B 内部分摊
│
└──► Group C (审计归档,1 个实例)
└── Consumer C1 ── 消费 P0, P1, P2两条核心规则:
- 同一个 Group 内:每个 Partition 只能被 1 个 Consumer 消费 —— 单播 / 负载均衡。
- 不同 Group 之间:每个 Partition 被每个 Group 独立 消费一次 —— 广播 / 多订阅。
生活类比:Consumer Group 就像「一家三口共用一个邮箱」 —— 邮件来了任何一个人取走就行,不会被取两次(Group 内);但「张家、李家、王家」各自有邮箱,同一封广告信会发 3 份(Group 间)。
应用模式:
- 完全广播:每个消费者一个独立 Group。
- 完全单播 / 负载均衡:所有消费者同一个 Group。
- 分组广播 + 组内负载均衡:每个业务方一个 Group(最常用模式)。
加分项:
- 能讲「Consumer 数量 ≤ Partition 数量」 —— 多了的 Consumer 会闲置,因为分区分配不到。生产建议两者相等或 Consumer 数 = Partition 数 / 2。
- 能讲「Group Coordinator 在哪台 Broker 上」 —— 通过
__consumer_offsetsTopic 的 Partition 选举,hash(group.id) % 50决定。 - 能讲「Static Membership(KIP-345)」:通过
group.instance.id让 Consumer 重启不触发 Rebalance。
易错点:
- 不要说「Consumer Group 是 Kafka 服务端的概念」 —— Group 元数据是服务端管理的,但逻辑上 Group 是消费者侧的协议约定,Broker 只是协调器。
- 不要把
__consumer_offsets当作「Consumer 的状态存储」 —— 它只存 Offset 提交记录,Group 成员状态是 Coordinator 内存里的。
📌 下一章预告:第 2 章我们正式拆解 Kafka 集群 —— 把 Producer / Broker / Consumer / Topic / Partition / Replica / Coordinator / Controller / KRaft 全部术语「首次出现 → 解释 → 类比」一遍,并用一张大图讲清读 / 写链路。然后用
AdminClient探针打印我们 docker compose 起的集群拓扑。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
"""
第 1 章 · Hello Kafka —— 你与 Kafka 的「第一次握手」
完整闭环:
1. 用 AdminClient 创建 Topic `learn.01.hello`(3 partition × 3 replica)。
如果 Topic 已经存在则跳过(幂等)。
2. 用 Producer 同步生产 5 条消息(带 key,便于观察分区路由)。
3. 用 Consumer 从头消费,打印 partition / offset / key / value。
4. 程序自动结束(最多等 10 秒)。
前置条件:
cd learnNote/kafka && docker compose up -d # 启动集群
pip install -r requirements.txt # 安装依赖
运行:
python 01_intro/code/hello_kafka.py
预期输出(partition 顺序可能不同,offset 必定从 0 起):
[INFO ] Topic learn.01.hello 已就绪 (partitions=3, RF=3)
[SEND ] key=0 → partition=2 offset=0
[SEND ] key=1 → partition=0 offset=0
[SEND ] key=2 → partition=1 offset=0
[SEND ] key=3 → partition=2 offset=1
[SEND ] key=4 → partition=0 offset=1
[RECV ] partition=0 offset=0 key=1 value=hello-1
[RECV ] partition=0 offset=1 key=4 value=hello-4
...
[DONE ] 共消费 5 条消息
阅读建议:
- 关注「同 key 必定到同一分区」(Murmur2 Hash)
- 关注 partition 内部 offset 单调递增
- 关注 group.id + auto.offset.reset=earliest 的组合是怎么"从头消费"的
"""
from __future__ import annotations
import logging
import sys
import time
from typing import Optional
from confluent_kafka import Consumer, KafkaError, KafkaException, Producer
from confluent_kafka.admin import AdminClient, NewTopic
# ----------------------- 全局配置 -----------------------
BOOTSTRAP_SERVERS = "localhost:9092"
TOPIC_NAME = "learn.01.hello"
NUM_PARTITIONS = 3
REPLICATION_FACTOR = 3
MESSAGE_COUNT = 5
CONSUMER_GROUP = "learn.01.hello.group"
logging.basicConfig(
level=logging.INFO,
format="[%(levelname)-5s] %(message)s",
stream=sys.stdout,
)
log = logging.getLogger("hello_kafka")
# ----------------------- 1) 建 Topic -----------------------
def ensure_topic(admin: AdminClient, name: str, partitions: int, rf: int) -> None:
"""幂等创建 Topic:已存在则跳过。"""
existing = admin.list_topics(timeout=10).topics
if name in existing:
meta = existing[name]
log.info(
"Topic %s 已存在 (partitions=%d, partitions_meta=%s)",
name,
len(meta.partitions),
sorted(meta.partitions.keys()),
)
return
log.info("Topic %s 不存在,正在创建 ...", name)
new_topic = NewTopic(
topic=name,
num_partitions=partitions,
replication_factor=rf,
config={
# 教学用:保留 1 小时即可,方便反复重置
"retention.ms": str(60 * 60 * 1000),
# acks=all 时至少 2 副本同步成功才算写入成功
"min.insync.replicas": "2",
},
)
futures = admin.create_topics([new_topic])
# create_topics 返回 dict[str, Future],必须遍历 .result() 才会真正阻塞等结果
for topic, fut in futures.items():
try:
fut.result(timeout=15)
log.info("Topic %s 已就绪 (partitions=%d, RF=%d)", topic, partitions, rf)
except KafkaException as e:
# 并发创建时可能撞到 TOPIC_ALREADY_EXISTS,可视为成功
if e.args and e.args[0].code() == KafkaError.TOPIC_ALREADY_EXISTS:
log.info("Topic %s 已存在(并发创建撞车),继续", topic)
else:
raise
# ----------------------- 2) 生产消息 -----------------------
def delivery_callback(err: Optional[KafkaError], msg) -> None:
"""每条消息发送结果的异步回调(在 producer.poll/flush 时被触发)。"""
if err is not None:
log.error("发送失败: %s", err)
else:
log.info(
"[SEND ] key=%s → partition=%d offset=%d",
msg.key().decode() if msg.key() else "<none>",
msg.partition(),
msg.offset(),
)
def produce_messages(producer: Producer, topic: str, count: int) -> None:
log.info("开始生产 %d 条消息到 %s ...", count, topic)
for i in range(count):
# key 决定分区:confluent-kafka 默认 Murmur2 Hash 算法
# 同 key → 同 partition;这里能观察到 partition 路由分布
producer.produce(
topic=topic,
key=str(i).encode(),
value=f"hello-{i}".encode(),
on_delivery=delivery_callback,
)
# 每次 produce 后调用 poll 处理回调队列;不调用回调就不会触发
producer.poll(0)
# flush 阻塞直到所有消息都收到 ack 或 timeout(30s)
remaining = producer.flush(timeout=30)
if remaining > 0:
raise RuntimeError(f"还有 {remaining} 条消息未发送成功,请检查 broker 状态")
# ----------------------- 3) 消费消息 -----------------------
def consume_messages(consumer: Consumer, topic: str, expected: int) -> int:
"""从头消费 expected 条,打印后返回实际拿到的条数。"""
log.info("订阅 %s,开始消费(最多等待 10s)...", topic)
consumer.subscribe([topic])
received = 0
deadline = time.time() + 10.0
while received < expected and time.time() < deadline:
# poll(timeout) 拉取一条消息;timeout 单位秒
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
# _PARTITION_EOF 是「分区没新消息」的友好提示,不是错误
if msg.error().code() == KafkaError._PARTITION_EOF:
continue
log.error("消费错误: %s", msg.error())
continue
log.info(
"[RECV ] partition=%d offset=%d key=%s value=%s",
msg.partition(),
msg.offset(),
msg.key().decode() if msg.key() else "<none>",
msg.value().decode(),
)
received += 1
return received
# ----------------------- main -----------------------
def main() -> int:
log.info("Hello Kafka! bootstrap.servers = %s", BOOTSTRAP_SERVERS)
common = {"bootstrap.servers": BOOTSTRAP_SERVERS}
# ---- AdminClient ----
admin = AdminClient(common)
try:
ensure_topic(admin, TOPIC_NAME, NUM_PARTITIONS, REPLICATION_FACTOR)
except Exception as e:
log.error("建 Topic 失败:%s", e)
log.error("请确认 docker compose up -d 已启动,并能 ping 通 localhost:9092")
return 1
# ---- Producer ----
producer_conf = {
**common,
"client.id": "hello-kafka-producer",
"acks": "all", # 等所有 ISR 副本确认(最强一致)
"enable.idempotence": True, # 幂等 Producer:失败重试不会重复
"linger.ms": 5, # 攒批 5ms 提升吞吐
"compression.type": "lz4", # Producer 端压缩(教学使用 lz4 兼顾速度与压缩比)
"max.in.flight.requests.per.connection": 5,
}
producer = Producer(producer_conf)
try:
produce_messages(producer, TOPIC_NAME, MESSAGE_COUNT)
except Exception as e:
log.error("生产失败:%s", e)
return 2
# ---- Consumer ----
consumer_conf = {
**common,
"group.id": CONSUMER_GROUP,
"client.id": "hello-kafka-consumer",
"enable.auto.commit": False, # 教学环境关掉自动提交,主动控制
"auto.offset.reset": "earliest", # 没有 offset 时从最早开始
}
consumer = Consumer(consumer_conf)
try:
got = consume_messages(consumer, TOPIC_NAME, MESSAGE_COUNT)
log.info("[DONE ] 共消费 %d 条消息", got)
if got < MESSAGE_COUNT:
log.warning("预期 %d 条,实际只拿到 %d 条 —— 检查是否之前 Group 已提交过 offset", MESSAGE_COUNT, got)
log.warning("可手动重置:kafka-consumer-groups.sh --bootstrap-server %s "
"--group %s --reset-offsets --to-earliest --topic %s --execute",
BOOTSTRAP_SERVERS, CONSUMER_GROUP, TOPIC_NAME)
finally:
consumer.close()
log.info("👋 完成。下一步可以打开 http://localhost:8080 在 Kafka UI 里查看 Topic 与 Group。")
return 0
if __name__ == "__main__":
sys.exit(main())