Merge branch '4.1.x'

Closes gh-51370
This commit is contained in:
Stéphane Nicoll
2026-08-11 17:49:09 +02:00
2 changed files with 32 additions and 2 deletions
@@ -213,7 +213,7 @@ public final class KafkaAutoConfiguration {
KafkaConnectionDetails connectionDetails) {
Configuration consumer = connectionDetails.getConsumer();
properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, consumer.getBootstrapServers());
applySecurityProtocol(properties, connectionDetails.getSecurityProtocol());
applySecurityProtocol(properties, consumer.getSecurityProtocol());
applySslBundle(properties, consumer.getSslBundle());
}
@@ -1037,7 +1037,37 @@ class KafkaAutoConfigurationTests {
}
@Test
void specificSecurityProtocolOverridesCommonSecurityProtocol() {
void specificProducerSecurityProtocolOverridesCommonSecurityProtocol() {
this.contextRunner
.withPropertyValues("spring.kafka.security.protocol=SSL",
"spring.kafka.producer.security.protocol=PLAINTEXT")
.run((context) -> {
DefaultKafkaConsumerFactory<?, ?> consumerFactory = context.getBean(DefaultKafkaConsumerFactory.class);
Map<String, Object> consumerConfigs = consumerFactory.getConfigurationProperties();
assertThat(consumerConfigs).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL");
DefaultKafkaProducerFactory<?, ?> producerFactory = context.getBean(DefaultKafkaProducerFactory.class);
Map<String, Object> producerConfigs = producerFactory.getConfigurationProperties();
assertThat(producerConfigs).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "PLAINTEXT");
});
}
@Test
void specificConsumerSecurityProtocolOverridesCommonSecurityProtocol() {
this.contextRunner
.withPropertyValues("spring.kafka.security.protocol=SSL",
"spring.kafka.consumer.security.protocol=PLAINTEXT")
.run((context) -> {
DefaultKafkaProducerFactory<?, ?> producerFactory = context.getBean(DefaultKafkaProducerFactory.class);
Map<String, Object> producerConfigs = producerFactory.getConfigurationProperties();
assertThat(producerConfigs).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL");
DefaultKafkaConsumerFactory<?, ?> consumerFactory = context.getBean(DefaultKafkaConsumerFactory.class);
Map<String, Object> consumerConfigs = consumerFactory.getConfigurationProperties();
assertThat(consumerConfigs).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "PLAINTEXT");
});
}
@Test
void specificAdminSecurityProtocolOverridesCommonSecurityProtocol() {
this.contextRunner
.withPropertyValues("spring.kafka.security.protocol=SSL", "spring.kafka.admin.security.protocol=PLAINTEXT")
.run((context) -> {