diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java index 6c304ff4d..1d0cbbcbd 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java @@ -23,7 +23,9 @@ import java.util.Map; * @author Marius Bogoevici * @author Ilayaperumal Gopinathan * - *

Thanks to Laszlo Szabo for providing the initial patch for generic property support.

+ *

+ * Thanks to Laszlo Szabo for providing the initial patch for generic property support. + *

*/ public class KafkaConsumerProperties { @@ -33,8 +35,6 @@ public class KafkaConsumerProperties { private Boolean autoCommitOnError; - private boolean resetOffsets; - private StartOffset startOffset; private boolean enableDlq; @@ -53,14 +53,6 @@ public class KafkaConsumerProperties { this.autoCommitOffset = autoCommitOffset; } - public boolean isResetOffsets() { - return this.resetOffsets; - } - - public void setResetOffsets(boolean resetOffsets) { - this.resetOffsets = resetOffsets; - } - public StartOffset getStartOffset() { return this.startOffset; } @@ -102,7 +94,8 @@ public class KafkaConsumerProperties { } public enum StartOffset { - earliest(-2L), latest(-1L); + earliest(-2L), + latest(-1L); private final long referencePoint; diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc index 878dd5d4d..5a760da4a 100644 --- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc @@ -146,12 +146,8 @@ recoveryInterval:: The interval between connection recovery attempts, in milliseconds. + Default: `5000`. -resetOffsets:: - Whether to reset offsets on the consumer to the value provided by `startOffset`. -+ -Default: `false`. startOffset:: - The starting offset for new groups, or when `resetOffsets` is `true`. + The starting offset for new groups. Allowed values: `earliest`, `latest`. If the consumer group is set explicitly for the consumer 'binding' (via `spring.cloud.stream.bindings..group`), then 'startOffset' is set to `earliest`; otherwise it is set to `latest` for the `anonymous` consumer group. + diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index a01cf4594..50aab787f 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -34,7 +34,6 @@ import org.apache.kafka.common.serialization.LongDeserializer; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.assertj.core.api.Condition; -import org.junit.Ignore; import org.junit.Rule; import org.junit.Test; import org.junit.rules.ExpectedException; @@ -86,16 +85,16 @@ import static org.junit.Assert.assertTrue; * @author Ilayaperumal Gopinathan * @author Henryk Konsek */ -public abstract class KafkaBinderTests extends PartitionCapableBinderTests, - ExtendedProducerProperties> { +public abstract class KafkaBinderTests extends + PartitionCapableBinderTests, ExtendedProducerProperties> { @Rule public ExpectedException expectedProvisioningException = ExpectedException.none(); @Override protected ExtendedConsumerProperties createConsumerProperties() { - final ExtendedConsumerProperties kafkaConsumerProperties = - new ExtendedConsumerProperties<>(new KafkaConsumerProperties()); + final ExtendedConsumerProperties kafkaConsumerProperties = new ExtendedConsumerProperties<>( + new KafkaConsumerProperties()); // set the default values that would normally be propagated by Spring Cloud Stream kafkaConsumerProperties.setInstanceCount(1); kafkaConsumerProperties.setInstanceIndex(0); @@ -104,7 +103,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests createProducerProperties() { - ExtendedProducerProperties producerProperties = new ExtendedProducerProperties<>(new KafkaProducerProperties()); + ExtendedProducerProperties producerProperties = new ExtendedProducerProperties<>( + new KafkaProducerProperties()); producerProperties.getExtension().setSync(true); return producerProperties; } @@ -120,10 +120,10 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests dlqConsumerProperties = createConsumerProperties(); dlqConsumerProperties.setMaxAttempts(1); QueueChannel dlqChannel = new QueueChannel(); - Binding dlqConsumerBinding = binder.bindConsumer(dlqName, null, dlqChannel, dlqConsumerProperties); + Binding dlqConsumerBinding = binder.bindConsumer(dlqName, null, dlqChannel, + dlqConsumerProperties); String testMessagePayload = "test." + UUID.randomUUID().toString(); Message testMessage = MessageBuilder.withPayload(testMessagePayload).build(); @@ -369,9 +370,9 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerBinding = binder.bindProducer(testTopicName, output, - createProducerProperties()); - String testPayload1 = "foo1-" + UUID.randomUUID().toString(); - output.send(new GenericMessage<>(testPayload1.getBytes())); - ExtendedConsumerProperties properties = createConsumerProperties(); - properties.getExtension().setResetOffsets(true); - properties.getExtension().setStartOffset(KafkaConsumerProperties.StartOffset.earliest); - Binding consumerBinding = binder.bindConsumer(testTopicName, "startOffsets", input1, - properties); - Message receivedMessage1 = (Message) receive(input1); - assertThat(receivedMessage1).isNotNull(); - assertThat(new String(receivedMessage1.getPayload())).isEqualTo(testPayload1); - String testPayload2 = "foo2-" + UUID.randomUUID().toString(); - output.send(new GenericMessage<>(testPayload2.getBytes())); - Message receivedMessage2 = (Message) receive(input1); - assertThat(receivedMessage2).isNotNull(); - assertThat(new String(receivedMessage2.getPayload())).isEqualTo(testPayload2); - consumerBinding.unbind(); - - String testPayload3 = "foo3-" + UUID.randomUUID().toString(); - output.send(new GenericMessage<>(testPayload3.getBytes())); - - ExtendedConsumerProperties properties2 = createConsumerProperties(); - properties2.getExtension().setResetOffsets(true); - properties2.getExtension().setStartOffset(KafkaConsumerProperties.StartOffset.earliest); - consumerBinding = binder.bindConsumer(testTopicName, "startOffsets", input1, properties2); - Message receivedMessage4 = (Message) receive(input1); - assertThat(receivedMessage4).isNotNull(); - assertThat(new String(receivedMessage4.getPayload())).isEqualTo(testPayload1); - Message receivedMessage5 = (Message) receive(input1); - assertThat(receivedMessage5).isNotNull(); - assertThat(new String(receivedMessage5.getPayload())).isEqualTo(testPayload2); - Message receivedMessage6 = (Message) receive(input1); - assertThat(receivedMessage6).isNotNull(); - assertThat(new String(receivedMessage6.getPayload())).isEqualTo(testPayload3); - consumerBinding.unbind(); - producerBinding.unbind(); - } - @Test @SuppressWarnings("unchecked") public void testResume() throws Exception { @@ -742,7 +694,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests[] messages = new Message[2]; messages[0] = receive(moduleInputChannel); messages[1] = receive(moduleInputChannel); @@ -768,13 +719,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerBinding = binder.bindProducer("testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", moduleOutputChannel, + Binding producerBinding = binder.bindProducer( + "testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", moduleOutputChannel, createProducerProperties()); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.getExtension().setAutoCommitOffset(false); - Binding consumerBinding = binder.bindConsumer("testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", "test", moduleInputChannel, + Binding consumerBinding = binder.bindConsumer( + "testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", "test", moduleInputChannel, consumerProperties); String testPayload1 = "foo" + UUID.randomUUID().toString(); @@ -788,7 +741,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests receivedMessage = receive(moduleInputChannel); assertThat(receivedMessage).isNotNull(); assertThat(receivedMessage.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT)).isNotNull(); - Acknowledgment acknowledgment = receivedMessage.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class); + Acknowledgment acknowledgment = receivedMessage.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, + Acknowledgment.class); try { acknowledgment.acknowledge(); } @@ -810,12 +764,14 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerBinding = binder.bindProducer("testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", moduleOutputChannel, + Binding producerBinding = binder.bindProducer( + "testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", moduleOutputChannel, createProducerProperties()); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); - Binding consumerBinding = binder.bindConsumer("testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", "test", moduleInputChannel, + Binding consumerBinding = binder.bindConsumer( + "testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", "test", moduleInputChannel, consumerProperties); String testPayload1 = "foo" + UUID.randomUUID().toString(); @@ -995,8 +951,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests outputBinding = binder.bindProducer("partJ.0", output, producerProperties); if (usesExplicitRouting()) { Object endpoint = extractEndpoint(outputBinding); - assertThat(getEndpointRouting(endpoint)). - contains(getExpectedRoutingBaseDestination("partJ.0", "test") + "-' + headers['partition']"); + assertThat(getEndpointRouting(endpoint)) + .contains(getExpectedRoutingBaseDestination("partJ.0", "test") + "-' + headers['partition']"); } output.send(new GenericMessage<>(2)); @@ -1044,7 +1000,7 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests binding2 = binder.bindConsumer("defaultGroup.0", null, input2, consumerProperties); - //Since we don't provide any topic info, let Kafka bind the consumer successfully + // Since we don't provide any topic info, let Kafka bind the consumer successfully Thread.sleep(1000); String testPayload1 = "foo-" + UUID.randomUUID().toString(); output.send(new GenericMessage<>(testPayload1.getBytes())); @@ -1063,7 +1019,7 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests(testPayload2.getBytes())); binding2 = binder.bindConsumer("defaultGroup.0", null, input2, consumerProperties); - //Since we don't provide any topic info, let Kafka bind the consumer successfully + // Since we don't provide any topic info, let Kafka bind the consumer successfully Thread.sleep(1000); String testPayload3 = "foo-" + UUID.randomUUID().toString(); output.send(new GenericMessage<>(testPayload3.getBytes())); @@ -1226,7 +1182,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties); DirectFieldAccessor consumerAccessor = new DirectFieldAccessor(getKafkaConsumer(binding)); - assertTrue("Expected StringDeserializer as a custom key deserializer", consumerAccessor.getPropertyValue("keyDeserializer") instanceof StringDeserializer); - assertTrue("Expected LongDeserializer as a custom value deserializer", consumerAccessor.getPropertyValue("valueDeserializer") instanceof LongDeserializer); + assertTrue("Expected StringDeserializer as a custom key deserializer", + consumerAccessor.getPropertyValue("keyDeserializer") instanceof StringDeserializer); + assertTrue("Expected LongDeserializer as a custom value deserializer", + consumerAccessor.getPropertyValue("valueDeserializer") instanceof LongDeserializer); } finally { if (binding != null) { @@ -1367,11 +1326,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerProperties = createProducerProperties(); producerProperties.setUseNativeEncoding(true); - producerProperties.getExtension().getConfiguration().put("value.serializer", "org.apache.kafka.common.serialization.IntegerSerializer"); + producerProperties.getExtension().getConfiguration().put("value.serializer", + "org.apache.kafka.common.serialization.IntegerSerializer"); producerBinding = binder.bindProducer(testTopicName, moduleOutputChannel, producerProperties); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.getExtension().setAutoRebalanceEnabled(false); - consumerProperties.getExtension().getConfiguration().put("value.deserializer", "org.apache.kafka.common.serialization.IntegerDeserializer"); + consumerProperties.getExtension().getConfiguration().put("value.deserializer", + "org.apache.kafka.common.serialization.IntegerDeserializer"); consumerBinding = binder.bindConsumer(testTopicName, "test", moduleInputChannel, consumerProperties); // Let the consumer actually bind to the producer before sending a msg binderBindUnbindLatency(); @@ -1497,9 +1462,9 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests input2Binding = binder.bindConsumer("partJ.raw.0", "test", input2, consumerProperties); - output.send(new GenericMessage<>(new byte[]{(byte) 0})); - output.send(new GenericMessage<>(new byte[]{(byte) 1})); - output.send(new GenericMessage<>(new byte[]{(byte) 2})); + output.send(new GenericMessage<>(new byte[] { (byte) 0 })); + output.send(new GenericMessage<>(new byte[] { (byte) 1 })); + output.send(new GenericMessage<>(new byte[] { (byte) 2 })); Message receive0 = receive(input0); assertThat(receive0).isNotNull(); @@ -1538,7 +1503,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); consumerProperties.setConcurrency(2); consumerProperties.setInstanceIndex(0); @@ -1558,13 +1522,13 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests input2Binding = binder.bindConsumer("part.raw.0", "test", input2, consumerProperties); - Message message2 = org.springframework.integration.support.MessageBuilder.withPayload(new byte[]{2}) + Message message2 = org.springframework.integration.support.MessageBuilder.withPayload(new byte[] { 2 }) .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, "kafkaBinderTestCommonsDelegate") .setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, 42) .setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE, 43).build(); output.send(message2); - output.send(new GenericMessage<>(new byte[]{1})); - output.send(new GenericMessage<>(new byte[]{0})); + output.send(new GenericMessage<>(new byte[] { 1 })); + output.send(new GenericMessage<>(new byte[] { 0 })); Message receive0 = receive(input0); assertThat(receive0).isNotNull(); Message receive1 = receive(input1); @@ -1593,7 +1557,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerBinding = binder.bindConsumer("raw.0", "test", moduleInputChannel, consumerProperties); - Message message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithRawMode".getBytes()).build(); + Message message = org.springframework.integration.support.MessageBuilder + .withPayload("testSendAndReceiveWithRawMode".getBytes()).build(); // Let the consumer actually bind to the producer before sending a msg binderBindUnbindLatency(); moduleOutputChannel.send(message); @@ -1618,13 +1583,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerBinding = binder.bindConsumer("raw.string.0", "test", moduleInputChannel, consumerProperties); - Message message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithRawModeAndStringPayload").build(); + Message message = org.springframework.integration.support.MessageBuilder + .withPayload("testSendAndReceiveWithRawModeAndStringPayload").build(); // Let the consumer actually bind to the producer before sending a msg binderBindUnbindLatency(); moduleOutputChannel.send(message); Message inbound = receive(moduleInputChannel); assertThat(inbound).isNotNull(); - assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("testSendAndReceiveWithRawModeAndStringPayload"); + assertThat(new String((byte[]) inbound.getPayload())) + .isEqualTo("testSendAndReceiveWithRawModeAndStringPayload"); producerBinding.unbind(); consumerBinding.unbind(); } @@ -1656,14 +1623,16 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests input3Binding = binder.bindConsumer(barTapName, "tap2", module3InputChannel, consumerProperties); - Message message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithExplicitConsumerGroupWithRawMode".getBytes()).build(); + Message message = org.springframework.integration.support.MessageBuilder + .withPayload("testSendAndReceiveWithExplicitConsumerGroupWithRawMode".getBytes()).build(); boolean success = false; boolean retried = false; while (!success) { moduleOutputChannel.send(message); Message inbound = receive(module1InputChannel); assertThat(inbound).isNotNull(); - assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode"); + assertThat(new String((byte[]) inbound.getPayload())) + .isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode"); Message tapped1 = receive(module2InputChannel); Message tapped2 = receive(module3InputChannel); @@ -1674,12 +1643,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests message2 = org.springframework.integration.support.MessageBuilder.withPayload("bar".getBytes()).build(); + Message message2 = org.springframework.integration.support.MessageBuilder.withPayload("bar".getBytes()) + .build(); moduleOutputChannel.send(message2); // other tap still receives messages