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 130e975028c..167ea7758d4 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 @@ -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()); } 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 f4cf275aa63..7e8bf4cf386 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 @@ -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 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) -> {