From ada7ac7c4847aaa02ffb30c9559ca70a7568865b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?St=C3=A9phane=20Nicoll?= Date: Tue, 11 Aug 2026 17:46:39 +0200 Subject: [PATCH] Apply consumer-specific security protocol Closes gh-51365 --- .../autoconfigure/KafkaAutoConfiguration.java | 2 +- .../KafkaAutoConfigurationTests.java | 32 ++++++++++++++++++- 2 files changed, 32 insertions(+), 2 deletions(-) diff --git a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java index 4e370eeb341..441ebec4cc9 100644 --- a/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java +++ b/module/spring-boot-kafka/src/main/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfiguration.java @@ -206,7 +206,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()); } diff --git a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java index d6b87983c36..ab9be0b8f74 100644 --- a/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java +++ b/module/spring-boot-kafka/src/test/java/org/springframework/boot/kafka/autoconfigure/KafkaAutoConfigurationTests.java @@ -993,7 +993,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 consumerConfigs = consumerFactory.getConfigurationProperties(); + assertThat(consumerConfigs).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL"); + DefaultKafkaProducerFactory producerFactory = context.getBean(DefaultKafkaProducerFactory.class); + Map 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 producerConfigs = producerFactory.getConfigurationProperties(); + assertThat(producerConfigs).containsEntry(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL"); + DefaultKafkaConsumerFactory consumerFactory = context.getBean(DefaultKafkaConsumerFactory.class); + Map 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) -> {