Skip to content

第 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.hellolearn.04.orders。不污染业务命名空间。
默认连接bootstrap.servers=localhost:9092对应根目录 docker-compose.yml 暴露的 broker1。
配置项写法官方点分小写bootstrap.serversenable.idempotenceacks=all,不缩写不大小写混写。Python 与 Java 的配置键对照见 §1.10。
可视化Kafka UIhttp://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)有两个致命缺陷:

  1. 吞吐撑不住:消息一旦被消费就丢,而 LinkedIn 需要保留 7 天供任意下游回溯。
  2. 不能多订阅: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 章
2PageCacheBroker 不在 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 的最小单元:消费者组重平衡也是按分区分配。
  • 存储的物理单元:每个分区对应磁盘上一个目录,里面装一堆 .log Segment。

📌 记住一句话:「Kafka = Topic + Partition + 不可变追加日志 + 多副本 + 消费组」。后面所有章节都是在这五个词基础上展开。


1.4 与其它消息中间件的全面对比

这是面试高频题,必须能脱口而出。

1.4.1 五大消息中间件全景对比表

维度KafkaRabbitMQRocketMQPulsarRedis Stream / NATS
诞生年份2010(LinkedIn)2007(Rabbit Tech)2012(阿里)2016(Yahoo)Redis 5.0 / 2010
协议自研二进制AMQP 0.9.1 / STOMP / MQTT自研自研(兼容 Kafka)RESP / NATS
语言Scala / JavaErlangJavaJava(Broker)+ C++C / Go
存储模型磁盘顺序日志 + 多副本内存队列 + 持久化(lazy queue)磁盘日志(CommitLog 单文件)存算分离:Broker 计算 + BookKeeper 存储内存(Stream 持久化为 RDB / AOF)
消费模型拉(pull)+ Offset 提交推(push)+ ACK拉为主,支持长轮询推 + ACK / 拉皆可
顺序保证单分区严格有序单队列有序单队列有序(顺序消费需指定 Queue)单分区有序Stream 内有序
路由仅按 Key Hash 到 PartitionExchange + Binding + Routing Key(最灵活)Topic + TagTopic + SubscriptionStream Key
吞吐(单机 4C8G 参考)百万级 QPS1~10 万 QPS10~50 万 QPS50~100 万 QPS10~50 万 QPS
延迟(p99)5~50 ms< 5 ms< 10 ms< 10 ms< 1 ms(内存)
是否可重放(按 Offset / 时间戳)否(消息消费即删)是(按 Offset)是(按 ID)
事务 / EOS幂等 Producer + 事务简单事务(性能差)事务消息(半消息 + 二次确认)事务(接近 Kafka)不支持
延迟 / 定时消息不原生支持(需外部调度)原生(TTL + DLX)原生支持(多个延迟级别)原生不支持
死信队列 DLQ第三方 / 自建原生原生原生自建
多租户 / 命名空间ACL + QuotavhostTopic GroupTenant / Namespace 三层DB Index
生态最丰富(Connect / Streams / ksqlDB / Schema Registry)中等中等(中文文档好)增长中Redis 生态 / 轻量
典型场景大数据日志总线 / 流处理 / CDC / 事件溯源业务异步 / RPC 风格队列电商交易(订单 / 支付)多租户云原生 / 跨地域缓存里的轻量队列 / 微服务
架构特点Broker = 存储 + 计算 耦合内存为主,集群方案弱NameServer + Broker(去 ZK)Broker 无状态 + BookKeeper 存储单进程内存

