KafkaNotes

第 04 章:第一条消息:命令行实战

zjc 于 2026-01-04 发布

这是《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

增加分区(只能增,不能减):

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

重置位移(必须先停止该组所有消费者):

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 动手实验:亲眼看见分区与副本

按下面步骤做一遍,比读十遍文档更有效。

  1. 创建 demo 主题:6 分区、3 副本(伪集群);
  2. 用带 key 的控制台生产者发送 20 条 user-1user-5 的消息;
  3. print.partition=true 的消费者观察:相同 key 是否总落在同一分区?
  4. describe 该主题,记下每个分区的 Leader;
  5. kill 掉某个 Leader 所在 Broker,再 describe:Leader 变了吗?ISR 少了谁?
  6. 重新启动被 kill 的 Broker,观察 ISR 恢复过程;
  7. 在数据目录中找到 demo-0demo-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

本章小结

思考题

  1. 为什么增加分区后,相同 key 的消息可能落到新分区,从而导致该 key 的历史顺序被“切断”?
  2. --from-beginning--group 同时使用时,行为是什么?为什么?
  3. 如果 LAG 一直是 0 但业务说“没收到消息”,你会按什么顺序排查?