Skip to content

第 3 章 命令行与客户端基础:先把 Kafka「玩起来」

目标读者:第 1、2 章已经搞清楚 Topic / Partition / Broker / Offset / Consumer Group 是什么,现在想自己动手「敲一敲、看一看」的同学。

学完你会:闭着眼睛用 kafka-topics.sh 建/删/改 Topic,用 kafka-console-producer/consumer.sh 收发消息,用 kafka-consumer-groups.sh 查 Lag、重置 Offset,用 kafka-dump-log.sh 解开一个真实的 segment 文件;用 Python confluent-kafka 写出第一个 Producer / Consumer / AdminClient 程序;并能在 Kafka UI 浏览器里观察集群的实时状态。


0. 导读:命令行是「Kafka 的瑞士军刀」

很多教程一上来就让你写 Java / Python 代码,看似很高大上,其实绕过了最重要的「肌肉记忆」环节:Kafka 自带的 bin/kafka-*.sh 工具是排障、运维、压测、教学的事实标准,不会用命令行 = 没法调试线上集群

本章我们走「命令行 → Python → 浏览器 UI」三段式:

  1. 先用命令行手敲一遍核心动作:建 Topic、发消息、收消息、查消费组、重置 Offset、查看 segment 文件。
  2. 再用 Python 写等价的最小程序,让你从「人发消息」过渡到「程序发消息」。
  3. 最后用 Kafka UI(provectus/kafka-ui) 观察集群,建立「直觉化」的画面感。

📌 本章命令全部在「单机 KRaft 模式」上验证(bootstrap.servers=127.0.0.1:9092,单 Broker,3 分区,1 副本)。所有输出都来自真实运行结果,已脱敏。


1. Kafka 命令行家族总览

$KAFKA_HOME/bin/ 下常用脚本梳理一遍,按职责分四类:

Kafka 自带 CLI 全家桶
┌─────────────────────────────────────────────────────────────────────┐
│ 元数据管理                                                          │
│   kafka-topics.sh           Topic 增删改查                          │
│   kafka-configs.sh          Broker / Topic / User / Client 动态配置 │
│   kafka-acls.sh             ACL 权限                                │
├─────────────────────────────────────────────────────────────────────┤
│ 消息收发(人肉版)                                                  │
│   kafka-console-producer.sh   stdin -> Kafka                        │
│   kafka-console-consumer.sh   Kafka -> stdout                       │
│   kafka-verifiable-producer.sh / -consumer.sh   带统计的压测版      │
├─────────────────────────────────────────────────────────────────────┤
│ 消费组与 Offset                                                     │
│   kafka-consumer-groups.sh   list / describe / reset-offsets        │
│   kafka-get-offsets.sh       按时间戳查 offset                      │
├─────────────────────────────────────────────────────────────────────┤
│ 存储与运维                                                          │
│   kafka-dump-log.sh          解析 .log / .index / .timeindex        │
│   kafka-log-dirs.sh          查看每个 Broker 的磁盘占用             │
│   kafka-reassign-partitions.sh   分区迁移                           │
│   kafka-leader-election.sh   主动触发 Leader 选举                   │
│   kafka-features.sh          查看 / 升级 KRaft feature 等级         │
└─────────────────────────────────────────────────────────────────────┘

🍱 生活类比:把 Kafka 想象成一个邮局,那么——

  • kafka-topics.sh = 「设置信箱柜」(开新柜子、换隔板、贴新标签);
  • kafka-console-producer.sh = 「人手往信箱里塞信」;
  • kafka-console-consumer.sh = 「人手从信箱里抽信看」;
  • kafka-consumer-groups.sh = 「查家庭成员各自看到第几封信、把读到的进度往回拨/往前推」;
  • kafka-configs.sh = 「不停业改邮局规章制度」;
  • kafka-dump-log.sh = 「拆开邮局的归档箱看里面到底装了什么字节」。

2. kafka-topics.sh:Topic 的增删改查

所有命令默认追加 --bootstrap-server 127.0.0.1:9092,下面的输出来自一个干净的单 Broker KRaft 集群。

2.1 创建 Topic:--create

bash
kafka-topics.sh --bootstrap-server 127.0.0.1:9092 \
  --create \
  --topic learn.03.hello \
  --partitions 3 \
  --replication-factor 1

输出

Created topic learn.03.hello.

参数解读:

参数含义常见取值
--topicTopic 名推荐 learn.<chapter>.<scene> 风格
--partitions分区数写吞吐 / 并行消费上限的主控因素
--replication-factor副本数1=单点;3=生产推荐;不能 > Broker 数
--config k=v创建时同时下发 Topic 级配置可重复多次,例如 --config retention.ms=3600000

带配置创建的完整示例:

bash
kafka-topics.sh --bootstrap-server 127.0.0.1:9092 \
  --create --topic learn.03.orders \
  --partitions 6 --replication-factor 1 \
  --config retention.ms=86400000 \
  --config segment.bytes=104857600 \
  --config cleanup.policy=delete \
  --config min.insync.replicas=1

⚠️ 常见报错InvalidReplicationFactorException: Replication factor: 3 larger than available brokers: 1. —— 单 Broker 集群只能开 1 副本,加 Broker 才能开多副本。

2.2 列出 Topic:--list

bash
kafka-topics.sh --bootstrap-server 127.0.0.1:9092 --list

输出

__consumer_offsets
learn.03.hello
learn.03.orders

__consumer_offsets 是 Kafka 自己管理消费 Offset 的内部 Topic(第 5 章重点讲),默认不显示,只有指定 --exclude-internal=false(默认就这样)才会出现。如果想隐藏内部 Topic,加 --exclude-internal

2.3 查看 Topic 详情:--describe

bash
kafka-topics.sh --bootstrap-server 127.0.0.1:9092 \
  --describe --topic learn.03.hello

输出

Topic: learn.03.hello   TopicId: pX2H1ZzMSi6...   PartitionCount: 3   ReplicationFactor: 1   Configs: segment.bytes=1073741824
        Topic: learn.03.hello   Partition: 0    Leader: 1   Replicas: 1   Isr: 1   Elr:    LastKnownElr:
        Topic: learn.03.hello   Partition: 1    Leader: 1   Replicas: 1   Isr: 1   Elr:    LastKnownElr:
        Topic: learn.03.hello   Partition: 2    Leader: 1   Replicas: 1   Isr: 1   Elr:    LastKnownElr:

