You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot KafkaTemplate发送成功但主题无消息问题排查

Kafka发送成功但消费者无法读取消息的排查与解决

问题背景

使用Apache Kafka 4.1.2(KRaft单节点Docker部署)+ Spring Boot 4.0.6,向local-health-status-v0主题发送消息:

  • 主题通过NewTopic配置创建成功,可通过kafka-topics.sh查询到
  • KafkaTemplate发送日志显示成功,返回了Producer Record和Record Metadata
  • 执行kafka-console-consumer.sh --topic local-health-status-v0 --from-beginning --bootstrap-server localhost:9092时,始终显示Processed a total of 0 messages
  • Kafka日志反复出现INFO Sent auto-creation request for Set(__consumer_offsets) to the active controller.

核心原因

单节点KRaft模式下,Kafka内部主题(如__consumer_offsets)默认副本数为3,但单节点无法满足多副本部署要求,导致内部主题创建失败。消费者依赖__consumer_offsets存储消费偏移量,内部主题异常会导致消费者无法正常初始化或读取消息;同时如果自定义主题副本数配置大于1,也会引发副本同步问题。

解决方案

1. 修改Kafka Docker启动参数,适配单节点环境

在启动命令中添加内部主题副本数配置,确保内部主题能正常创建:

KAFKA_VERSION="4.1.2"
CONTAINER_NAME="kafka-beaufort"
CLUSTER_ID="4L6G9nxwQreSVeTuSsh_Hg"

docker run -d \
  --name $CONTAINER_NAME \
  -p 9092:9092 \
  -e KAFKA_PROCESS_ROLES=broker,controller \
  -e KAFKA_NODE_ID=1 \
  -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@localhost:9093 \
  -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \
  -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \
  -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
  -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \
  -e CLUSTER_ID=$CLUSTER_ID \
  # 添加以下适配单节点的配置
  -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
  -e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 \
  -e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 \
  apache/kafka:$KAFKA_VERSION

注意:需要先停止并删除旧容器,清理关联的数据卷(如果有),再重新启动。

2. 确认自定义主题副本数匹配单节点环境

检查NewTopic配置中的replicas参数,确保healthTopicReplicas的值为1:

@Bean
public NewTopic healthStatusTopic() {
   String healthStatusTopicName = String.format(HEALTH_STATUS_TOPIC_FMT, this.environmentService.topicEnvironment(), this.environmentService.version());

   return TopicBuilder.name(healthStatusTopicName)
      .partitions(this.healthTopicPartitions)
      .replicas(1) // 单节点环境下副本数必须设为1
      .build();
}

3. 验证生产者与消费者配置

  • 生产者:确认spring.kafka.producer.acks配置(默认值为1),如果设置为all,需同步设置主题的min.insync.replicas为1,避免因副本同步问题导致消息无法确认。
  • 消费者:使用控制台消费者时,可显式指定group.id避免自动生成的组偏移量干扰,命令示例:
    kafka-console-consumer.sh --topic local-health-status-v0 --from-beginning --bootstrap-server localhost:9092 --group test-consumer-group
    

验证步骤

  1. 重新启动Kafka容器后,用kafka-topics.sh --list --bootstrap-server localhost:9092确认__consumer_offsets主题已创建
  2. 重新发送消息,再用控制台消费者读取,确认能正常获取消息

内容的提问来源于stack exchange,提问作者Marc Le Bihan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.05 00:07:31