Spring Boot 3.2.6与Spring Cloud 2023.0.2中Kafka Producer的bootstrap.servers配置问题
环境
- Spring Boot:3.2.6
- Spring Cloud:2023.0.2
- Spring Cloud Config Server
问题描述
升级前(Spring Boot 2.7.8 + Spring Cloud 2021.0.8),org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration类的kafkaProducerFactory方法直接通过properties.buildProducerProperties()创建DefaultKafkaProducerFactory,能正确从Spring Config Server加载bootstrap.servers属性。
旧版本核心代码:
@Bean @ConditionalOnMissingBean(ProducerFactory.class) public DefaultKafkaProducerFactory<?, ?> kafkaProducerFactory( ObjectProvider<DefaultKafkaProducerFactoryCustomizer> customizers) { DefaultKafkaProducerFactory<?, ?> factory = new DefaultKafkaProducerFactory<>( this.properties.buildProducerProperties()); String transactionIdPrefix = this.properties.getProducer().getTransactionIdPrefix(); if (transactionIdPrefix != null) { factory.setTransactionIdPrefix(transactionIdPrefix); } customizers.orderedStream().forEach((customizer) -> customizer.customize(factory)); return factory; }
升级后,kafkaProducerFactory方法新增KafkaConnectionDetails参数,且新增applyKafkaConnectionDetailsForProducer调用,会将从Config Server加载的bootstrap.servers覆盖为默认的localhost:9092。
新版本核心代码:
@Bean @ConditionalOnMissingBean({ProducerFactory.class}) public DefaultKafkaProducerFactory<?, ?> kafkaProducerFactory(KafkaConnectionDetails connectionDetails, ObjectProvider<DefaultKafkaProducerFactoryCustomizer> customizers, ObjectProvider<SslBundles> sslBundles) { Map<String, Object> properties = this.properties.buildProducerProperties((SslBundles)sslBundles.getIfAvailable()); this.applyKafkaConnectionDetailsForProducer(properties, connectionDetails); DefaultKafkaProducerFactory<?, ?> factory = new DefaultKafkaProducerFactory(properties); String transactionIdPrefix = this.properties.getProducer().getTransactionIdPrefix(); if (transactionIdPrefix != null) { factory.setTransactionIdPrefix(transactionIdPrefix); } customizers.orderedStream().forEach((customizer) -> { customizer.customize(factory); }); return factory; }
观察到properties.buildProducerProperties已正确填充bootstrap.servers为my.eastus2.azure.confluent.cloud:9092,但后续被applyKafkaConnectionDetailsForProducer覆盖为默认值。
配置片段(旧版本有效,新版本失效)
# application.yml spring: config: import: optional:configserver:https://my.config.server.com/v1/config/ kafka: producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer properties: ### Auth ### sasl.jaas.config: "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required \ clientId='${my_client_id}' \ clientSecret='${my_client_secret}' \ scope='' \ extension_logicalCluster='${my_kafka_cluster_id}' \ extension_identityPoolId='${my_identity_pool}';" ### Schema Registry ### bearer.auth.client.id: ${my_client_id} bearer.auth.client.secret: ${my_client_secret} bearer.auth.identity.pool.id: ${my_identity_pool:pool-xyz} kekcipherkms.provider.identity.client.id: ${my_client_id} kekcipherkms.provider.identity.client.secret: ${my_client_secret} kekcipherkms.provider.identity.scope: my.scope.w producer: properties: topic: ${topic}
观察到的行为
- 旧版本:
DefaultKafkaProducerFactory正确使用Config Server提供的bootstrap.servers值 - 新版本:
bootstrap.servers被覆盖为localhost:9092
期望行为
DefaultKafkaProducerFactory继续使用Spring Config Server提供的bootstrap.servers属性
问题
如何配置Spring Boot应用,使其在新版本中正确使用Spring Config Server的bootstrap.servers值,避免被KafkaConnectionDetails的默认值覆盖?
解决方案
方法1:适配KafkaConnectionDetails属性规则
Spring Boot 3.x引入的KafkaConnectionDetails优先级高于传统spring.kafka.*配置,需确保Config Server中使用正确的属性名:
- 将原
spring.kafka.bootstrap.servers(点分隔)改为spring.kafka.bootstrap-servers(短横线分隔),或对应环境变量SPRING_KAFKA_BOOTSTRAP_SERVERS - 启动时查看日志,确认Config Server返回的配置中
bootstrap-servers已被正确加载
方法2:自定义KafkaConnectionDetails Bean
通过自定义Bean覆盖默认实现,直接复用传统配置值:
import org.springframework.boot.autoconfigure.kafka.KafkaConnectionDetails; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.List; @Configuration public class CustomKafkaConfig { @Bean public KafkaConnectionDetails kafkaConnectionDetails(KafkaProperties kafkaProperties) { return new KafkaConnectionDetails() { @Override public List<String> getBootstrapServers() { return kafkaProperties.getBootstrapServers(); } }; } }
方法3:禁用KafkaConnectionDetails自动配置
在application.yml中添加配置,回到旧版本配置逻辑:
spring: autoconfigure: exclude: org.springframework.boot.autoconfigure.kafka.KafkaConnectionDetailsAutoConfiguration
注意:此方法可能影响SSL bundles等新特性集成,需谨慎使用
方法4:通过自定义配置器修正属性
使用DefaultKafkaProducerFactoryCustomizer在工厂创建后重置bootstrap.servers:
import org.springframework.boot.autoconfigure.kafka.DefaultKafkaProducerFactoryCustomizer; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.stereotype.Component; @Component public class KafkaProducerConfigCustomizer implements DefaultKafkaProducerFactoryCustomizer { private final KafkaProperties kafkaProperties; public KafkaProducerConfigCustomizer(KafkaProperties kafkaProperties) { this.kafkaProperties = kafkaProperties; } @Override public void customize(DefaultKafkaProducerFactory<?, ?> factory) { factory.getConfigurationProperties().put("bootstrap.servers", String.join(",", kafkaProperties.getBootstrapServers())); } }
内容的提问来源于stack exchange,提问作者Fisher_Coder

