SpringBoot Kafka生产者无法发送消息,Admin Client已注销求助
Spring Boot集成Kafka消息发送失败问题排查与解决
问题描述
开发Spring Boot集成Kafka的消息发送应用,启动后无法完成消息发送,日志输出如下:
Kafka commitId: 5e3c2b738d253ff5
Kafka startTimeMs: 1712071898984
App info kafka.admin.client for adminclient-1 unregistered
代码示例
主启动类
@SpringBootApplication public class SpringKafkaApplication { public static void main(String[] args) { SpringApplication.run(SpringKafkaApplication.class, args); } @Autowired private KafkaTemplate<String, String> kafkaTemplate; private String ms = "Hello Kafka"; public void sendMessage(String msg) { kafkaTemplate.send("main", msg); } }
Kafka主题配置类
@Configuration public class KafkaTopicConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServers; @Bean public KafkaAdmin kafkaAdmin() { Map<String, Object> configs = new HashMap<>(); configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); return new KafkaAdmin(configs); } @Bean public NewTopic topic1() { return new NewTopic("main", 1, (short) 1); } }
Kafka生产者配置类
@Configuration public class KafkaProducerConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServers; public Map<String,Object> producerConfig() { Map<String,Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return props; } @Bean ProducerFactory<String, String> producerFactory(){ return new DefaultKafkaProducerFactory<>(producerConfig()); } @Bean KafkaTemplate<String, String> kafkaTemplate(){ return new KafkaTemplate<>(producerFactory()); } }
application.properties
spring.kafka.bootstrap-servers=localhost:9092
问题分析与修复
1. 核心问题定位
日志中的App info kafka.admin.client for adminclient-1 unregistered是正常生命周期日志——KafkaAdmin在完成主题创建后会自动关闭AdminClient,并非报错。真正导致消息未发送的原因是:
- 主启动类仅定义了
sendMessage方法,但没有触发执行的逻辑,应用启动后不会主动调用该方法发送消息; - 生产者配置中硬编码了
bootstrap.servers,虽当前值与配置文件一致,但不利于后续维护。
2. 修复步骤
(1)添加消息发送触发逻辑
修改主启动类,实现CommandLineRunner接口,在应用启动时自动调用发送方法,同时添加回调查看发送结果:
@SpringBootApplication public class SpringKafkaApplication implements CommandLineRunner { @Autowired private KafkaTemplate<String, String> kafkaTemplate; private String ms = "Hello Kafka"; public static void main(String[] args) { SpringApplication.run(SpringKafkaApplication.class, args); } public void sendMessage(String msg) { kafkaTemplate.send("main", msg) .addCallback( success -> System.out.println("消息发送成功:" + success.getRecordMetadata()), failure -> System.err.println("消息发送失败:" + failure.getMessage()) ); } @Override public void run(String... args) throws Exception { sendMessage(ms); } }
(2)统一配置来源
修改KafkaProducerConfig的producerConfig方法,使用注入的配置值替代硬编码:
public Map<String,Object> producerConfig() { Map<String,Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return props; }
(3)验证Kafka服务状态
确保本地Kafka Broker(localhost:9092)处于运行状态,可通过以下命令验证主题main是否存在:
kafka-topics.sh --list --bootstrap-server localhost:9092
3. 验证结果
启动应用后,查看控制台输出:
- 若显示
消息发送成功:...,则说明问题解决; - 若显示失败信息,可根据回调中的错误提示进一步排查(如网络连接、主题权限、Broker配置等)。
内容的提问来源于stack exchange,提问作者Sky
相关产品推荐
相关产品推荐