1.4.2 5 个最容易踩坑的差异

  1. Kafka 是拉模型,RabbitMQ 是推模型:拉模型让消费者按自己节奏(防止打挂自己),推模型反应快但消费者一旦慢下来 Broker 就要做「流控」很复杂。
  2. Kafka 路由能力弱:只能按 Key Hash,做不了「订单金额 > 1000 走 VIP 队列」这种规则。要么客户端自己路由,要么用 Streams 做拆流。这不是 bug,是为了换吞吐做的取舍
  3. Kafka 没有原生延迟消息:要做「30 分钟未支付自动取消」这种业务,得借助外部调度器或 Confluent 的 KIP-815。RocketMQ 一行配置搞定。
  4. Pulsar 不是 Kafka 的替代品:Pulsar 强在「存算分离 + 多租户」,适合云厂商;Kafka 强在「生态 + 简单」,适合自建。
  5. 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-kafkaJava kafka-clients
集群地址bootstrap.serversbootstrap.servers
Producer ack 策略acksacks
启用幂等 Producerenable.idempotenceenable.idempotence
攒批等待linger.mslinger.ms
单批最大字节batch.sizebatch.size
压缩算法compression.typecompression.type
Consumer 组group.idgroup.id
起始位置策略auto.offset.resetauto.offset.reset
是否自动提交enable.auto.commitenable.auto.commit
心跳间隔heartbeat.interval.msheartbeat.interval.ms
单次 poll 最大数max.poll.records (注: librdkafka 中通过 queued.max.messages.kbytes/fetch.message.max.bytes 控制)max.poll.records
客户端 IDclient.idclient.id

📌 核心一致confluent-kafka 是对 librdkafka 的 Python 封装,librdkafka 又是 Kafka 协议的 C 重写,配置键基本与 Java 一一对应。少数 librdkafka 特有的(如 queued.max.messages.kbytes)会在用到的章节单独说明。


1.11 面试高频题

Q1:Kafka 和 RabbitMQ 的本质区别是什么?

考察点:是否真的理解两者解决的不是同一个问题,而不是只会背「Kafka 吞吐高 / RabbitMQ 延迟低」。

标准答案(按重要性分点):

  1. 设计定位完全不同
    • RabbitMQ 是消息队列(Queue),核心抽象是「队列里的消息被谁取走就消失了」,关注「可靠投递」。
    • Kafka 是分布式提交日志(Distributed Commit Log),核心抽象是「磁盘上一份不可变的追加日志」,关注「海量数据持久化 + 多次重放」。
  2. 存储模型不同:RabbitMQ 内存为主 + 持久化补充;Kafka 磁盘顺序写 + PageCache 为核心。
  3. 消费模型不同:RabbitMQ 是 Broker 推(push) + 客户端 ACK;Kafka 是客户端拉(pull) + 自己提交 Offset。Kafka 拉模型让消费者按自己速率消费,永远不会被压垮。
  4. 路由能力不同:RabbitMQ 有 Exchange + Binding + Routing Key 三层灵活路由(topic / direct / fanout / headers);Kafka 只能按 Key Hash 到 Partition,路由极简,换吞吐。
  5. 顺序保证不同:RabbitMQ 单队列内有序;Kafka 单分区内有序,跨分区不保证。
  6. 吞吐量级不同:单机 RabbitMQ 1~10 万 QPS,Kafka 百万 QPS。
  7. 消息生命周期不同:RabbitMQ 消费完即删;Kafka 按时间(默认 7 天)或 Compaction 策略清理,消费完不删,可被任意消费组重复读。

加分项:能讲一句「RabbitMQ 是消息队列,Kafka 是日志平台」,并举例:要做 RPC 风格异步任务(订单 → 短信通知)选 RabbitMQ;要做大数据日志总线 / CDC / 流处理选 Kafka。

易错点:不要说「Kafka 性能更好」 —— 在「单条消息 + 复杂路由 + 低延迟」场景,RabbitMQ p99 延迟(< 5ms)甚至比 Kafka(5-50ms)更低。


Q2:为什么 Kafka 这么快?吞吐量为什么能比传统 MQ 高一个数量级?

考察点:「五大性能基石」必须能脱口而出,并能讲出每一条的「为什么」。

