mirror of
https://github.com/spring-projects/spring-boot.git
synced 2026-09-29 21:59:02 +00:00
Apply consumer-specific security protocol
Closes gh-51365
This commit is contained in:
+1
-1
@@ -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());
|
||||
}
|
||||
|
||||
|
||||
+31
-1
@@ -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<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) -> {
|
||||
|
||||
Reference in New Issue
Block a user