逐字段拆解:

字段含义
TopicIdKRaft 引入的 UUID,删除重建后也是新 ID(用来防错乱)
PartitionCount分区数
ReplicationFactor副本数
Configs当前生效的 Topic 级配置(只列与默认值不同的)
Partition分区号(从 0 起)
Leader当前 Leader 副本所在的 Broker ID
Replicas该分区的所有副本 Broker ID 列表(含 Leader)
IsrIn-Sync Replicas,与 Leader 保持同步的副本集合
Elr / LastKnownElrEligible Leader Replicas(KIP-966,4.x 新特性,可暂时忽略)

只看「副本不健康」的分区

bash
kafka-topics.sh --bootstrap-server 127.0.0.1:9092 \
  --describe --under-replicated-partitions

正常时输出为空。线上巡检必备——任何输出都意味着「有副本掉队 ISR」。

类似地:

  • --unavailable-partitions:列出无 Leader 的分区(数据不可读写)。
  • --under-min-isr-partitions:列出 ISR 少于 min.insync.replicas 的分区。
  • --at-min-isr-partitions:列出刚好等于 min.insync.replicas 的分区(再掉一台就要写失败了)。

2.4 修改分区数:--alter

bash
kafka-topics.sh --bootstrap-server 127.0.0.1:9092 \
  --alter --topic learn.03.hello --partitions 6

输出

Adding partitions succeeded!

⚠️ 永远只能加分区,不能减!而且:

  • 加分区后,已经按 Key Hash 分布的消息会「错乱」(同一个 Key 之前去 0 号分区,加完后可能去 4 号),破坏顺序保证。
  • 真要「减分区」或「重新分布」,必须新建 Topic + 双写迁移,没有捷径。

2.5 删除 Topic:--delete

bash
kafka-topics.sh --bootstrap-server 127.0.0.1:9092 \
  --delete --topic learn.03.hello

输出

(无输出,命令立即返回)

⚠️ 注意几点:

  1. Broker 必须开启 delete.topic.enable=true(默认即开)。
  2. 命令是异步的:返回成功 ≠ 数据已经物理删除,要等 log.retention.check.interval.ms(默认 5 分钟)的清理线程扫到。
  3. KRaft 模式下删除走的是元数据日志,元数据立刻不可见,但磁盘上的 segment 文件会延迟清。

2.6 一次性查看「Topic 全量配置」

--describe 默认只显示与默认值不同的配置。要看一个 Topic 的全部生效配置,得用下一节的 kafka-configs.sh


3. kafka-console-producer.sh:人肉发消息

3.1 最简:发普通消息

bash
kafka-console-producer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.hello

进入交互式 stdin,每行一条消息:

> hello kafka
> 你好 卡夫卡
> {"order_id":1001,"amount":99.5}
> ^D            (Ctrl+D 退出)

💡 没有任何提示符 >,是 Kafka 自带 producer 的「裸输入」体验。要带提示符也能加 --property,下文细说。

3.2 带 Key 发送

很多业务需要按 Key 分区(让同一个用户 / 订单 ID 的消息进入同一个分区,保证顺序)。

bash
kafka-console-producer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders \
  --property "parse.key=true" \
  --property "key.separator=:"

输入:

user-1:{"action":"login"}
user-1:{"action":"add_to_cart"}
user-2:{"action":"login"}
user-1:{"action":"checkout"}

user-1 的 4 条消息会全部落到同一个分区(按 Key Murmur2 Hash),user-2 的消息可能落到别的分区。这就是「Key = 顺序保证的边界」。

3.3 调整可靠性参数

bash
kafka-console-producer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders \
  --producer-property acks=all \
  --producer-property enable.idempotence=true \
  --producer-property compression.type=zstd

--producer-property 后面可以跟任何标准 Producer 配置(详见第 4 章),把命令行变成一个简易的「可配置发送器」。

3.4 一次性发指定文件

bash
kafka-console-producer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders < orders.txt

把每行作为一条消息发出。脚本压测、灌数据时常用。


4. kafka-console-consumer.sh:人肉收消息

4.1 默认只收「订阅之后」的新消息

bash
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.hello

启动后只能看到启动后新发的消息。Kafka 默认 auto.offset.reset=latest

4.2 --from-beginning:从头消费

bash
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.hello \
  --from-beginning

输出(假设 Topic 里已经有 3 条消息):

hello kafka
你好 卡夫卡
{"order_id":1001,"amount":99.5}

⚠️ --from-beginning 等价于 --consumer-property auto.offset.reset=earliest仅在该消费组没有已提交 Offset 时才生效。已经有 Offset 的话仍然从上次位置继续。

4.3 同时打印 Key / 时间戳 / 分区 / Offset

bash
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders \
  --from-beginning \
  --property print.key=true \
  --property print.value=true \
  --property print.partition=true \
  --property print.offset=true \
  --property print.timestamp=true \
  --property key.separator=" | "

输出示例

CreateTime:1713344000123	Partition:2	Offset:0	user-1 | {"action":"login"}
CreateTime:1713344000456	Partition:2	Offset:1	user-1 | {"action":"add_to_cart"}
CreateTime:1713344000789	Partition:0	Offset:0	user-2 | {"action":"login"}
CreateTime:1713344001012	Partition:2	Offset:2	user-1 | {"action":"checkout"}

注意 user-1 的 3 条都在 Partition 2,Offset 严格 0、1、2 递增——单分区有序的直观证明。

4.4 加入消费组:--group

不带 --group 时,每次启动都是一个临时随机消费组,不会提交 Offset。带 --group 后才是「正经消费组」:

bash
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders \
  --group team-A

启动两个终端用同一个 --group team-A,会发现:

  • 同一条消息只被其中一个终端收到(同组分摊)。
  • 关掉一个,另一个会自动接管它的分区(Rebalance,第 5 章细讲)。

4.5 限量消费:--max-messages

bash
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders \
  --from-beginning \
  --max-messages 5

收满 5 条立即退出。线上排障神器:「我只想看最近的几条样本,看完就走」。

4.6 只消费某个分区:--partition

bash
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders \
  --partition 2 \
  --offset earliest

--partition + --offset(earliest / latest / 整数)的组合等价于 Java API 的 assign()绕过消费组分配直接读指定分区——专治「只想验证某个分区里有什么」的场景。

4.7 按时间戳消费

需要先用 kafka-get-offsets.sh 把时间戳转成 Offset:

