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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 19:05:05