标准答案:5 个核心机制协同作用:

  1. 顺序写磁盘:Producer 把消息只追加到 Partition 日志末尾,机械盘顺序写吞吐能到 600 MB/s(vs 随机写 0.4 MB/s,差 1500 倍),SSD 上仍有 30 倍差距。Kafka 强制 append-only,从根本上规避了寻道 / GC 开销
  2. PageCache(操作系统页缓存):Broker 不在 JVM 堆里缓存消息,全部交给 OS 的 PageCache 管理。好处:① 避免 JVM GC 压力;② 进程重启缓存不丢;③ 享受 OS 自带的智能预读 / 后写;④ 与读写共用一份内存,利用率极高。
  3. 零拷贝 sendfile:消费者拉消息时,Broker 用 sendfile 系统调用直接把 PageCache 的数据 → 网卡,绕过用户态 JVM 堆。从传统的「4 次拷贝 + 2 次切换」减到「2 次拷贝 + 0~1 次切换」,CPU 几乎不参与,单台 Broker 能服务上万个 Consumer。
  4. 批量化:Producer 端 linger.ms + batch.size 攒一波再发;Broker 端按批写盘;Consumer 端按批拉取。单条消息开销摊薄到几百上千条上,吞吐提升 5-50 倍。
  5. 压缩: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(同步调用)解决不了的:

  1. 解耦:上游不需要知道下游有谁、在哪、活没活。新增一个下游消费者只需要订阅 Topic,上游不改一行代码。RPC 模式下 N 个上游 × M 个下游 = N×M 条耦合关系。
  2. 异步:上游写完消息立刻返回(毫秒级),不需要等下游处理。用户体验大幅提升(下单 → 立即提示成功,发短信 / 加积分 / 推荐召回慢慢做)。RPC 必须串行等所有下游响应。
  3. 削峰:流量洪峰先进 Kafka 这个「水库」,下游按自己稳定速率消费。保护下游不被瞬时高 QPS 打挂(秒杀场景必备)。
  4. 日志总线 / 数据集成:一份数据写一次到 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 个优势

  1. 消费者按自己速率消费,永远不会被压垮:消费慢就少拉一点,消费快就多拉一点。Push 模型下 Broker 不知道消费者能扛多少,必须做复杂的「流控 / 反压」(RabbitMQ 的 prefetch 配置就是干这个的)。
  2. 天然支持批量拉取:一次 fetch 拿一批,吞吐高。Push 是 Broker 一条一条推,批量化困难。
  3. 消费者宕机 / 恢复对 Broker 透明:Broker 不维护「消息发给谁了」状态,Consumer 自己维护 Offset。Push 模型 Broker 必须记 ACK 状态。
  4. 支持「重置消费位置 / 重放历史」: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) % N Hash 到固定分区,自然顺序。
  • 千万不要把整个 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

两条核心规则

  1. 同一个 Group 内:每个 Partition 只能被 1 个 Consumer 消费 —— 单播 / 负载均衡。
  2. 不同 Group 之间:每个 Partition 被每个 Group 独立 消费一次 —— 广播 / 多订阅。

生活类比:Consumer Group 就像「一家三口共用一个邮箱」 —— 邮件来了任何一个人取走就行,不会被取两次(Group 内);但「张家、李家、王家」各自有邮箱,同一封广告信会发 3 份(Group 间)。

应用模式

  • 完全广播:每个消费者一个独立 Group。
  • 完全单播 / 负载均衡:所有消费者同一个 Group。
  • 分组广播 + 组内负载均衡:每个业务方一个 Group(最常用模式)。

加分项

  • 能讲「Consumer 数量 ≤ Partition 数量」 —— 多了的 Consumer 会闲置,因为分区分配不到。生产建议两者相等或 Consumer 数 = Partition 数 / 2。
  • 能讲「Group Coordinator 在哪台 Broker 上」 —— 通过 __consumer_offsets Topic 的 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())

hello_kafka.py ↗