这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 环境就绪后,本章用 Kafka 自带命令行完成“建 topic、发消息、收消息、看进度”,并把每个命令的输出读懂。命令行是最好的 Kafka “体检工具”,生产排查也离不开它们。
约定:本章统一使用 --bootstrap-server localhost:9092(伪集群请换成 19092)。
4.1 主题管理
创建一个三分区、双副本的主题:
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--create \
--topic order-events \
--partitions 3 \
--replication-factor 2
查看主题列表与详情:
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--describe --topic order-events
describe 输出解读:
Topic: order-events TopicId: xxxx PartitionCount: 3 ReplicationFactor: 2
Topic: order-events Partition: 0 Leader: 1 Replicas: 1,2 Isr: 1,2
Topic: order-events Partition: 1 Leader: 2 Replicas: 2,3 Isr: 2,3
Topic: order-events Partition: 2 Leader: 3 Replicas: 3,1 Isr: 3,1
Leader:当前处理读写的副本所在 Broker;Replicas:副本分配方案(含不在同步状态的);Isr:同步副本集合。Isr 数量小于 Replicas 数量就是告警信号(第 19 章)。
增加分区(只能增,不能减):
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--alter --topic order-events --partitions 6
删除主题:
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--delete --topic order-events
若删除卡住,通常是 delete.topic.enable=false 或主题正被使用。
4.2 控制台生产者
最简单的生产者:
bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic order-events
逐行输入内容,回车即发送一条消息。Ctrl+C 退出。
带 key 的消息(key 与 value 用分隔符分开):
bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic order-events \
--property parse.key=true \
--property key.separator=:
输入示例:
user-1001:{"type":"login"}
user-1002:{"type":"order"}
user-1001:{"type":"logout"}
相同 key 会进同一分区,这是后面验证顺序性的基础。
4.3 控制台消费者
从头消费所有历史消息:
bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic order-events \
--from-beginning
显示 key、分区与 offset:
bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic order-events \
--from-beginning \
--property print.key=true \
--property print.partition=true \
--property print.offset=true
输出类似:
Partition:0 Offset:5 key:user-1001 value:{"type":"login"}
Partition:2 Offset:9 key:user-1002 value:{"type":"order"}
指定消费组(生产推荐始终带 group):
bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic order-events \
--group demo-group
再次执行同一命令时不会重复消费,因为位移已提交到 __consumer_offsets。
4.4 消费组与位移管理
列出所有消费组:
bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 --list
查看组详情与堆积:
bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe --group demo-group
关键输出:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
demo order-events 0 120 125 5
demo order-events 1 88 88 0
demo order-events 2 61 63 2
CURRENT-OFFSET:组已提交位移;LOG-END-OFFSET:分区最新位移(可近似理解为“写入进度”);LAG:两者之差,即堆积量。这是日常监控最重要的数字。
重置位移(必须先停止该组所有消费者):
bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--group demo-group \
--topic order-events \
--reset-offsets --to-earliest --execute
常用目标还有 --to-latest、--to-offset 100、--shift-by -10、--to-datetime 2026-08-25T00:00:00.000。
查看某主题最新位移:
bin/kafka-get-offsets.sh \
--bootstrap-server localhost:9092 --topic order-events
4.5 动手实验:亲眼看见分区与副本
按下面步骤做一遍,比读十遍文档更有效。
- 创建
demo主题:6 分区、3 副本(伪集群); - 用带 key 的控制台生产者发送 20 条
user-1到user-5的消息; - 用
print.partition=true的消费者观察:相同 key 是否总落在同一分区? describe该主题,记下每个分区的 Leader;- kill 掉某个 Leader 所在 Broker,再
describe:Leader 变了吗?ISR 少了谁? - 重新启动被 kill 的 Broker,观察 ISR 恢复过程;
- 在数据目录中找到
demo-0、demo-1等目录,看看.log、.index、.timeindex文件长什么样。
这个实验覆盖了分区路由、副本、Leader 选举、ISR 与存储布局,是第 10、11 章的预习。
4.6 查看日志文件内部
kafka-dump-log.sh 能把二进制日志转成可读格式:
bin/kafka-dump-log.sh \
--files /tmp/kraft-1/demo-0/00000000000000000000.log \
--print-data-log
你会看到每个 batch 的 offset 范围、压缩算法、时间戳,以及(带 --print-data-log 时)消息内容。这个命令在生产排查“消息到底写没写进去”时非常常用。
查看某个消费组在内部主题中的记录:
bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic __consumer_offsets \
--from-beginning \
--formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter" \
--isolation-level read_committed
能直观看到“提交位移”到底存了什么。
4.7 命令速查表
| 目的 | 命令要点 |
|---|---|
| 建主题 | kafka-topics.sh --create --topic X --partitions N --replication-factor R |
| 查主题 | kafka-topics.sh --describe --topic X |
| 加分区 | kafka-topics.sh --alter --topic X --partitions N |
| 发消息 | kafka-console-producer.sh --topic X |
| 收消息 | kafka-console-consumer.sh --topic X --from-beginning |
| 看组 | kafka-consumer-groups.sh --describe --group G |
| 重置位移 | kafka-consumer-groups.sh --reset-offsets --to-earliest --execute |
| 看最新位移 | kafka-get-offsets.sh --topic X |
| 转储日志 | kafka-dump-log.sh --files <log文件> |
| 选举 Leader | kafka-leader-election.sh --election-type preferred --all-topic-partitions |
本章小结
- 主题创建后
describe是最直观的健康检查:看 Leader、Replicas、ISR; - 控制台生产者/消费者支持 key、分区、offset 打印,适合做实验和临时排查;
- 消费进度看
CURRENT-OFFSET / LOG-END-OFFSET / LAG三件套; - 位移重置前必须停组,
--execute才真正生效; kafka-dump-log.sh可以直接透视磁盘上的消息日志。
思考题
- 为什么增加分区后,相同 key 的消息可能落到新分区,从而导致该 key 的历史顺序被“切断”?
--from-beginning和--group同时使用时,行为是什么?为什么?- 如果
LAG一直是 0 但业务说“没收到消息”,你会按什么顺序排查?