主题
第 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 文件;用 Pythonconfluent-kafka写出第一个 Producer / Consumer / AdminClient 程序;并能在 Kafka UI 浏览器里观察集群的实时状态。
0. 导读:命令行是「Kafka 的瑞士军刀」
很多教程一上来就让你写 Java / Python 代码,看似很高大上,其实绕过了最重要的「肌肉记忆」环节:Kafka 自带的 bin/kafka-*.sh 工具是排障、运维、压测、教学的事实标准,不会用命令行 = 没法调试线上集群。
本章我们走「命令行 → Python → 浏览器 UI」三段式:
- 先用命令行手敲一遍核心动作:建 Topic、发消息、收消息、查消费组、重置 Offset、查看 segment 文件。
- 再用 Python 写等价的最小程序,让你从「人发消息」过渡到「程序发消息」。
- 最后用 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.参数解读:
| 参数 | 含义 | 常见取值 |
|---|---|---|
--topic | Topic 名 | 推荐 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:逐字段拆解:
| 字段 | 含义 |
|---|---|
TopicId | KRaft 引入的 UUID,删除重建后也是新 ID(用来防错乱) |
PartitionCount | 分区数 |
ReplicationFactor | 副本数 |
Configs | 当前生效的 Topic 级配置(只列与默认值不同的) |
Partition | 分区号(从 0 起) |
Leader | 当前 Leader 副本所在的 Broker ID |
Replicas | 该分区的所有副本 Broker ID 列表(含 Leader) |
Isr | In-Sync Replicas,与 Leader 保持同步的副本集合 |
Elr / LastKnownElr | Eligible 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输出:
(无输出,命令立即返回)⚠️ 注意几点:
- Broker 必须开启
delete.topic.enable=true(默认即开)。 - 命令是异步的:返回成功 ≠ 数据已经物理删除,要等
log.retention.check.interval.ms(默认 5 分钟)的清理线程扫到。 - 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-Bconsole-consumer-XXXXX 是 kafka-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(即末尾) |
LAG | LOG-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 2STATE 取值:
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 05.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 PnDTnHnMnS | ISO 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 --execute5.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.ms6.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 章细讲)。 - 每条 Record:
offset/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.08.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:
- 点左侧
Topics→ 右上Add a Topic。 - 填
Name、Number of partitions、Replication factor。 - 在
Custom parameters加retention.ms/cleanup.policy等。 - 点
Create Topic,立即在列表中出现。
浏览消息:
- 点 Topic 名 →
MessagesTab。 - 上方有
Seek Type(Offset / Timestamp / Earliest / Latest)和Filters(Key contains / Value contains / Header)。 - 实时滚动展示分区、Offset、Key、Value、Timestamp、Headers。可点单条展开看完整 JSON。
查看消费组:
- 左侧
Consumers→ 选组 → 看到每个 Topic + Partition 的 Lag、当前 Offset、分配的消费者实例。 - 右上
Reset offset一键重置(带 dry-run 视图)。
动态修改配置:
- Topic 详情 →
ConfigsTab → 找到要改的配置 → 点Edit→ 输入新值 →Save。 - 改完会在右侧用红框标出「与默认值不同」的项,与
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 全功能 GUI12. 面试高频题(5 题)
Q1:kafka-topics.sh --create 时 --partitions 和 --replication-factor 怎么选?为什么不能后悔减分区?
考察点:分区设计、Kafka 存储模型。
答案:
--partitions:决定该 Topic 的最大并行消费度与生产者总写入吞吐。- 经验值:
分区数 = 期望峰值吞吐 / 单分区吞吐能力(约 10-50 MB/s),并保留 2-4 倍冗余。 - 太少 → 消费瓶颈、扩消费者无效;太多 → 元数据放大、Controller 压力大、文件句柄爆炸。
- 经验值:
--replication-factor:决定可靠性。生产环境普遍=3(容忍 1 台 Broker 完全失效)。必须 ≤ Broker 数。- 不能减分区的原因:
- Kafka 分区数据是「按分区独立存储 + 按 Key Hash 分布」的,删一个分区意味着该分区所有消息丢失,且其他分区的 Key 分布也要重新映射。
- 加分区也会破坏已有 Key 的顺序保证(同一 Key Hash 后落到的分区编号变了),所以「加分区也是有代价的」,不是无副作用操作。
- 正确做法:上线前充分预估,宁可多开 2 倍分区;如果真要「重新分布」,必须新建 Topic + 双写迁移。
加分项:提到 min.insync.replicas 必须 ≤ replication-factor,且 acks=all + min.insync.replicas=2 才是真正的「不丢消息」组合。
Q2:kafka-console-consumer.sh --from-beginning 一定能从头消费吗?
考察点:auto.offset.reset 与已有 Offset 的关系。
答案:
- 不一定。
--from-beginning等价于--consumer-property auto.offset.reset=earliest。 auto.offset.reset只在「该消费组没有该分区的已提交 Offset」时生效。如果消费组之前已经消费过、提交过 Offset,会直接从已提交位置继续,跟--from-beginning无关。- 真正想从头的方法:
- 不带
--group启动(每次随机生成临时组,没有历史 Offset,肯定从头)。 - 带
--group+ 先用kafka-consumer-groups.sh --reset-offsets --to-earliest --execute重置。 - 改用一个没用过的新 group.id。
- 不带
- 类似地,
auto.offset.reset=latest(默认)会让一个全新组只看到启动后的新消息,「初次上线少了历史」就是这么踩的。 - 第三个取值
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 有哪些坑?
考察点:运维实操、协议细节。
答案:
- 必须
--execute才生效。默认--dry-run只打印「将要改成的样子」。生产上永远先--dry-run再执行。 - 目标消费组必须 inactive(state = Empty / Dead)。否则报
Assignments can only be reset if the group is inactive。要么先停消费者,要么用--all-groups慎用。 - 重置策略多种:
--to-earliest/--to-latest/--to-offset N/--shift-by ±N/--by-duration PT1H/--to-datetime/--from-file CSV。注意--to-datetime是「该时间之后的第一条」,时区按本地时区解析。 - 重置范围:
--topic T是该 topic 全部分区;--topic T:0,2只对 0、2 分区;--all-topics重置该组的所有 topic。 - 重置后不会自动消费,要重新启动消费者才会从新位置 poll。
- 常见用例:
- 业务回放:
--by-duration PT1H --execute重新消费 1 小时数据。 - 跳过坏消息:
--shift-by 1 --execute跳过当前阻塞的一条。 - 全量重算:
--to-earliest --execute。
- 业务回放:
- 替代方案:在程序里调
consumer.seek()或AdminClient.alter_consumer_group_offsets()也能达到一样的效果,且能编程化管控。
加分项:提到 KRaft 模式下重置 Offset 走的是 __consumer_offsets 写入,原理与 ZK 时代一致;但删除消费组在 KRaft 下走元数据日志。
Q4:kafka-dump-log.sh 能告诉你哪些信息?怎么排查「磁盘上消息和我发的不一样」?
考察点:Kafka 存储格式、排障思路。
答案:
- 能看到的信息:
- RecordBatch 级:
baseOffset/lastOffset/count/producerId/producerEpoch/isTransactional/isControl/compresscodec/position(在 .log 文件中的字节偏移)/size/magic(消息格式版本)。 - Record 级:
offset/key/value/headers/timestamp/sequence(幂等用的序列号)。
- RecordBatch 级:
- 常用参数:
--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时用。
- 排障套路「我发的消息和读到的不一样」:
- 找到 Broker 的
log.dirs,定位<topic>-<partition>目录。 - 跑
kafka-dump-log.sh --files xxx.log --print-data-log --deep-iteration,看磁盘上实际存的字节。 - 与 Producer 发送时的
payload对比:- 字节级一致 → 是消费端解码 / 反序列化 / 字符集问题。
- 不一致 → 检查 Producer 端是否做了 SerDe / 拦截器修改;或者中途经过了 Connect / SMT / Streams 处理。
- 关注
compresscodec字段:如果 batch 是 zstd,不带--deep-iteration你看到的是压缩后的二进制乱码。
- 找到 Broker 的
- 进阶:分析事务消息时用
--transaction-log-decoder看__transaction_state内的Ongoing/PrepareCommit/CompleteCommit,能精确还原一次事务的状态机。
Q5:confluent-kafka 的 Producer.flush() 到底等什么?什么时候必须调?
考察点:异步 Producer 模型、可靠性边界。
答案:
Producer.produce()是纯异步入队:消息只是放到 librdkafka 内部的发送队列,立刻返回。不等 broker 任何响应。- librdkafka 内部有个后台 IO 线程根据
linger.ms/batch.size攒批,向 broker 发Produce请求;收到响应后回调你的on_delivery。 flush(timeout)的语义:- 阻塞直到队列里所有未完成消息都收到响应(成功或失败),或超时。
- 期间会触发 delivery callback(
produce时注册的回调)。
- 必须调
flush的场景:- 程序退出前(否则未发完的消息就丢了)。
- 一批关键消息发完后,需要确认结果再继续业务(例如订单写完才允许下游计费)。
- 单元测试里需要同步语义。
- 配套技巧:
flush(0)不等待,只把已就绪的回调触发一遍。- 长生命周期的 Producer 不需要每次都 flush,只在「业务边界」flush;高频 flush 会严重退化为同步发送,吞吐崩盘。
- Python 还要定期
poll(0)(或在每次produce后顺手poll(0))才能让回调被触发;flush内部其实就是循环poll。
- 错误兜底:
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 ↗