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

Spring Boot Kafka消费者始终默认使用localhost:9092问题求助

解决Spring Boot Kafka消费者默认使用localhost:9092的问题

问题描述

开发了一个包含Kafka生产者和消费者的Spring Boot应用,两者使用相同的配置属性和bootstrap servers。生产者能正确获取配置的bootstrap服务器地址,但消费者始终默认使用localhost:9092作为bootstrap servers。尝试过添加KafkaAdmin配置,也将配置属性改为spring.kafka.bootstrap-servers,但消费者仍指向localhost:9092,而生产者用相同配置可正常连接指定broker。


相关代码

消费者配置类

@Configuration
@RequiredArgsConstructor
@Slf4j
@EnableKafka
public class KafkaConsumerConfig {
    private final KafkaConfig kafkaConfig;
     Map<String, Object> consumerProperties() {
        log.info("INSIDE consumerProperties:");
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfig.getBootstrapServers());
        log.info("ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG:"+props.get(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG));
        props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaConfig.getConsumerGroupId());
        props.put("ssl.truststore.location", kafkaConfig.getSslTrustStoreLocation());
        props.put("ssl.truststore.password", kafkaConfig.getSslTrustStorePassword());
        props.put("ssl.keystore.location", kafkaConfig.getSslTrustStoreLocation());
        props.put("ssl.keystore.password", kafkaConfig.getSslTrustStorePassword());
        props.put("ssl.key.location", kafkaConfig.getSslTrustStoreLocation());
        props.put("ssl.key.password", kafkaConfig.getSslTrustStorePassword());
        props.put("security.protocol", kafkaConfig.getSecurityProtocol());
        props.put("ssl.endpoint.identification.algorithm", "");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, kafkaConfig.getConsumerAutoOffsetReset());
        log.info("props:"+props);
        return props;
    }
    @Bean
    public ConsumerFactory<String, String> consumerFactory()
    {
        return new DefaultKafkaConsumerFactory<>(consumerProperties());
    }
    @Bean
    public ConcurrentKafkaListenerContainerFactory concurrentKafkaListenerContainerFactory()
    {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.getContainerProperties().setIdleBetweenPolls(10000);
        return factory;
    }
}

KafkaConfig类

@ConfigurationProperties(prefix = "proj.kafka")
@Getter
@Setter
public class KafkaConfig {
  private String schemaRegistryUrl;
  private String bootstrapServers;
  private String securityProtocol;
  private String sslTrustStoreLocation;
  private String sslTrustStorePassword;
  private String consumerGroupId;
  private String consumerEnableAutoCommit;
  private String consumerAutoOffsetReset;
  private String consumerSessionTimeoutMs;
  private String producerRetries;
  private long consumerRetries;
  private long consumerBackoff;
  private String producerMaxInflightConnections;
  private String primaryCluster;
  private String secondaryCluster;
  private long delayMs;
  private long deleteMappingDelay;
  private String groupId;
}

KafkaConsumerService类

@Slf4j
@RequiredArgsConstructor
@Service
public class KafkaConsumerService {
    @Value("${proj.kafka.topic}")
    private String topicName;

    @KafkaListener(topics = "${proj.kafka.topic}", groupId = "${proj.kafka.consumer-group-id}")
    public void consumeMessage(ConsumerRecord<String, String> message) throws InterruptedException {
        log.info("INSIDE consumeMessage:");
        log.info("Tópic:", topicName);
        log.info("Headers:", message.headers());
        log.info("Partion:", message.partition());
        log.info("key:", message.key());
        log.info("Order:", message.value())
        }
}

application.properties配置

proj.kafka.bootstrap-servers=SSL://server1.com:9092,SSL://server2.com:9092,SSL://server3.com:9092

日志信息

消费者日志

2022-10-06 18:10:12.969  INFO 34048 --- [           main] c.v.a.v.e.config.KafkaConsumerConfig     : INSIDE consumerFactory:
2022-10-06 18:10:12.969  INFO 34048 --- [           main] c.v.a.v.e.config.KafkaConsumerConfig     : INSIDE consumerProperties:
2022-10-06 18:10:12.970  INFO 34048 --- [           main] c.v.a.v.e.config.KafkaConsumerConfig     : 
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG:SSL://server1.com:9092,SSL://server2.com:9092,SSL://server3.com:9092

2022-10-06 18:10:14.500  INFO 34048 --- [           main] o.a.k.clients.consumer.ConsumerConfig    : ConsumerConfig values: 
    allow.auto.create.topics = true
    auto.commit.interval.ms = 5000
    auto.offset.reset = latest
    bootstrap.servers = [localhost:9092]
    check.crcs = true
    client.dns.lookup = use_all_dns_ips
    client.id = consumer-proj-1
    client.rack = 
    connections.max.idle.ms = 540000
    default.api.timeout.ms = 60000
    enable.auto.commit = false

生产者日志

2022-10-06 21:09:33.949  INFO 9 --- [nio-8443-exec-3] o.a.k.clients.producer.ProducerConfig    : ProducerConfig values: 
acks = 1
batch.size = 16384
bootstrap.servers = [SSL://server1.com:9092, SSL://server2.com:9092, SSL://server3.com:9092]

解决方案

1. 修正ConcurrentKafkaListenerContainerFactory泛型声明

当前工厂方法未指定泛型,Spring可能无法识别自定义工厂,转而使用默认自动配置的工厂。修改返回类型为带泛型的版本:

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> concurrentKafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.getContainerProperties().setIdleBetweenPolls(10000);
    return factory;
}

2. 显式指定@KafkaListener使用自定义工厂

如果Spring仍未自动选择自定义工厂,在注解中指定容器工厂名称:

@KafkaListener(
    topics = "${proj.kafka.topic}", 
    groupId = "${proj.kafka.consumer-group-id}",
    containerFactory = "concurrentKafkaListenerContainerFactory"
)

3. 禁用Spring Kafka自动配置

若上述方法无效,在启动类中排除默认自动配置类,强制使用自定义配置:

@SpringBootApplication(exclude = KafkaAutoConfiguration.class)
public class YourApplication {
    public static void main(String[] args) {
        SpringApplication.run(YourApplication.class, args);
    }
}

4. 确保KafkaConfig正确注入Spring容器

给KafkaConfig添加@Component注解,或者在启动类添加@EnableConfigurationProperties(KafkaConfig.class),确保配置属性能被正确绑定:

@Component
@ConfigurationProperties(prefix = "proj.kafka")
@Getter
@Setter
public class KafkaConfig {
    // 原有属性
}

5. 修正SSL配置错误

当前代码中ssl.keystore.location和ssl.key.location都复用了信任库路径,这是错误的。需在KafkaConfig中补充对应属性:

private String sslKeystoreLocation;
private String sslKeyLocation;

然后修改消费者配置中的对应代码:

props.put("ssl.keystore.location", kafkaConfig.getSslKeystoreLocation());
props.put("ssl.key.location", kafkaConfig.getSslKeyLocation());

同时在application.properties中添加对应的配置项。


内容的提问来源于stack exchange,提问作者user1326784

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 15:25:19