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 e1e87d677..e771c13cd 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,9 +153,9 @@ 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 boolean DEFAULT_RESET_OFFSETS = false; - private static final StartOffset DEFAULT_START = StartOffset.latest; + private static final StartOffset DEFAULT_START_OFFSET = StartOffset.latest; private RetryOperations retryOperations; @@ -278,9 +278,9 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { private Mode mode = Mode.embeddedHeaders; - private boolean resetOffsets = false; + private boolean resetOffsets = DEFAULT_RESET_OFFSETS; - private StartOffset startOffset = StartOffset.latest; + private StartOffset startOffset = DEFAULT_START_OFFSET; public KafkaMessageChannelBinder(ZookeeperConnect zookeeperConnect, String brokers, String zkAddress, String... headersToMap) { @@ -460,11 +460,11 @@ public class KafkaMessageChannelBinder extends MessageChannelBinderSupport { @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 resetOffsets at the latest time, which allows them to receive only messages sent after - // they've been bound + // If the caller provides a consumer group, use it; otherwise an anonymous consumer group + // is generated each time, such that each anonymous binding will receive all messages. + // Consumers reset offsets at the latest time by default, which allows them to receive only + // messages sent after they've been bound. That behavior can be changed with the + // "resetOffsets" and "startOffset" properties. String consumerGroup = group == null ? "anonymous." + UUID.randomUUID().toString() : group; long referencePoint = this.startOffset != null ? startOffset.getReferencePoint() : OffsetRequest.LatestTime(); 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 870591f77..86d976dc7 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 @@ -1,5 +1,5 @@ /* - * Copyright 2015 the original author or authors. + * Copyright 2015-2016 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -31,7 +31,7 @@ import org.springframework.context.annotation.PropertySource; */ @Configuration @ConditionalOnMissingBean(Binder.class) -@EnableConfigurationProperties({KafkaBinderConfigurationProperties.class,KafkaMessageChannelBinderConfiguration.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 7279c23c9..af96d4699 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 @@ -172,7 +172,6 @@ public class KafkaBinderTests extends PartitionCapableBinderTests { Arrays.fill(ratherBigPayload, (byte) 65); KafkaTestBinder binder = (KafkaTestBinder) getBinder(); - DirectChannel moduleOutputChannel = new DirectChannel(); QueueChannel moduleInputChannel = new QueueChannel(); Properties producerProperties = new Properties(); @@ -352,7 +351,6 @@ public class KafkaBinderTests extends PartitionCapableBinderTests { assertThat(new String(receivedMessage2.getPayload()), equalTo(testPayload2)); } - @Test @SuppressWarnings("unchecked") public void testEarliest() throws Exception { @@ -427,13 +425,9 @@ public class KafkaBinderTests extends PartitionCapableBinderTests { 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 { @@ -471,6 +465,6 @@ public class KafkaBinderTests extends PartitionCapableBinderTests { assertThat(receivedMessage3, not(nullValue())); assertThat(new String(receivedMessage3.getPayload()), equalTo(testPayload3)); binder.unbind(consumerBinding); - + binder.unbind(producerBinding); } } diff --git a/spring-cloud-stream-samples/source/pom.xml b/spring-cloud-stream-samples/source/pom.xml index bc6fb0149..3ce3d1df0 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-kafka + spring-cloud-stream-binder-redis org.springframework.boot