From 46aea174a3b64e61257a6fc652c2ba5f6d042032 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Sat, 23 Feb 2019 11:36:34 -0500 Subject: [PATCH] Fix Consumer Property Overrides If the consumer properties object contains `default` properties, they were ignored. `Properties.forEach()` does not return default properties. Use `stringPropertyNames()` instead. **cherry-pick to 2.2.x** --- .../kafka/core/DefaultKafkaConsumerFactory.java | 14 ++++++++------ .../listener/AbstractMessageListenerContainer.java | 3 ++- .../kafka/listener/ContainerProperties.java | 4 ++-- .../KafkaMessageListenerContainerTests.java | 9 +++++++-- 4 files changed, 19 insertions(+), 11 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java index 15e10a30..7bf45e3b 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaConsumerFactory.java @@ -126,7 +126,9 @@ public class DefaultKafkaConsumerFactory implements ConsumerFactory } boolean shouldModifyClientId = (this.configs.containsKey(ConsumerConfig.CLIENT_ID_CONFIG) && StringUtils.hasText(clientIdSuffix)) || overrideClientIdPrefix; - if (groupId == null && properties == null && !shouldModifyClientId) { + if (groupId == null + && (properties == null || properties.stringPropertyNames().size() == 0) + && !shouldModifyClientId) { return createKafkaConsumer(this.configs); } else { @@ -149,11 +151,11 @@ public class DefaultKafkaConsumerFactory implements ConsumerFactory : modifiedConfigs.get(ConsumerConfig.CLIENT_ID_CONFIG)) + clientIdSuffix); } if (properties != null) { - properties.forEach((k, v) -> { - if (!k.equals(ConsumerConfig.CLIENT_ID_CONFIG) && !k.equals(ConsumerConfig.GROUP_ID_CONFIG)) { - modifiedConfigs.put((String) k, v); - } - }); + properties.stringPropertyNames() + .stream() + .filter(name -> !name.equals(ConsumerConfig.CLIENT_ID_CONFIG) + && !name.equals(ConsumerConfig.GROUP_ID_CONFIG)) + .forEach(name -> modifiedConfigs.put(name, properties.getProperty(name))); } return createKafkaConsumer(modifiedConfigs); } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java index b694b5f8..123620ea 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java @@ -282,7 +282,8 @@ public abstract class AbstractMessageListenerContainer if (this.containerProperties.isMissingTopicsFatal() && this.containerProperties.getTopicPattern() == null) { try (Consumer consumer = this.consumerFactory.createConsumer(this.containerProperties.getGroupId(), - this.containerProperties.getClientId(), null)) { + this.containerProperties.getClientId(), null, + this.containerProperties.getConsumerProperties())) { if (consumer != null) { String[] topics = this.containerProperties.getTopics(); if (topics == null) { diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ContainerProperties.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ContainerProperties.java index 88e81db7..52b5c897 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ContainerProperties.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ContainerProperties.java @@ -617,7 +617,7 @@ public class ContainerProperties { * name(s) in the consumer factory. * {@code group.id} and {@code client.id} are ignored. * @return the properties. - * @since 2.1.4 + * @since 2.2.4 * @see org.apache.kafka.clients.consumer.ConsumerConfig * @see #setGroupId(String) * @see #setClientId(String) @@ -633,7 +633,7 @@ public class ContainerProperties { * name(s) in the consumer factory. * {@code group.id} and {@code client.id} are ignored. * @param consumerProperties the properties. - * @since 2.1.4 + * @since 2.2.4 * @see org.apache.kafka.clients.consumer.ConsumerConfig * @see #setGroupId(String) * @see #setClientId(String) diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java index e69ed227..a5425410 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java @@ -1213,8 +1213,8 @@ public class KafkaMessageListenerContainerTests { DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory(props) { @Override - public Consumer createConsumer(String groupId, String clientIdPrefix, - String clientIdSuffix, Properties properties) { + protected KafkaConsumer createKafkaConsumer(Map configs) { + assertThat(configs).containsKey(ConsumerConfig.MAX_POLL_RECORDS_CONFIG); return new KafkaConsumer(props) { @Override @@ -1238,6 +1238,10 @@ public class KafkaMessageListenerContainerTests { logger.info("defined part: " + message); latch1.countDown(); }); + Properties defaultProperties = new Properties(); + defaultProperties.setProperty(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "42"); + Properties consumerProperties = new Properties(defaultProperties); + container1Props.setConsumerProperties(consumerProperties); CountDownLatch stubbingComplete1 = new CountDownLatch(1); KafkaMessageListenerContainer container1 = spyOnContainer( new KafkaMessageListenerContainer<>(cf, container1Props), stubbingComplete1); @@ -1266,6 +1270,7 @@ public class KafkaMessageListenerContainerTests { logger.info("defined part: " + message); latch2.countDown(); }); + container2Props.setConsumerProperties(consumerProperties); CountDownLatch stubbingComplete2 = new CountDownLatch(1); KafkaMessageListenerContainer container2 = spyOnContainer( new KafkaMessageListenerContainer<>(cf, container2Props), stubbingComplete2);