1. Kafka命令行工具入门指南

第一次接触Kafka命令行工具时,我完全被那一长串参数搞懵了。后来在实际项目中摸爬滚打才发现,这些命令行工具其实是Kafka运维最实用的"瑞士军刀"。不同于图形界面工具,命令行操作不仅响应更快,还能轻松集成到自动化脚本中。今天我就带大家从零开始,手把手掌握这套生产力工具。

Kafka的命令行工具主要存放在安装目录的bin文件夹下,所有脚本都以.sh结尾。最常用的有三个:kafka-topics.sh负责主题管理,kafka-console-producer.sh用于生产数据,kafka-console-consumer.sh则用来消费数据。这三个工具配合使用,就能完成从主题创建到消息收发的完整流程。建议新手先在测试环境练习这些命令,毕竟生产环境的操作需要更加谨慎。

2. 主题管理全攻略

2.1 查看主题信息

查看主题是最基础的操作,但很多人不知道其中的门道。执行bin/kafka-topics.sh --list --bootstrap-server localhost:9092可以列出所有主题,但更实用的方式是加上--exclude-internal参数过滤掉内部主题:

bin/kafka-topics.sh --list \
    --bootstrap-server localhost:9092 \
    --exclude-internal

如果想查看主题详情,--describe参数会给你惊喜。它会显示分区数、副本分布等关键信息。我常用这个命令检查分区是否均衡:

bin/kafka-topics.sh --describe \
    --topic my_topic \
    --bootstrap-server localhost:9092

输出中的Leader表示负责读写的节点,Replicas是所有副本所在节点,ISR是保持同步的副本集合。当发现ISR数量小于副本数时,就说明有节点掉队了,需要及时排查。

2.2 创建主题的学问

创建主题看似简单,但参数设置直接影响后续性能。这个命令我至少用过上百次:

bin/kafka-topics.sh --create \
    --bootstrap-server localhost:9092 \
    --replication-factor 3 \
    --partitions 6 \
    --topic orders

这里有三个关键点需要注意:

  1. 副本数通常设为3,保证高可用但又不至于过多影响性能
  2. 分区数要根据吞吐量预估,我一般按每秒万条消息配10个分区
  3. 主题命名避免使用点号或下划线混用,可能引发监控指标问题

2.3 主题配置调整

实际运维中经常需要调整主题配置。比如发现某个主题流量激增,可以通过增加分区来提升吞吐:

bin/kafka-topics.sh --alter \
    --bootstrap-server localhost:9092 \
    --topic orders \
    --partitions 12

但要注意两点:

  1. 分区只能增加不能减少
  2. 对于有key的消息,增加分区会影响消息顺序

删除主题前务必三思,这个操作不可逆:

bin/kafka-topics.sh --delete \
    --bootstrap-server localhost:9092 \
    --topic temp_topic

如果发现主题没被真正删除,记得检查server.properties中delete.topic.enable=true是否设置。

3. 数据生产实战技巧

3.1 基础消息生产

控制台生产者是我调试时的好帮手,这个命令简单却实用:

bin/kafka-console-producer.sh \
    --broker-list localhost:9092 \
    --topic logs \
    --property "parse.key=true" \
    --property "key.separator=:"

输入消息时可以指定key,这对分区选择很重要:

user123:{"action":"login","time":"2023-01-01"}

--property参数还有很多实用配置,比如设置消息压缩方式、调整批次大小等。生产环境建议开启压缩:

--compression-codec snappy

3.2 高级生产配置

通过配置文件可以设置更复杂的参数,这是我常用的producer.properties:

acks=all
retries=10
max.in.flight.requests.per.connection=1
compression.type=snappy
linger.ms=20
batch.size=65536

使用时指定配置文件:

bin/kafka-console-producer.sh \
    --broker-list localhost:9092 \
    --topic important_data \
    --producer.config producer.properties

其中acks=all确保消息被所有副本确认,适合对可靠性要求高的场景。如果追求吞吐量,可以设为1或0。

4. 数据消费完全手册

4.1 基础消费操作

最基础的消费命令是从最新位置开始:

bin/kafka-console-consumer.sh \
    --bootstrap-server localhost:9092 \
    --topic orders

但调试时我更喜欢从头消费:

--from-beginning

