Learn
Kafka/04-cli-tools

命令行工具实战

Kafka 自带的脚本工具是日常开发与排障的瑞士军刀。本章用一条完整的工作流串起最常用的六个工具:建 topic → 生产 → 消费 → 查消费组 → 改配置 → 查水位。建议开着一个本地集群跟着敲。

以下命令假设在 Kafka 安装目录执行,localhost:9092 换成你的地址。

1. kafka-topics.sh:主题管理

# 创建: 3 分区, 2 副本
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic orders --partitions 3 --replication-factor 2
 
# 建 topic 时附带配置
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic events --partitions 6 --replication-factor 3 \
  --config retention.ms=259200000 --config cleanup.policy=delete
 
# 列出所有 topic
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
 
# 查看详情: 分区分布 / Leader / ISR
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders
 
# 扩分区 (只能增不能减!)
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --alter --topic orders --partitions 6
 
# 删除
bin/kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic orders
 
# 只看有问题的分区: 缺副本 / 无 leader
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --under-replicated-partitions
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --unavailable-partitions
⚠️扩分区会打乱 key 的路由

分区从 3 扩到 6 后,同一个 key 的哈希落点会变化——扩容前 user-1 的消息在 p1,扩容后可能进 p4。依赖"同 key 同分区保证顺序"的业务,扩分区等于破坏历史顺序语义,要在业务低峰、消费者清空积压后谨慎操作(详见第 10 章)。

2. kafka-console-producer.sh:命令行生产

# 交互式生产: 每行一条消息, Ctrl-D 退出
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic orders
 
# 发送带 key 的消息: 格式 "key:value"
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 \
  --topic orders \
  --property parse.key=true --property key.separator=:
# 输入示例:
#   user-1:{"order_id":1001,"amount":99.9}
#   user-2:{"order_id":1002,"amount":15.0}
 
# 从文件批量灌入
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 \
  --topic orders < orders.txt
 
# 带可靠性参数
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 \
  --topic orders --request-required-acks all \
  --producer-property enable.idempotence=true

3. kafka-console-consumer.sh:命令行消费

# 从最新位置开始消费 (默认), 只看新消息
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic orders
 
# 从头消费全部历史
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders --from-beginning
 
# 显示 key、时间戳、分区、offset —— 排障必备
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders --from-beginning \
  --property print.key=true --property print.timestamp=true \
  --property print.partition=true --property print.offset=true
 
# 只消费指定分区、从指定 offset 开始
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders --partition 0 --offset 100 --max-messages 10
 
# 以指定消费组身份消费 (会提交 offset, 影响该组进度!)
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders --group inventory-service
💡排障时别用业务的 group.id

不带 --group 时 console consumer 会生成一次性临时组,不影响任何人。带上业务的 group.id 去"看一眼",会把人家的 offset 提交掉,导致业务服务漏消费。查数据永远用 --from-beginning + 临时组。

4. kafka-consumer-groups.sh:消费组管理

这是排查"消息积压"最重要的工具:

# 列出所有消费组
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
 
# 查看消费组详情: 每个分区的进度与 LAG
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group inventory-service

输出解读:

GROUP             TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CONSUMER-ID
inventory-service orders  0          1500            1520            20   consumer-1-xxx
inventory-service orders  1          980             2400            1420 consumer-2-xxx
  • CURRENT-OFFSET:该组已提交的消费位置
  • LOG-END-OFFSET:分区最新消息位置(LEO)
  • LAG = 两者之差,即积压量。p1 积压 1420 条,该看看 consumer-2 是不是卡了。

重置 offset(必须先停掉该组所有消费者):

# 预览: 重置到最早 (--dry-run 只显示不执行)
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group inventory-service --topic orders \
  --reset-offsets --to-earliest --dry-run
 
# 真正执行
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group inventory-service --topic orders \
  --reset-offsets --to-earliest --execute
 
# 其他重置方式
#   --to-latest                    跳到最新, 放弃积压
#   --to-datetime 2026-07-30T00:00:00.000  按时间点
#   --shift-by -100                回退 100 条
#   --to-offset 5000               指定绝对位置
 
# 删除消费组 (组内无活跃成员时)
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --delete --group old-service

5. kafka-configs.sh:动态配置

不重启就能改 topic / broker / 客户端配额等配置:

# 修改 topic 保留时间为 3 天
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --entity-type topics --entity-name orders \
  --add-config retention.ms=259200000
 
# 查看 topic 的非默认配置
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
  --describe --entity-type topics --entity-name orders
 
# 删除覆盖, 恢复默认值
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --entity-type topics --entity-name orders \
  --delete-config retention.ms
 
# 动态调 broker 日志清理线程数 (entity-name 是 node.id)
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --entity-type brokers --entity-name 1 \
  --add-config log.cleaner.threads=2

6. kafka-get-offsets.sh:查询水位

查看每个分区的最早/最新 offset,常用于估算数据量和验证保留策略:

# 最新 offset (LEO), -1 表示 latest
bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 \
  --topic orders --time -1
 
# 最早 offset, -2 表示 earliest
bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 \
  --topic orders --time -2
 
# 按时间戳查: 该时间之后的第一条消息 offset
bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 \
  --topic orders --time 1722297600000

输出 orders:0:1520 表示分区 0 的目标 offset 是 1520。最新减最早即当前分区实际存储的消息条数。

7. 一条完整工作流

T=demo-$(date +%s)
# 1. 建 topic
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic $T --partitions 3 --replication-factor 1
# 2. 灌 5 条带 key 的消息
printf 'a:1\nb:2\na:3\nc:4\nb:5\n' | bin/kafka-console-producer.sh \
  --bootstrap-server localhost:9092 --topic $T \
  --property parse.key=true --property key.separator=:
# 3. 带分区信息消费, 验证同 key 进同分区
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic $T --from-beginning --max-messages 5 \
  --property print.key=true --property print.partition=true
# 4. 查水位
bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic $T --time -1
# 5. 清理
bin/kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic $T

小结

  • kafka-topics.sh:建/查/扩/删 topic,describe 看 Leader 与 ISR
  • kafka-console-producer/consumer.sh:手工收发,排障时打印 key/分区/offset
  • kafka-consumer-groups.sh:看 LAG、重置 offset——积压排查第一入口
  • kafka-configs.sh:动态改配置不用重启
  • kafka-get-offsets.sh:查询分区水位
  • 重置 offset 前必须停消费者;查数据别用业务 group.id
🎯练习
  1. 创建 4 分区 topic,用带 key 生产写入 20 条消息(5 个不同 key),消费时打印分区号,验证同 key 落同分区。
  2. 用一个固定 group.id 消费 10 条后停止,再用 kafka-consumer-groups 查看该组的 CURRENT-OFFSET 与 LAG。
  3. 把该组 offset 重置到 earliest(先 dry-run 再 execute),重新消费验证消息被重复消费。