bash
kafka-get-offsets.sh --bootstrap-server 127.0.0.1:9092 \
  --topic learn.03.orders \
  --time 1713344000000           # Unix 毫秒

输出

learn.03.orders:0:0
learn.03.orders:1:5
learn.03.orders:2:12

每行 topic:partition:offset 表示该时间戳之后的第一条消息位置,再喂给 --offset 即可。


5. kafka-consumer-groups.sh:消费组的「驾驶舱」

5.1 列出所有消费组

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 --list

输出

team-A
console-consumer-23498
team-B

console-consumer-XXXXXkafka-console-consumer.sh 不带 --group 时随机生成的临时组,会有大量残留,正常忽略。

5.2 查看 Lag(最常用!

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
  --describe --group team-A

输出

GROUP   TOPIC              PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CONSUMER-ID                     HOST            CLIENT-ID
team-A  learn.03.orders    0          5               5               0    consumer-team-A-1-abc123        /127.0.0.1      consumer-team-A-1
team-A  learn.03.orders    1          12              15              3    consumer-team-A-1-abc123        /127.0.0.1      consumer-team-A-1
team-A  learn.03.orders    2          80              80              0    consumer-team-A-2-def456        /127.0.0.1      consumer-team-A-2

字段:

字段含义
CURRENT-OFFSET该消费组在该分区已提交到的 Offset(下条要消费的位置)
LOG-END-OFFSET分区当前最大 Offset + 1(即末尾)
LAGLOG-END-OFFSET - CURRENT-OFFSET还没消费完的消息数
CONSUMER-ID当前持有该分区的消费者实例 ID(消费组内唯一)
HOST该消费者所在主机
CLIENT-ID客户端配置的 client.id(不配则自动生成)

排障要点:

  • CONSUMER-ID 为空 → 该分区没有任何活跃消费者,可能消费者全挂了。
  • LAG 持续上涨 → 消费速度跟不上生产,需要扩消费者 / 优化业务逻辑 / 扩分区。
  • 多个 CONSUMER-ID 频繁变动 → Rebalance 风暴,检查 max.poll.interval.ms 与业务处理时间。

5.3 查看消费组的成员构成

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
  --describe --group team-A --members --verbose

输出

GROUP   CONSUMER-ID                    HOST          CLIENT-ID            #PARTITIONS  ASSIGNMENT
team-A  consumer-team-A-1-abc123       /127.0.0.1    consumer-team-A-1    2            learn.03.orders(0,1)
team-A  consumer-team-A-2-def456       /127.0.0.1    consumer-team-A-2    1            learn.03.orders(2)

直观看到「谁拿了哪些分区」。

5.4 查看消费组的状态

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
  --describe --group team-A --state

输出

GROUP   COORDINATOR (ID)        ASSIGNMENT-STRATEGY  STATE     #MEMBERS
team-A  127.0.0.1:9092 (1)      range                Stable    2

STATE 取值:

  • PreparingRebalance / CompletingRebalance:正在再平衡(应是瞬态)。
  • Stable:正常工作中。
  • Empty:组存在但无成员。
  • Dead:组已被删除。

5.5 重置 Offset:--reset-offsets

最容易出事的命令,看清楚再敲

5.5.1 演练(不真改)

--dry-run 永远先演练一遍:

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
  --group team-A --topic learn.03.orders \
  --reset-offsets --to-earliest --dry-run

输出

GROUP   TOPIC              PARTITION  NEW-OFFSET
team-A  learn.03.orders    0          0
team-A  learn.03.orders    1          0
team-A  learn.03.orders    2          0

5.5.2 执行

--dry-run 换成 --execute

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
  --group team-A --topic learn.03.orders \
  --reset-offsets --to-earliest --execute

⚠️ 执行前必须确保该消费组没有活跃消费者(CONSUMER-ID 为空),否则会报:

Error: Assignments can only be reset if the group 'team-A' is inactive, but the current state is Stable.

5.5.3 重置策略汇总

选项含义
--to-earliest重置到分区最早 Offset
--to-latest重置到分区末尾(跳过所有未消费
--to-offset N重置到指定 Offset N
--to-current保持当前 Offset 不变(一般搭配 --shift-by 用)
--shift-by N在当前基础上前进/后退 N 条(N 可负)
--to-datetime YYYY-MM-DDTHH:mm:ss.SSS重置到该时间之后的第一条消息
--by-duration PnDTnHnMnSISO 8601 时间间隔,例如 PT1H 回退 1 小时
--from-file file.csv从 CSV(topic,partition,offset)批量设置

实战示例:「消费者写出 Bug,把 1 小时的数据重新跑一遍」:

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
  --group team-A --topic learn.03.orders \
  --reset-offsets --by-duration PT1H --execute

5.6 删除消费组

bash
kafka-consumer-groups.sh --bootstrap-server 127.0.0.1:9092 \
  --delete --group team-A

输出

Deletion of requested consumer groups ('team-A') was successful.

清理临时测试组的好工具。组里如果还有活跃成员会失败。


6. kafka-configs.sh:动态配置中心

6.1 查看 Topic 全量配置

bash
kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \
  --entity-type topics --entity-name learn.03.orders \
  --describe --all

输出(截断,实际有 70+ 行):

All configs for topic learn.03.orders are:
  cleanup.policy=delete sensitive=false synonyms={DYNAMIC_TOPIC_CONFIG:cleanup.policy=delete, ...}
  retention.ms=86400000 sensitive=false synonyms={DYNAMIC_TOPIC_CONFIG:retention.ms=86400000, DEFAULT_CONFIG:log.retention.ms=604800000}
  segment.bytes=104857600 sensitive=false synonyms={DYNAMIC_TOPIC_CONFIG:segment.bytes=104857600, ...}
  min.insync.replicas=1 sensitive=false synonyms={STATIC_BROKER_CONFIG:min.insync.replicas=1, DEFAULT_CONFIG:min.insync.replicas=1}
  ...

synonyms 揭示了配置生效优先级

DYNAMIC_TOPIC_CONFIG > DYNAMIC_BROKER_CONFIG > STATIC_BROKER_CONFIG > DEFAULT_CONFIG
   (Topic 级动态)        (Broker 级动态)        (server.properties)     (Kafka 默认值)

6.2 动态修改 Topic 配置

不重启就能改:

bash
kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \
  --entity-type topics --entity-name learn.03.orders \
  --alter --add-config retention.ms=3600000,segment.bytes=33554432

输出

Completed updating config for topic learn.03.orders.

6.3 删除动态配置(恢复默认)

bash
kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \
  --entity-type topics --entity-name learn.03.orders \
  --alter --delete-config retention.ms

6.4 改 Broker 级动态配置

bash
kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \
  --entity-type brokers --entity-name 1 \
  --alter --add-config "log.cleaner.threads=4,num.replica.fetchers=8"

部分参数标记为「read-only」(如 broker.id),改它们必须重启。

6.5 给 Client 设 Quota

bash
# 给 client.id=batch-importer 限速到 10MB/s 写入 + 5MB/s 读取
kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \
  --entity-type clients --entity-name batch-importer \
  --alter --add-config "producer_byte_rate=10485760,consumer_byte_rate=5242880"

第 15 章「安全与多租户」会回到 Quota 的细节。


7. kafka-dump-log.sh:拆开 Segment 看字节

本工具是「Kafka 存储原理」的最强可视化工具,第 7 章会重度使用。本章先入门:知道怎么调出来、看什么。

7.1 找到 Segment 文件

Broker 的日志目录由 log.dirs 决定(默认 /tmp/kraft-combined-logs/)。每个分区一个子目录:

bash
ls /tmp/kraft-combined-logs/learn.03.orders-0/

输出

00000000000000000000.index
00000000000000000000.log
00000000000000000000.timeindex
leader-epoch-checkpoint
partition.metadata
  • .log:消息正文(顺序追加写入的二进制文件)。
  • .index:Offset → 文件位置的稀疏索引。
  • .timeindex:时间戳 → Offset 的稀疏索引。
  • leader-epoch-checkpoint:Leader 任期变更记录(用来防截断丢数据)。
  • partition.metadata:分区 UUID 等元信息。

文件名是该 segment 的起始 Offset(base offset),所以 00000000000000000000.log 表示「从 Offset 0 起」。下一个 segment 文件名就是它的最大 Offset+1。

7.2 解析 .log 文件

bash
kafka-dump-log.sh --files /tmp/kraft-combined-logs/learn.03.orders-0/00000000000000000000.log \
  --print-data-log

输出(截断):

Dumping /tmp/kraft-combined-logs/learn.03.orders-0/00000000000000000000.log
Log starting offset: 0
baseOffset: 0 lastOffset: 2 count: 3 baseSequence: -1 lastSequence: -1 producerId: -1 producerEpoch: -1 partitionLeaderEpoch: 0 isTransactional: false isControl: false deleteHorizonMs: OptionalLong.empty position: 0 CreateTime: 1713344000123 size: 178 magic: 2 compresscodec: none crc: 1857384021 isvalid: true
| offset: 0 CreateTime: 1713344000123 keySize: 6 valueSize: 24 sequence: -1 headerKeys: [] key: user-1 payload: {"action":"login"}
| offset: 1 CreateTime: 1713344000456 keySize: 6 valueSize: 30 sequence: -1 headerKeys: [] key: user-1 payload: {"action":"add_to_cart"}
| offset: 2 CreateTime: 1713344001012 keySize: 6 valueSize: 28 sequence: -1 headerKeys: [] key: user-1 payload: {"action":"checkout"}
baseOffset: 3 lastOffset: 3 count: 1 ... position: 178 CreateTime: 1713344002000 size: 95 magic: 2 compresscodec: zstd ...
| offset: 3 CreateTime: 1713344002000 keySize: 6 valueSize: 18 sequence: -1 headerKeys: [] key: user-2 payload: {"action":"login"}

逐字段重点:

  • RecordBatch 头baseOffset / lastOffset / count / position / compresscodec —— 一个 batch 是 Producer 攒起来一起发的一组消息,是 Kafka 性能的核心单元(第 4 章细讲)。
  • 每条 Recordoffset / key / payload / sequence(幂等生产时的序列号)/ producerId(事务/幂等用的 PID)。
  • magic: 2 表示 V2 消息格式(当代 Kafka 默认)。
  • isTransactional / isControl:事务消息相关,第 13 章细讲。

7.3 解析 .index 文件

bash
kafka-dump-log.sh --files /tmp/kraft-combined-logs/learn.03.orders-0/00000000000000000000.index

输出

Dumping /tmp/kraft-combined-logs/learn.03.orders-0/00000000000000000000.index
offset: 0 position: 0
offset: 4096 position: 524160
offset: 8192 position: 1047840

每个条目 8 字节 offset + 4 字节 position(文件偏移)。稀疏索引:默认每写满 4KB 才记一条(index.interval.bytes=4096)。查找时先二分定位最接近的索引项,再顺序扫几条 batch 即可命中目标 offset。

7.4 校验数据完整性

bash
kafka-dump-log.sh --files .../00000000000000000000.log --verify-index-only

会校验 .index 与 .log 的对应关系,CRC 是否正确。怀疑磁盘坏块时直接跑这条。


8. Python confluent-kafka 入门:三大 API

confluent-kafka 是基于 librdkafka(C 实现)的 Python 客户端,性能与语义完整度都是 Python 阵营第一。安装:

bash
pip install confluent-kafka==2.4.0

8.1 Producer 三连:produce + flush + 回调

python
from confluent_kafka import Producer

p = Producer({
    "bootstrap.servers": "127.0.0.1:9092",
    "client.id": "py-hello-producer",
    "acks": "all",
    "enable.idempotence": True,
})

def on_delivery(err, msg):
    if err is not None:
        print(f"❌ 发送失败: {err}")
    else:
        print(f"✅ 已写入 {msg.topic()}-{msg.partition()}@offset={msg.offset()}")

for i in range(5):
    p.produce(
        topic="learn.03.hello",
        key=f"k{i}".encode(),
        value=f"hello-{i}".encode(),
        on_delivery=on_delivery,
    )

p.flush(timeout=10)

输出

✅ 已写入 learn.03.hello-2@offset=0
✅ 已写入 learn.03.hello-2@offset=1
✅ 已写入 learn.03.hello-1@offset=0
✅ 已写入 learn.03.hello-0@offset=0
✅ 已写入 learn.03.hello-2@offset=2

要点:

  • produce() 是异步入队:消息只是放到本地 buffer,不等 broker 确认。
  • flush() 才会等所有未完成消息出 buffer + ack,并触发回调。
  • on_delivery唯一的成功 / 失败通知通道,必须实现。

📌 与 Java KafkaProducer.send().get()(同步等 Future)的区别:confluent-kafka-python 没有「单条同步」语法,只能整批 flush,要单条同步就用 flush(1) 或自己 wrap。

8.2 Consumer 主循环:subscribe + poll

python
from confluent_kafka import Consumer

c = Consumer({
    "bootstrap.servers": "127.0.0.1:9092",
    "group.id": "py-hello-group",
    "auto.offset.reset": "earliest",
    "enable.auto.commit": True,
})

c.subscribe(["learn.03.hello"])

try:
    while True:
        msg = c.poll(timeout=1.0)
        if msg is None:
            continue
        if msg.error():
            print(f"⚠️ 消费错误: {msg.error()}")
            continue
        print(f"[{msg.topic()}-{msg.partition()}@{msg.offset()}] "
              f"key={msg.key()} value={msg.value()}")
except KeyboardInterrupt:
    pass
finally:
    c.close()

输出

[learn.03.hello-2@0] key=b'k0' value=b'hello-0'
[learn.03.hello-2@1] key=b'k1' value=b'hello-1'
[learn.03.hello-1@0] key=b'k2' value=b'hello-2'
[learn.03.hello-0@0] key=b'k3' value=b'hello-3'
[learn.03.hello-2@2] key=b'k4' value=b'hello-4'

要点:

  • poll(timeout)同步阻塞,超时返回 None
  • 每次循环都要检查 msg.error():很多「软错误」(分区 EOF、Rebalance 进行中)也会通过 msg 返回。
  • c.close() 会触发最后一次同步 commit + LeaveGroup,必须在 finally 里调,否则要等 session.timeout.ms 才能被踢出组(第 5 章细讲)。

8.3 AdminClient:用代码管理 Topic / Group

python
from confluent_kafka.admin import AdminClient, NewTopic, ConfigResource

admin = AdminClient({"bootstrap.servers": "127.0.0.1:9092"})

new_topic = NewTopic(
    topic="learn.03.admin_demo",
    num_partitions=3,
    replication_factor=1,
    config={"retention.ms": "3600000"},
)
fs = admin.create_topics([new_topic])
for t, f in fs.items():
    try:
        f.result()
        print(f"✅ 创建 Topic: {t}")
    except Exception as e:
        print(f"❌ {t}: {e}")

md = admin.list_topics(timeout=5)
print("📋 当前 Topic:", list(md.topics.keys()))

AdminClient 提供的核心方法:

方法用途
create_topics(list)批量创建 Topic
delete_topics(list)批量删除 Topic
list_topics()列出所有 Topic + 集群元数据
describe_configs([ConfigResource])查询 Topic / Broker 配置
alter_configs([...])整体替换配置(注意会清掉别的动态配置,4.0 起推荐用 incremental_alter_configs
incremental_alter_configs(...)增量改配置(推荐)
list_consumer_groups()列出所有消费组
describe_consumer_groups([...])查消费组成员、状态
list_consumer_group_offsets([...])查 Offset & Lag
alter_consumer_group_offsets([...])重置 Offset(等价于 kafka-consumer-groups.sh --reset-offsets

配套脚本 code/admin_demo.py 把上述 API 串成一个完整的「建 Topic → 灌数据 → 看 Lag → 重置 Offset → 删 Topic」演示。


9. Kafka UI:浏览器里的「驾驶舱」

provectus/kafka-ui 是开源 Kafka 控制台,零配置启动一个 Docker 容器即可。

9.1 启动

bash
docker run -d --name kafka-ui \
  -p 8080:8080 \
  -e KAFKA_CLUSTERS_0_NAME=local \
  -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=host.docker.internal:9092 \
  provectusinc/kafka-ui:latest

打开 http://localhost:8080,左侧导航有:

Dashboard      整集群一览(Broker 数 / Topic 数 / Partition 数 / 在线状态)
Brokers        每个 Broker 的角色、版本、磁盘占用
Topics         所有 Topic 的列表、分区、副本、消息条数
Consumers      所有消费组、Lag、成员
Schema Registry / KSQL / Connect    可选模块

9.2 操作步骤(无截图,纯描述)

新建 Topic

  1. 点左侧 Topics → 右上 Add a Topic
  2. NameNumber of partitionsReplication factor
  3. Custom parametersretention.ms / cleanup.policy 等。
  4. Create Topic,立即在列表中出现。

浏览消息

  1. 点 Topic 名 → Messages Tab。
  2. 上方有 Seek Type(Offset / Timestamp / Earliest / Latest)和 Filters(Key contains / Value contains / Header)。
  3. 实时滚动展示分区、Offset、Key、Value、Timestamp、Headers。可点单条展开看完整 JSON。

查看消费组

  1. 左侧 Consumers → 选组 → 看到每个 Topic + Partition 的 Lag、当前 Offset、分配的消费者实例。
  2. 右上 Reset offset 一键重置(带 dry-run 视图)。

动态修改配置

  1. Topic 详情 → Configs Tab → 找到要改的配置 → 点 Edit → 输入新值 → Save
  2. 改完会在右侧用红框标出「与默认值不同」的项,与 kafka-configs.sh --describe --all 看到的一致。

与命令行的关系:UI 后端就是调 Kafka AdminClient + Consumer 协议,不会做任何命令行做不到的事。但直观、可视、人人会用,运维和产品同学非常喜欢。


10. 实战脚本一览

本章配套目录 03_cli_basics/

  • init.sh —— 一键创建本章用到的所有 Topic。
  • code/admin_demo.py —— AdminClient 全家桶演示。
  • code/simple_producer.py —— 30 行写完一个生产者。
  • code/simple_consumer.py —— 30 行写完一个消费者。
  • demo.html —— 浏览器里的「Web 版 kafka-cli」,鼠标点点就能学命令。

11. 小结

命令行与客户端基础
├─ Topic 管理
│  ├─ kafka-topics.sh --create / --list / --describe / --alter / --delete
│  └─ ⚠️ 分区只能加不能减;删除是异步的
├─ 收发消息
│  ├─ kafka-console-producer.sh    人肉 stdin → Kafka
│  ├─ kafka-console-consumer.sh    --from-beginning / --group / --max-messages / --partition
│  └─ kafka-get-offsets.sh         按时间戳查 Offset
├─ 消费组
│  ├─ --list / --describe / --members / --state
│  ├─ Lag = LOG-END-OFFSET - CURRENT-OFFSET
│  └─ --reset-offsets:earliest / latest / to-offset / by-duration / from-file
│      必须 --dry-run 先演练,组必须 inactive
├─ 动态配置
│  └─ kafka-configs.sh   topics / brokers / clients / users 四类实体
├─ 存储观察
│  ├─ kafka-dump-log.sh --files xxx.log --print-data-log
│  └─ 读懂 RecordBatch + Record 的字段
├─ Python confluent-kafka
│  ├─ Producer:produce + flush + delivery callback
│  ├─ Consumer:subscribe + poll 循环 + close
│  └─ AdminClient:建/删/改 Topic、查 Lag、重置 Offset
└─ Kafka UI
   └─ Topic / Messages / Consumers / Configs 全功能 GUI

12. 面试高频题(5 题)

Q1:kafka-topics.sh --create--partitions--replication-factor 怎么选?为什么不能后悔减分区?

考察点:分区设计、Kafka 存储模型。

答案

  1. --partitions:决定该 Topic 的最大并行消费度生产者总写入吞吐
    • 经验值:分区数 = 期望峰值吞吐 / 单分区吞吐能力(约 10-50 MB/s),并保留 2-4 倍冗余。
    • 太少 → 消费瓶颈、扩消费者无效;太多 → 元数据放大、Controller 压力大、文件句柄爆炸。
  2. --replication-factor:决定可靠性。生产环境普遍 =3(容忍 1 台 Broker 完全失效)。必须 ≤ Broker 数。
  3. 不能减分区的原因:
    • Kafka 分区数据是「按分区独立存储 + 按 Key Hash 分布」的,删一个分区意味着该分区所有消息丢失,且其他分区的 Key 分布也要重新映射。
    • 加分区也会破坏已有 Key 的顺序保证(同一 Key Hash 后落到的分区编号变了),所以「加分区也是有代价的」,不是无副作用操作。
  4. 正确做法:上线前充分预估,宁可多开 2 倍分区;如果真要「重新分布」,必须新建 Topic + 双写迁移。

加分项:提到 min.insync.replicas 必须 ≤ replication-factor,且 acks=all + min.insync.replicas=2 才是真正的「不丢消息」组合。


Q2:kafka-console-consumer.sh --from-beginning 一定能从头消费吗?

考察点auto.offset.reset 与已有 Offset 的关系。

答案

  1. 不一定--from-beginning 等价于 --consumer-property auto.offset.reset=earliest
  2. auto.offset.reset 只在「该消费组没有该分区的已提交 Offset」时生效。如果消费组之前已经消费过、提交过 Offset,会直接从已提交位置继续,跟 --from-beginning 无关。
  3. 真正想从头的方法:
    • 不带 --group 启动(每次随机生成临时组,没有历史 Offset,肯定从头)。
    • --group + 先用 kafka-consumer-groups.sh --reset-offsets --to-earliest --execute 重置。
    • 改用一个没用过的新 group.id
  4. 类似地,auto.offset.reset=latest(默认)会让一个全新组只看到启动后的新消息,「初次上线少了历史」就是这么踩的。
  5. 第三个取值 auto.offset.reset=none 表示「找不到 Offset 就报错而不是兜底」,事务 / EOS 场景常用。

加分项:提到 __consumer_offsets 内部 Topic 才是「有没有已提交 Offset」的真实存储位置,可以用 kafka-console-consumer.sh --topic __consumer_offsets --formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter" 直接看(第 5 章细讲)。


Q3:kafka-consumer-groups.sh --reset-offsets 有哪些坑?

考察点:运维实操、协议细节。

答案

  1. 必须 --execute 才生效。默认 --dry-run 只打印「将要改成的样子」。生产上永远先 --dry-run 再执行
  2. 目标消费组必须 inactive(state = Empty / Dead)。否则报 Assignments can only be reset if the group is inactive。要么先停消费者,要么用 --all-groups 慎用。
  3. 重置策略多种:--to-earliest / --to-latest / --to-offset N / --shift-by ±N / --by-duration PT1H / --to-datetime / --from-file CSV。注意 --to-datetime 是「该时间之后的第一条」,时区按本地时区解析。
  4. 重置范围--topic T 是该 topic 全部分区;--topic T:0,2 只对 0、2 分区;--all-topics 重置该组的所有 topic。
  5. 重置后不会自动消费,要重新启动消费者才会从新位置 poll。
  6. 常见用例
    • 业务回放:--by-duration PT1H --execute 重新消费 1 小时数据。
    • 跳过坏消息:--shift-by 1 --execute 跳过当前阻塞的一条。
    • 全量重算:--to-earliest --execute
  7. 替代方案:在程序里调 consumer.seek()AdminClient.alter_consumer_group_offsets() 也能达到一样的效果,且能编程化管控。

加分项:提到 KRaft 模式下重置 Offset 走的是 __consumer_offsets 写入,原理与 ZK 时代一致;但删除消费组在 KRaft 下走元数据日志。


Q4:kafka-dump-log.sh 能告诉你哪些信息?怎么排查「磁盘上消息和我发的不一样」?

考察点:Kafka 存储格式、排障思路。

答案

  1. 能看到的信息
    • RecordBatch 级baseOffset / lastOffset / count / producerId / producerEpoch / isTransactional / isControl / compresscodec / position(在 .log 文件中的字节偏移)/ size / magic(消息格式版本)。
    • Record 级offset / key / value / headers / timestamp / sequence(幂等用的序列号)。
  2. 常用参数
    • --print-data-log:把消息正文也打印出来(默认只打 batch 头)。
    • --deep-iteration:解压后打印(compresscodec ≠ none 时必用)。
    • --verify-index-only:校验 .index 与 .log 的对应一致性。
    • --offsets-decoder:当 dump 的是 __consumer_offsets 时,用 OffsetCommit 的解码器还原成 (group, topic, partition) → offset
    • --transaction-log-decoder:dump __transaction_state 时用。
  3. 排障套路「我发的消息和读到的不一样」:
    1. 找到 Broker 的 log.dirs,定位 <topic>-<partition> 目录。
    2. kafka-dump-log.sh --files xxx.log --print-data-log --deep-iteration,看磁盘上实际存的字节
    3. 与 Producer 发送时的 payload 对比:
      • 字节级一致 → 是消费端解码 / 反序列化 / 字符集问题。
      • 不一致 → 检查 Producer 端是否做了 SerDe / 拦截器修改;或者中途经过了 Connect / SMT / Streams 处理。
    4. 关注 compresscodec 字段:如果 batch 是 zstd,不带 --deep-iteration 你看到的是压缩后的二进制乱码。
  4. 进阶:分析事务消息时用 --transaction-log-decoder__transaction_state 内的 Ongoing/PrepareCommit/CompleteCommit,能精确还原一次事务的状态机。

Q5:confluent-kafkaProducer.flush() 到底等什么?什么时候必须调?

考察点:异步 Producer 模型、可靠性边界。

答案

  1. Producer.produce()纯异步入队:消息只是放到 librdkafka 内部的发送队列,立刻返回。不等 broker 任何响应。
  2. librdkafka 内部有个后台 IO 线程根据 linger.ms / batch.size 攒批,向 broker 发 Produce 请求;收到响应后回调你的 on_delivery
  3. flush(timeout) 的语义:
    • 阻塞直到队列里所有未完成消息都收到响应(成功或失败),或超时。
    • 期间会触发 delivery callbackproduce 时注册的回调)。
  4. 必须调 flush 的场景:
    • 程序退出前(否则未发完的消息就丢了)。
    • 一批关键消息发完后,需要确认结果再继续业务(例如订单写完才允许下游计费)。
    • 单元测试里需要同步语义。
  5. 配套技巧
    • flush(0) 不等待,只把已就绪的回调触发一遍。
    • 长生命周期的 Producer 不需要每次都 flush,只在「业务边界」flush;高频 flush 会严重退化为同步发送,吞吐崩盘。
    • Python 还要定期 poll(0)(或在每次 produce 后顺手 poll(0))才能让回调被触发;flush 内部其实就是循环 poll
  6. 错误兜底produce 时如果队列已满(queue.buffering.max.messages 默认 100k)会抛 BufferError,要 catch 后 poll(1) 让队列腾空再重试。

加分项:提到 Java KafkaProducer.send().get() 是「同步」语义但本质上还是 batch;真正的同步发送(一条一发)会让吞吐降到 1/100。


下一章 04_producer.md 我们会深入 Producer 内部:批次聚合、压缩、acks 三种语义、幂等生产者的 PID + Epoch + 序列号机制——把命令行里看到的「一条 produce 命令」拆成几十个内部步骤。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
"""
第 3 章 - AdminClient 全家桶演示
==================================

本脚本依次演示:
1. 用 AdminClient 创建 Topic(带 Topic 级配置)
2. 列出当前所有 Topic + 集群元数据
3. describe_configs 看 Topic 的全量配置
4. 灌入一些消息,再用 list_consumer_group_offsets 看 Lag
5. alter_consumer_group_offsets 重置 Offset 到 earliest
6. 删除 Topic 收尾

依赖:
    pip install confluent-kafka==2.4.0

运行:
    python admin_demo.py [bootstrap]
    默认 bootstrap = 127.0.0.1:9092
"""

from __future__ import annotations

import sys
import time
import uuid
from typing import Iterable

from confluent_kafka import Consumer, Producer, TopicPartition
from confluent_kafka.admin import (
    AdminClient,
    ConfigResource,
    ConsumerGroupTopicPartitions,
    NewTopic,
)

BOOTSTRAP = sys.argv[1] if len(sys.argv) > 1 else "127.0.0.1:9092"
TOPIC = "learn.03.admin_demo"
GROUP = f"admin-demo-grp-{uuid.uuid4().hex[:6]}"


def banner(title: str) -> None:
    print("\n" + "=" * 60)
    print(f"  {title}")
    print("=" * 60)


def safe_run(name: str, futures: dict) -> None:
    for k, f in futures.items():
        try:
            f.result()
            print(f"  ✅ {name}: {k}")
        except Exception as e:
            print(f"  ❌ {name}: {k} -> {e}")


def main() -> None:
    admin = AdminClient({"bootstrap.servers": BOOTSTRAP})

    banner("1) create_topics")
    new_topic = NewTopic(
        topic=TOPIC,
        num_partitions=3,
        replication_factor=1,
        config={
            "retention.ms": "3600000",
            "segment.bytes": "33554432",
            "cleanup.policy": "delete",
        },
    )
    safe_run("create", admin.create_topics([new_topic], request_timeout=10))
    time.sleep(1)

    banner("2) list_topics + 集群元数据")
    md = admin.list_topics(timeout=10)
    print(f"  Cluster ID         : {md.cluster_id}")
    print(f"  Controller Broker  : {md.controller_id}")
    print(f"  Brokers            :")
    for b in md.brokers.values():
        print(f"    - id={b.id} host={b.host}:{b.port}")
    print(f"  Topics ({len(md.topics)}):")
    for t in sorted(md.topics.keys()):
        if t.startswith("__"):
            continue
        partitions = md.topics[t].partitions
        print(f"    - {t} (partitions={len(partitions)})")

    banner("3) describe_configs(看 Topic 全量配置)")
    res = ConfigResource(ConfigResource.Type.TOPIC, TOPIC)
    fs = admin.describe_configs([res])
    cfgs = list(fs.values())[0].result()
    interesting = [
        "cleanup.policy",
        "retention.ms",
        "segment.bytes",
        "min.insync.replicas",
        "compression.type",
        "max.message.bytes",
    ]
    for k in interesting:
        if k in cfgs:
            v = cfgs[k]
            print(f"  {k:24s} = {v.value}  (source={v.source.name})")

    banner(f"4) 灌入 30 条消息 + 用消费组 {GROUP} 消费一半")
    p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
    for i in range(30):
        p.produce(TOPIC, key=f"k{i % 5}".encode(), value=f"msg-{i}".encode())
    p.flush(10)
    print("  ✅ 30 条消息已发出")

    c = Consumer(
        {
            "bootstrap.servers": BOOTSTRAP,
            "group.id": GROUP,
            "auto.offset.reset": "earliest",
            "enable.auto.commit": False,
        }
    )
    c.subscribe([TOPIC])
    consumed = 0
    deadline = time.time() + 10
    while consumed < 15 and time.time() < deadline:
        msg = c.poll(1.0)
        if msg is None or msg.error():
            continue
        consumed += 1
    c.commit(asynchronous=False)
    print(f"  ✅ 已消费并提交 {consumed} 条")

    banner("5) list_consumer_group_offsets(查 Lag)")
    req = ConsumerGroupTopicPartitions(GROUP)
    fs = admin.list_consumer_group_offsets([req])
    res = list(fs.values())[0].result()

    md = admin.list_topics(TOPIC, timeout=5)
    end_offsets = {}
    for part_id in md.topics[TOPIC].partitions.keys():
        low, high = c.get_watermark_offsets(TopicPartition(TOPIC, part_id), timeout=5)
        end_offsets[part_id] = high

    print(f"  Group = {GROUP}")
    print(f"  {'Partition':<10} {'Committed':<12} {'End':<12} {'Lag':<6}")
    total_lag = 0
    for tp in res.topic_partitions:
        committed = tp.offset
        end = end_offsets.get(tp.partition, 0)
        lag = max(0, end - committed)
        total_lag += lag
        print(f"  {tp.partition:<10} {committed:<12} {end:<12} {lag:<6}")
    print(f"  TOTAL LAG = {total_lag}")

    c.close()

    banner("6) alter_consumer_group_offsets(重置到 earliest)")
    new_tps = [TopicPartition(TOPIC, p, 0) for p in md.topics[TOPIC].partitions.keys()]
    req = ConsumerGroupTopicPartitions(GROUP, new_tps)
    fs = admin.alter_consumer_group_offsets([req])
    safe_run("reset", {GROUP: list(fs.values())[0]})

    banner("7) 删除 Topic 收尾")
    safe_run("delete", admin.delete_topics([TOPIC], operation_timeout=10))

    print("\n🎉 admin_demo 演示完成!")


if __name__ == "__main__":
    main()
python
"""
第 3 章 - 最小可用 Consumer
==================================

30 行写完一个消费者,演示:
- group.id / auto.offset.reset / enable.auto.commit 配置
- subscribe + poll 循环模型
- 处理 None / 错误消息 / EOF
- finally 中调 close() 触发同步 commit + LeaveGroup

运行:
    bash ../init.sh
    # 启动两份会自动分摊分区,关掉一份会触发 Rebalance
    python simple_consumer.py
    python simple_consumer.py
"""

from __future__ import annotations

import signal
import sys
import time

from confluent_kafka import Consumer, KafkaError

BOOTSTRAP = sys.argv[1] if len(sys.argv) > 1 else "127.0.0.1:9092"
TOPIC = "learn.03.hello"
GROUP = "ch3-simple-group"

consumer = Consumer(
    {
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "client.id": f"ch3-simple-consumer-{int(time.time())}",
        "auto.offset.reset": "earliest",
        "enable.auto.commit": True,
        "auto.commit.interval.ms": 5000,
        "session.timeout.ms": 10000,
        "max.poll.interval.ms": 60000,
    }
)


def on_assign(c, partitions):
    print(f"👉 被分配到分区: {[(p.topic, p.partition) for p in partitions]}")


def on_revoke(c, partitions):
    print(f"👋 被收回分区: {[(p.topic, p.partition) for p in partitions]}")


running = True


def stop(*_):
    global running
    print("\n⏹  收到退出信号,正在关闭...")
    running = False


signal.signal(signal.SIGINT, stop)
signal.signal(signal.SIGTERM, stop)


def main() -> None:
    print(f"📥 订阅 {TOPIC}(group={GROUP})")
    consumer.subscribe([TOPIC], on_assign=on_assign, on_revoke=on_revoke)
    count = 0
    try:
        while running:
            msg = consumer.poll(timeout=1.0)
            if msg is None:
                continue
            if msg.error():
                if msg.error().code() == KafkaError._PARTITION_EOF:
                    continue
                print(f"⚠️  消费错误: {msg.error()}")
                continue
            count += 1
            print(
                f"[{count:>3}] {msg.topic()}-{msg.partition()}@{msg.offset()} "
                f"key={msg.key().decode() if msg.key() else None} "
                f"value={msg.value().decode()}"
            )
    finally:
        consumer.close()
        print(f"\n📊 共消费 {count} 条消息")


if __name__ == "__main__":
    main()
python
"""
第 3 章 - 最小可用 Producer
==================================

30 行写完一个生产者,演示:
- bootstrap.servers / client.id / acks / enable.idempotence 配置
- 带 Key 发送(同 Key 必去同分区)
- delivery callback 拿到 broker 的 ack 结果
- flush() 确保程序退出前所有消息送达

运行:
    bash ../init.sh                       # 先创建 Topic
    python simple_producer.py             # 发送 10 条到 learn.03.hello
    python simple_producer.py 50 myhost:9092
"""

from __future__ import annotations

import sys
import time

from confluent_kafka import Producer

BOOTSTRAP = sys.argv[2] if len(sys.argv) > 2 else "127.0.0.1:9092"
N = int(sys.argv[1]) if len(sys.argv) > 1 else 10
TOPIC = "learn.03.hello"

producer = Producer(
    {
        "bootstrap.servers": BOOTSTRAP,
        "client.id": "ch3-simple-producer",
        "acks": "all",
        "enable.idempotence": True,
        "linger.ms": 5,
        "compression.type": "zstd",
    }
)


def on_delivery(err, msg):
    if err is not None:
        print(f"❌ 发送失败 key={msg.key()}: {err}")
    else:
        print(
            f"✅ 已写入 {msg.topic()}-{msg.partition()}@offset={msg.offset()} "
            f"key={msg.key().decode() if msg.key() else None} "
            f"value={msg.value().decode()}"
        )


def main() -> None:
    print(f"🚀 向 {TOPIC} 发送 {N} 条消息(acks=all + 幂等)")
    start = time.perf_counter()
    for i in range(N):
        producer.produce(
            topic=TOPIC,
            key=f"user-{i % 3}".encode(),
            value=f"msg-#{i}-ts={int(time.time() * 1000)}".encode(),
            on_delivery=on_delivery,
        )
        # 主动驱动一次回调,避免回调被 buffer 攒住
        producer.poll(0)

    remaining = producer.flush(timeout=15)
    elapsed = time.perf_counter() - start
    print(
        f"\n📊 完成:{N - remaining}/{N} 成功,耗时 {elapsed * 1000:.1f}ms,"
        f"平均 {N / elapsed:.0f} msg/s"
    )
    if remaining > 0:
        print(f"⚠️  仍有 {remaining} 条未发出(flush 超时)")


if __name__ == "__main__":
    main()

admin_demo.py ↗ · simple_consumer.py ↗ · simple_producer.py ↗