对于大流量主题,建议用--max-messages限制数量,避免刷屏:

--max-messages 100

4.2 消费组管理

创建消费组是实际项目中的常规操作:

bin/kafka-console-consumer.sh \
    --bootstrap-server localhost:9092 \
    --topic payments \
    --group payment_processor

查看消费组情况用这个命令:

bin/kafka-consumer-groups.sh \
    --bootstrap-server localhost:9092 \
    --describe \
    --group payment_processor

输出中的CURRENT-OFFSET表示已消费位置,LOG-END-OFFSET是最新位置,LAG就是堆积量。当LAG持续增长时,就要考虑扩容消费者了。

4.3 消费偏移量控制

重置偏移量是处理消息堆积的终极手段,但需谨慎:

bin/kafka-consumer-groups.sh \
    --bootstrap-server localhost:9092 \
    --group payment_processor \
    --reset-offsets \
    --to-earliest \
    --execute \
    --topic payments

其他常用重置选项:

  • --to-latest 跳到最新位置
  • --to-datetime 2023-01-01T00:00:00.000 跳到指定时间
  • --shift-by -1000 往回调整1000条

5. 运维监控与排错

5.1 集群健康检查

这个组合命令是我的每日必查:

# 检查控制器
bin/kafka-topics.sh --describe \
    --bootstrap-server localhost:9092 \
    --under-replicated-partitions

# 检查ISR
bin/kafka-topics.sh --describe \
    --bootstrap-server localhost:9092 \
    | grep -v "Isr: 1,0,2"

发现under-replicated或ISR不完整时,通常意味着有节点故障或网络问题。

5.2 消息内容检查

当怀疑消息有问题时,可以用这个技巧查看原始数据:

bin/kafka-console-consumer.sh \
    --bootstrap-server localhost:9092 \
    --topic problem_topic \
    --from-beginning \
    --formatter "kafka.tools.DefaultMessageFormatter" \
    --property print.key=true \
    --property print.value=true \
    --property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \
    --property value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

5.3 性能测试工具

Kafka自带的性能测试工具很实用:

# 生产者测试
bin/kafka-producer-perf-test.sh \
    --topic benchmark \
    --num-records 1000000 \
    --record-size 1024 \
    --throughput -1 \
    --producer-props bootstrap.servers=localhost:9092

# 消费者测试
bin/kafka-consumer-perf-test.sh \
    --topic benchmark \
    --bootstrap-server localhost:9092 \
    --messages 1000000

测试结果主要看吞吐量(records/sec)和延迟(ms),这些数据对容量规划很有帮助。

6. 安全配置实践

6.1 SSL加密配置

在生产环境必须配置SSL加密,先准备server.properties:

listeners=SSL://:9093
ssl.keystore.location=/path/to/keystore.jks
ssl.keystore.password=keystore_password
ssl.key.password=key_password
ssl.truststore.location=/path/to/truststore.jks
ssl.truststore.password=truststore_password
ssl.client.auth=required

客户端连接时需要指定SSL配置:

bin/kafka-console-producer.sh \
    --broker-list localhost:9093 \
    --topic secure_topic \
    --producer.config ssl.properties

6.2 SASL认证配置

对于需要认证的场景,配置SASL_PLAINTEXT:

listeners=SASL_PLAINTEXT://:9094
sasl.mechanism.inter.broker.protocol=PLAIN
sasl.enabled.mechanisms=PLAIN

客户端需要提供jaas配置:

sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required \
    username="admin" \
    password="admin-secret";

7. 高级主题管理

7.1 分区重平衡

当新增节点后,需要手动触发分区重平衡:

bin/kafka-reassign-partitions.sh \
    --bootstrap-server localhost:9092 \
    --reassignment-json-file reassign.json \
    --execute

reassign.json文件格式如下:

{
  "partitions": [
    {
      "topic": "big_topic",
      "partition": 0,
      "replicas": [3,4,5]
    }
  ]
}

7.2 日志清理策略

对于数据量大的主题,配置日志清理策略很重要:

bin/kafka-configs.sh --alter \
    --bootstrap-server localhost:9092 \
    --entity-type topics \
    --entity-name sensor_data \
    --add-config cleanup.policy=compact,segment.ms=3600000

compact策略适合key-value型数据,能自动去重。retention.ms则控制数据保留时间。

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