diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index d4bc55797..e1e87d677 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -153,6 +153,10 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { private static final boolean DEFAULT_AUTO_COMMIT_ENABLED = true; + private static final boolean DEFAULT_RESET_ON_START = false; + + private static final StartOffset DEFAULT_START = StartOffset.latest; + private RetryOperations retryOperations; /** @@ -274,6 +278,10 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { private Mode mode = Mode.embeddedHeaders; + private boolean resetOffsets = false; + + private StartOffset startOffset = StartOffset.latest; + public KafkaMessageChannelBinder(ZookeeperConnect zookeeperConnect, String brokers, String zkAddress, String... headersToMap) { this.zookeeperConnect = zookeeperConnect; @@ -434,15 +442,33 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { this.mode = mode; } + public boolean isResetOffsets() { + return resetOffsets; + } + + public void setResetOffsets(boolean resetOffsets) { + this.resetOffsets = resetOffsets; + } + + public StartOffset getStartOffset() { + return startOffset; + } + + public void setStartOffset(StartOffset startOffset) { + this.startOffset = startOffset; + } + @Override protected Binding doBindConsumer(String name, String group, MessageChannel inputChannel, Properties properties) { // If the caller provides a group, use it; otherwise // usage of a different consumer group each time achieves pub-sub // but multiple instances of this binding will each get all messages - // PubSub consumers reset at the latest time, which allows them to receive only messages sent after + // PubSub consumers resetOffsets at the latest time, which allows them to receive only messages sent after // they've been bound String consumerGroup = group == null ? "anonymous." + UUID.randomUUID().toString() : group; - return createKafkaConsumer(name, inputChannel, properties, consumerGroup, OffsetRequest.LatestTime()); + long referencePoint = this.startOffset != null ? + startOffset.getReferencePoint() : OffsetRequest.LatestTime(); + return createKafkaConsumer(name, inputChannel, properties, consumerGroup, referencePoint); } @Override @@ -678,6 +704,9 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { // if we have less target partitions than target concurrency, adjust accordingly messageListenerContainer.setConcurrency(Math.min(maxConcurrency, listenedPartitions.size())); OffsetManager offsetManager = createOffsetManager(group, referencePoint); + if (resetOffsets) { + offsetManager.resetOffsets(listenedPartitions); + } messageListenerContainer.setOffsetManager(offsetManager); messageListenerContainer.setQueueSize(accessor.getProperty(QUEUE_SIZE, defaultQueueSize)); messageListenerContainer.setMaxFetch(accessor.getProperty(FETCH_SIZE, defaultFetchSize)); @@ -883,4 +912,19 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { embeddedHeaders } + public enum StartOffset { + earliest(OffsetRequest.EarliestTime()), + latest(OffsetRequest.LatestTime()); + + private final long referencePoint; + + StartOffset(long referencePoint) { + this.referencePoint = referencePoint; + } + + public long getReferencePoint() { + return referencePoint; + } + } + } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java index ef3800aa5..00fe867a8 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaMessageChannelBinderConfiguration.java @@ -19,7 +19,6 @@ package org.springframework.cloud.stream.binder.kafka.config; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.PropertyPlaceholderAutoConfiguration; import org.springframework.boot.context.properties.ConfigurationProperties; -import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder; import org.springframework.cloud.stream.config.codec.kryo.KryoCodecAutoConfiguration; import org.springframework.context.annotation.Bean; @@ -36,7 +35,6 @@ import org.springframework.util.StringUtils; * @author Mark Fisher */ @Configuration -@EnableConfigurationProperties(KafkaBinderConfigurationProperties.class) @Import({KryoCodecAutoConfiguration.class, PropertyPlaceholderAutoConfiguration.class}) @ConfigurationProperties(prefix = "spring.cloud.stream.binder.kafka") public class KafkaMessageChannelBinderConfiguration { @@ -73,6 +71,10 @@ public class KafkaMessageChannelBinderConfiguration { private int offsetUpdateShutdownTimeout; + private boolean resetOffsets = false; + + private KafkaMessageChannelBinder.StartOffset startOffset; + @Autowired private Codec codec; @@ -119,6 +121,9 @@ public class KafkaMessageChannelBinderConfiguration { .getReplicationFactor()); kafkaMessageChannelBinder.setDefaultRequiredAcks(kafkaBinderConfigurationProperties.getRequiredAcks()); + kafkaMessageChannelBinder.setResetOffsets(resetOffsets); + kafkaMessageChannelBinder.setStartOffset(startOffset); + return kafkaMessageChannelBinder; } @@ -202,6 +207,22 @@ public class KafkaMessageChannelBinderConfiguration { return toConnectionString(this.brokers, this.defaultBrokerPort); } + public KafkaMessageChannelBinder.StartOffset getStartOffset() { + return startOffset; + } + + public void setStartOffset(KafkaMessageChannelBinder.StartOffset startOffset) { + this.startOffset = startOffset; + } + + public boolean isResetOffsets() { + return resetOffsets; + } + + public void setResetOffsets(boolean resetOffsets) { + this.resetOffsets = resetOffsets; + } + /** * Converts an array of host values to a comma-separated String. * diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaServiceAutoConfiguration.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaServiceAutoConfiguration.java index 1ca61606c..870591f77 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaServiceAutoConfiguration.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaServiceAutoConfiguration.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.kafka.config; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.test.ImportAutoConfiguration; +import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.binder.Binder; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Profile; @@ -30,6 +31,7 @@ import org.springframework.context.annotation.PropertySource; */ @Configuration @ConditionalOnMissingBean(Binder.class) +@EnableConfigurationProperties({KafkaBinderConfigurationProperties.class,KafkaMessageChannelBinderConfiguration.class}) @ImportAutoConfiguration(KafkaMessageChannelBinderConfiguration.class) public class KafkaServiceAutoConfiguration { diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index f5648cb26..7279c23c9 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -16,6 +16,10 @@ package org.springframework.cloud.stream.binder.kafka; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.Matchers.not; +import static org.hamcrest.Matchers.nullValue; import static org.hamcrest.collection.IsCollectionWithSize.hasSize; import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertNotNull; @@ -29,6 +33,7 @@ import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.TimeUnit; +import kafka.api.OffsetRequest; import org.junit.ClassRule; import org.junit.Test; @@ -38,16 +43,17 @@ import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.PartitionCapableBinderTests; import org.springframework.cloud.stream.binder.Spy; import org.springframework.cloud.stream.test.junit.kafka.KafkaTestSupport; +import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.kafka.core.KafkaMessage; import org.springframework.integration.kafka.core.Partition; import org.springframework.integration.kafka.listener.KafkaMessageListenerContainer; import org.springframework.integration.kafka.listener.MessageListener; +import org.springframework.integration.kafka.support.ZookeeperConnect; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; - -import kafka.api.OffsetRequest; +import org.springframework.messaging.support.GenericMessage; /** @@ -319,4 +325,152 @@ public class KafkaBinderTests extends PartitionCapableBinderTests { binder.unbind(consumerBinding); } + @Test + @SuppressWarnings("unchecked") + public void testDefaultConsumerStartsAtLatest() throws Exception { + KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(new ZookeeperConnect(kafkaTestSupport.getZkConnectString()), + kafkaTestSupport.getBrokerAddress(), kafkaTestSupport.getZkConnectString()); + GenericApplicationContext context = new GenericApplicationContext(); + context.refresh(); + binder.setApplicationContext(context); + binder.afterPropertiesSet(); + DirectChannel output = new DirectChannel(); + Properties properties = new Properties(); + QueueChannel input1 = new QueueChannel(); + + String testTopicName = UUID.randomUUID().toString(); + binder.bindProducer(testTopicName,output,properties); + String testPayload1 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload1.getBytes())); + binder.bindConsumer(testTopicName, "startOffsets", input1, properties); + Message receivedMessage1 = (Message) input1.receive(1000); + assertThat(receivedMessage1, is(nullValue())); + String testPayload2 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload2.getBytes())); + Message receivedMessage2 = (Message) input1.receive(1000); + assertThat(receivedMessage2, not(nullValue())); + assertThat(new String(receivedMessage2.getPayload()), equalTo(testPayload2)); + } + + + @Test + @SuppressWarnings("unchecked") + public void testEarliest() throws Exception { + KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(new ZookeeperConnect(kafkaTestSupport.getZkConnectString()), + kafkaTestSupport.getBrokerAddress(), kafkaTestSupport.getZkConnectString()); + GenericApplicationContext context = new GenericApplicationContext(); + context.refresh(); + binder.setApplicationContext(context); + binder.afterPropertiesSet(); + binder.setStartOffset(KafkaMessageChannelBinder.StartOffset.earliest); + DirectChannel output = new DirectChannel(); + Properties properties = new Properties(); + QueueChannel input1 = new QueueChannel(); + + String testTopicName = UUID.randomUUID().toString(); + binder.bindProducer(testTopicName,output,properties); + String testPayload1 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload1.getBytes())); + binder.bindConsumer(testTopicName, "startOffsets", input1, properties); + Message receivedMessage1 = (Message) input1.receive(1000); + assertThat(receivedMessage1, not(nullValue())); + String testPayload2 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload2.getBytes())); + Message receivedMessage2 = (Message) input1.receive(1000); + assertThat(receivedMessage2, not(nullValue())); + assertThat(new String(receivedMessage2.getPayload()), equalTo(testPayload2)); + } + + @Test + @SuppressWarnings("unchecked") + public void testReset() throws Exception { + KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(new ZookeeperConnect(kafkaTestSupport.getZkConnectString()), + kafkaTestSupport.getBrokerAddress(), kafkaTestSupport.getZkConnectString()); + GenericApplicationContext context = new GenericApplicationContext(); + context.refresh(); + binder.setApplicationContext(context); + binder.setStartOffset(KafkaMessageChannelBinder.StartOffset.earliest); + binder.setResetOffsets(true); + binder.afterPropertiesSet(); + DirectChannel output = new DirectChannel(); + Properties properties = new Properties(); + QueueChannel input1 = new QueueChannel(); + + String testTopicName = UUID.randomUUID().toString(); + + Binding producerBinding = binder.bindProducer(testTopicName, output, properties); + String testPayload1 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload1.getBytes())); + Binding consumerBinding = + binder.bindConsumer(testTopicName, "startOffsets", input1, properties); + Message receivedMessage1 = (Message) input1.receive(1000); + assertThat(receivedMessage1, not(nullValue())); + String testPayload2 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload2.getBytes())); + Message receivedMessage2 = (Message) input1.receive(1000); + assertThat(receivedMessage2, not(nullValue())); + assertThat(new String(receivedMessage2.getPayload()), equalTo(testPayload2)); + binder.unbind(consumerBinding); + + String testPayload3 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload3.getBytes())); + + consumerBinding = + binder.bindConsumer(testTopicName, "startOffsets", input1, properties); + Message receivedMessage4 = (Message) input1.receive(1000); + assertThat(receivedMessage4, not(nullValue())); + assertThat(new String(receivedMessage4.getPayload()), equalTo(testPayload1)); + Message receivedMessage5 = (Message) input1.receive(1000); + assertThat(receivedMessage5, not(nullValue())); + assertThat(new String(receivedMessage5.getPayload()), equalTo(testPayload2)); + Message receivedMessage6 = (Message) input1.receive(1000); + assertThat(receivedMessage6, not(nullValue())); + assertThat(new String(receivedMessage6.getPayload()), equalTo(testPayload3)); + binder.unbind(consumerBinding); + + + binder.unbind(producerBinding); + } + + + + @Test + @SuppressWarnings("unchecked") + public void testResume() throws Exception { + KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(new ZookeeperConnect(kafkaTestSupport.getZkConnectString()), + kafkaTestSupport.getBrokerAddress(), kafkaTestSupport.getZkConnectString()); + GenericApplicationContext context = new GenericApplicationContext(); + context.refresh(); + binder.setApplicationContext(context); + binder.afterPropertiesSet(); + binder.setStartOffset(KafkaMessageChannelBinder.StartOffset.earliest); + DirectChannel output = new DirectChannel(); + Properties properties = new Properties(); + QueueChannel input1 = new QueueChannel(); + + String testTopicName = UUID.randomUUID().toString(); + Binding producerBinding = binder.bindProducer(testTopicName, output, properties); + String testPayload1 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload1.getBytes())); + Binding consumerBinding = binder.bindConsumer(testTopicName, "startOffsets", input1, properties); + Message receivedMessage1 = (Message) input1.receive(1000); + assertThat(receivedMessage1, not(nullValue())); + String testPayload2 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload2.getBytes())); + Message receivedMessage2 = (Message) input1.receive(1000); + assertThat(receivedMessage2, not(nullValue())); + assertThat(new String(receivedMessage2.getPayload()), equalTo(testPayload2)); + binder.unbind(consumerBinding); + + String testPayload3 = "foo-" + UUID.randomUUID().toString(); + output.send(new GenericMessage<>(testPayload3.getBytes())); + + consumerBinding = + binder.bindConsumer(testTopicName, "startOffsets", input1, properties); + Message receivedMessage3 = (Message) input1.receive(1000); + assertThat(receivedMessage3, not(nullValue())); + assertThat(new String(receivedMessage3.getPayload()), equalTo(testPayload3)); + binder.unbind(consumerBinding); + + } } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java index 7b7e9e452..732177917 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/PartitionCapableBinderTests.java @@ -51,6 +51,7 @@ import org.springframework.messaging.support.GenericMessage; * * @author Gary Russell * @author Mark Fisher + * @author Marius Bogoevici */ abstract public class PartitionCapableBinderTests extends BrokerBinderTests { diff --git a/spring-cloud-stream-samples/source/pom.xml b/spring-cloud-stream-samples/source/pom.xml index 3ce3d1df0..bc6fb0149 100644 --- a/spring-cloud-stream-samples/source/pom.xml +++ b/spring-cloud-stream-samples/source/pom.xml @@ -24,7 +24,7 @@ org.springframework.cloud - spring-cloud-stream-binder-redis + spring-cloud-stream-binder-kafka org.springframework.boot diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java index 39fb949af..41bd2ba4c 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java @@ -145,7 +145,7 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean if (useApplicationContextAsParent) { springApplicationBuilder.parent(context); } - else if (environment != null && binderConfiguration.isInheritEnvironment()) { + if (useApplicationContextAsParent || (environment != null && binderConfiguration.isInheritEnvironment())) { StandardEnvironment binderEnvironment = new StandardEnvironment(); binderEnvironment.merge(environment); springApplicationBuilder.environment(binderEnvironment);