命令行工具实战
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=true3. 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-service5. 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=26. 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 与 ISRkafka-console-producer/consumer.sh:手工收发,排障时打印 key/分区/offsetkafka-consumer-groups.sh:看 LAG、重置 offset——积压排查第一入口kafka-configs.sh:动态改配置不用重启kafka-get-offsets.sh:查询分区水位- 重置 offset 前必须停消费者;查数据别用业务 group.id
🎯练习
- 创建 4 分区 topic,用带 key 生产写入 20 条消息(5 个不同 key),消费时打印分区号,验证同 key 落同分区。
- 用一个固定 group.id 消费 10 条后停止,再用 kafka-consumer-groups 查看该组的 CURRENT-OFFSET 与 LAG。
- 把该组 offset 重置到 earliest(先 dry-run 再 execute),重新消费验证消息被重复消费。