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
验证步骤
- 重新启动Kafka容器后,用
kafka-topics.sh --list --bootstrap-server localhost:9092确认__consumer_offsets主题已创建 - 重新发送消息,再用控制台消费者读取,确认能正常获取消息
内容的提问来源于stack exchange,提问作者Marc Le Bihan
相关产品推荐
相关产品推荐

