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
相关产品推荐
相关产品推荐

