diff --git a/build.gradle b/build.gradle index 42705689..38305523 100644 --- a/build.gradle +++ b/build.gradle @@ -78,7 +78,7 @@ subprojects { subproject -> mockitoVersion = '1.9.5' scalaVersion = '2.11' springRetryVersion = '1.1.2.RELEASE' - springVersion = '4.2.5.RELEASE' + springVersion = '4.2.6.RELEASE' idPrefix = 'kafka' @@ -154,6 +154,7 @@ project ('spring-kafka-test') { description = 'Spring Kafka Test Support' dependencies { + compile "org.springframework:spring-beans:$springVersion" compile "org.springframework:spring-test:$springVersion" compile "org.springframework.retry:spring-retry:$springRetryVersion" compile "org.apache.kafka:kafka_$scalaVersion:$kafkaVersion" diff --git a/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java b/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java index e8aa1cff..fa4d492c 100644 --- a/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java +++ b/spring-kafka-test/src/main/java/org/springframework/kafka/test/utils/KafkaTestUtils.java @@ -31,7 +31,9 @@ import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; +import org.springframework.beans.DirectFieldAccessor; import org.springframework.kafka.test.rule.KafkaEmbedded; +import org.springframework.util.Assert; /** * Kafka testing utilities. @@ -129,4 +131,47 @@ public final class KafkaTestUtils { return received; } + /** + * Uses nested {@link DirectFieldAccessor}s to obtain a property using dotted notation to traverse fields; e.g. + * "foo.bar.baz" will obtain a reference to the baz field of the bar field of foo. Adopted from Spring Integration. + * @param root The object. + * @param propertyPath The path. + * @return The field. + */ + public static Object getPropertyValue(Object root, String propertyPath) { + Object value = null; + DirectFieldAccessor accessor = new DirectFieldAccessor(root); + String[] tokens = propertyPath.split("\\."); + for (int i = 0; i < tokens.length; i++) { + value = accessor.getPropertyValue(tokens[i]); + if (value != null) { + accessor = new DirectFieldAccessor(value); + } + else if (i == tokens.length - 1) { + return null; + } + else { + throw new IllegalArgumentException("intermediate property '" + tokens[i] + "' is null"); + } + } + return value; + } + + /** + * A typed version of {@link #getPropertyValue(Object, String)}. + * @param root the object. + * @param propertyPath the path. + * @param type the type to cast the object to + * @return the field value. + * @see #getPropertyValue(Object, String) + */ + @SuppressWarnings("unchecked") + public static T getPropertyValue(Object root, String propertyPath, Class type) { + Object value = getPropertyValue(root, propertyPath); + if (value != null) { + Assert.isAssignable(type, value.getClass()); + } + return (T) value; + } + } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerContainerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerContainerFactory.java index fc2d4bc5..a69a33c4 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerContainerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/SimpleKafkaListenerContainerFactory.java @@ -19,6 +19,7 @@ package org.springframework.kafka.config; import java.util.Collection; import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; +import org.apache.kafka.clients.consumer.OffsetCommitCallback; import org.apache.kafka.common.TopicPartition; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; @@ -47,6 +48,10 @@ public class SimpleKafkaListenerContainerFactory private ConsumerRebalanceListener consumerRebalanceListener; + private OffsetCommitCallback commitCallback; + + private Boolean syncCommits; + /** * Specify the container concurrency. * @param concurrency the number of consumers to create. @@ -75,6 +80,24 @@ public class SimpleKafkaListenerContainerFactory this.consumerRebalanceListener = consumerRebalanceListener; } + /** + * Specify the commit callback. + * @param commitCallback the callback to set. + * @see ConcurrentMessageListenerContainer#setCommitCallback(OffsetCommitCallback) + */ + public void setCommitCallback(OffsetCommitCallback commitCallback) { + this.commitCallback = commitCallback; + } + + /** + * Specifiy whether or not to use sync commits. + * @param syncCommits the sync commits to set. + * @see ConcurrentMessageListenerContainer#setSyncCommits(boolean) + */ + public void setSyncCommits(Boolean syncCommits) { + this.syncCommits = syncCommits; + } + @Override protected ConcurrentMessageListenerContainer createContainerInstance(KafkaListenerEndpoint endpoint) { Collection topicPartitions = endpoint.getTopicPartitions(); @@ -107,6 +130,12 @@ public class SimpleKafkaListenerContainerFactory if (this.consumerRebalanceListener != null) { instance.setConsumerRebalanceListener(this.consumerRebalanceListener); } + if (this.commitCallback != null) { + instance.setCommitCallback(this.commitCallback); + } + if (this.syncCommits != null) { + instance.setSyncCommits(this.syncCommits); + } } } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java index f8f00edc..be10da26 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java @@ -64,7 +64,7 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis private OffsetCommitCallback commitCallback; - private boolean syncCommits; + private boolean syncCommits = true; /** * Construct an instance with the supplied configuration properties and specific diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index 52bd25b6..f57f157a 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -17,6 +17,7 @@ package org.springframework.kafka.annotation; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; import java.util.Collection; import java.util.Map; @@ -26,7 +27,7 @@ import java.util.concurrent.TimeUnit; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; import org.apache.kafka.clients.consumer.ConsumerRecord; - +import org.apache.kafka.clients.consumer.OffsetCommitCallback; import org.junit.ClassRule; import org.junit.Test; import org.junit.runner.RunWith; @@ -139,6 +140,9 @@ public class EnableKafkaIntegrationTests { this.registry.start(); assertThat(listenerContainer.isRunning()).isTrue(); listenerContainer.stop(); + assertThat(KafkaTestUtils.getPropertyValue(listenerContainer, "syncCommits", Boolean.class)).isFalse(); + assertThat(KafkaTestUtils.getPropertyValue(listenerContainer, "commitCallback")).isNotNull(); + assertThat(KafkaTestUtils.getPropertyValue(listenerContainer, "consumerRebalanceListener")).isNotNull(); } @Test @@ -187,7 +191,7 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaListenerContainerFactory() { + kafkaListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; @@ -195,7 +199,7 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaJsonListenerContainerFactory() { + kafkaJsonListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setMessageConverter(new StringJsonMessageConverter()); @@ -204,7 +208,7 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaManualAckListenerContainerFactory() { + kafkaManualAckListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(manualConsumerFactory()); factory.setAckMode(AckMode.MANUAL_IMMEDIATE); @@ -213,16 +217,19 @@ public class EnableKafkaIntegrationTests { @Bean public KafkaListenerContainerFactory> - kafkaAutoStartFalseListenerContainerFactory() { + kafkaAutoStartFalseListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setAutoStartup(false); + factory.setSyncCommits(false); + factory.setCommitCallback(mock(OffsetCommitCallback.class)); + factory.setConsumerRebalanceListener(mock(ConsumerRebalanceListener.class)); return factory; } @Bean public KafkaListenerContainerFactory> - kafkaRebalanceListenerContainerFactory() { + kafkaRebalanceListenerContainerFactory() { SimpleKafkaListenerContainerFactory factory = new SimpleKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setConsumerRebalanceListener(consumerRebalanceListener());