From e2f91417c368c9973b48727c80516f0ca775e5e5 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Wed, 25 May 2016 19:42:15 -0400 Subject: [PATCH] Add setting for `autoCommitOnError` Fixes #542 - add option for `autoCommitOnError` - make `autoCommitOnError` follow the state of `enableDlq` - no need to suppress commits if messages will be sent elsewhere --- .../binder/kafka/KafkaConsumerProperties.java | 10 ++ .../kafka/KafkaMessageChannelBinder.java | 4 + .../stream/binder/kafka/KafkaBinderTests.java | 141 +++++++++++++++++- .../spring-cloud-stream-overview.adoc | 7 + 4 files changed, 157 insertions(+), 5 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java index 9ef6a5e0c..0f3ff40c2 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java @@ -23,6 +23,8 @@ public class KafkaConsumerProperties { private boolean autoCommitOffset = true; + private Boolean autoCommitOnError; + private boolean resetOffsets; private KafkaMessageChannelBinder.StartOffset startOffset; @@ -60,4 +62,12 @@ public class KafkaConsumerProperties { public void setEnableDlq(boolean enableDlq) { this.enableDlq = enableDlq; } + + public Boolean getAutoCommitOnError() { + return autoCommitOnError; + } + + public void setAutoCommitOnError(Boolean autoCommitOnError) { + this.autoCommitOnError = autoCommitOnError; + } } 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 c76aacdca..bda8dd3eb 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 @@ -480,6 +480,10 @@ public class KafkaMessageChannelBinder extends messageListenerContainer.setOffsetManager(offsetManager); messageListenerContainer.setQueueSize(configurationProperties.getQueueSize()); messageListenerContainer.setMaxFetch(configurationProperties.getFetchSize()); + boolean autoCommitOnError = properties.getExtension().getAutoCommitOnError() != null + ? properties.getExtension().getAutoCommitOnError() + : properties.getExtension().isAutoCommitOffset() && properties.getExtension().isEnableDlq(); + messageListenerContainer.setAutoCommitOnError(autoCommitOnError); int concurrency = Math.min(properties.getConcurrency(), listenedPartitions.size()); messageListenerContainer.setConcurrency(concurrency); 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 d7d84519d..5ffea6d6c 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 @@ -18,8 +18,11 @@ package org.springframework.cloud.stream.binder.kafka; import java.util.Arrays; import java.util.Collection; +import java.util.LinkedHashMap; import java.util.Properties; import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import kafka.admin.AdminUtils; import kafka.api.TopicMetadata; @@ -43,6 +46,7 @@ import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.kafka.core.Partition; import org.springframework.integration.kafka.core.TopicNotFoundException; +import org.springframework.integration.kafka.support.KafkaHeaders; import org.springframework.integration.kafka.support.ProducerConfiguration; import org.springframework.integration.kafka.support.ProducerMetadata; import org.springframework.messaging.Message; @@ -145,7 +149,6 @@ public class KafkaBinderTests extends @Test public void testDlqAndRetry() { KafkaTestBinder binder = getBinder(); - DirectChannel moduleOutputChannel = new DirectChannel(); DirectChannel moduleInputChannel = new DirectChannel(); QueueChannel dlqChannel = new QueueChannel(); @@ -183,6 +186,115 @@ public class KafkaBinderTests extends producerBinding.unbind(); } + @Test + public void testDefaultAutoCommitOnErrorWithoutDlq() throws Exception { + KafkaTestBinder binder = getBinder(); + DirectChannel moduleOutputChannel = new DirectChannel(); + DirectChannel moduleInputChannel = new DirectChannel(); + FailingInvocationCountingMessageHandler handler = new FailingInvocationCountingMessageHandler(); + moduleInputChannel.subscribe(handler); + ExtendedProducerProperties producerProperties = createProducerProperties(); + producerProperties.setPartitionCount(10); + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + consumerProperties.setMaxAttempts(1); + consumerProperties.setBackOffInitialInterval(100); + consumerProperties.setBackOffMaxInterval(150); + long uniqueBindingId = System.currentTimeMillis(); + Binding producerBinding = binder.bindProducer("retryTest." + uniqueBindingId + ".0", + moduleOutputChannel, producerProperties); + Binding consumerBinding = binder.bindConsumer("retryTest." + uniqueBindingId + ".0", + "testGroup", moduleInputChannel, consumerProperties); + + String testMessagePayload = "test." + UUID.randomUUID().toString(); + Message testMessage = MessageBuilder.withPayload(testMessagePayload).build(); + moduleOutputChannel.send(testMessage); + + assertTrue(handler.getLatch().await((int) (timeoutMultiplier * 1000), TimeUnit.MILLISECONDS)); + // first attempt fails + assertThat(handler.getReceivedMessages().entrySet(), hasSize(1)); + Message receivedMessage = handler.getReceivedMessages().entrySet().iterator().next().getValue(); + assertNotNull(receivedMessage); + assertEquals(testMessagePayload, receivedMessage.getPayload()); + assertThat(handler.getInvocationCount(), equalTo(consumerProperties.getMaxAttempts())); + consumerBinding.unbind(); + + // on the second attempt the message is redelivered + QueueChannel successfulInputChannel = new QueueChannel(); + consumerBinding = binder.bindConsumer("retryTest." + uniqueBindingId + ".0", "testGroup", + successfulInputChannel, consumerProperties); + String testMessage2Payload = "test." + UUID.randomUUID().toString(); + Message testMessage2 = MessageBuilder.withPayload(testMessage2Payload).build(); + moduleOutputChannel.send(testMessage2); + + Message firstReceived = receive(successfulInputChannel); + assertEquals(testMessagePayload, firstReceived.getPayload()); + Message secondReceived = receive(successfulInputChannel); + assertEquals(testMessage2Payload, secondReceived.getPayload()); + consumerBinding.unbind(); + producerBinding.unbind(); + } + + @Test + public void testDefaultAutoCommitOnErrorWithDlq() throws Exception { + KafkaTestBinder binder = getBinder(); + DirectChannel moduleOutputChannel = new DirectChannel(); + DirectChannel moduleInputChannel = new DirectChannel(); + FailingInvocationCountingMessageHandler handler = new FailingInvocationCountingMessageHandler(); + moduleInputChannel.subscribe(handler); + ExtendedProducerProperties producerProperties = createProducerProperties(); + producerProperties.setPartitionCount(10); + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + consumerProperties.setMaxAttempts(3); + consumerProperties.setBackOffInitialInterval(100); + consumerProperties.setBackOffMaxInterval(150); + consumerProperties.getExtension().setEnableDlq(true); + long uniqueBindingId = System.currentTimeMillis(); + Binding producerBinding = binder.bindProducer("retryTest." + uniqueBindingId + ".0", + moduleOutputChannel, producerProperties); + Binding consumerBinding = binder.bindConsumer("retryTest." + uniqueBindingId + ".0", + "testGroup", moduleInputChannel, consumerProperties); + + String testMessagePayload = "test." + UUID.randomUUID().toString(); + Message testMessage = MessageBuilder.withPayload(testMessagePayload).build(); + moduleOutputChannel.send(testMessage); + + ExtendedConsumerProperties dlqConsumerProperties = createConsumerProperties(); + dlqConsumerProperties.setMaxAttempts(1); + + QueueChannel dlqChannel = new QueueChannel(); + Binding dlqConsumerBinding = binder.bindConsumer( + "error.retryTest." + uniqueBindingId + ".0.testGroup", null, dlqChannel, dlqConsumerProperties); + + Message dlqMessage = receive(dlqChannel, 3); + assertNotNull(dlqMessage); + assertEquals(testMessagePayload, dlqMessage.getPayload()); + + // first attempt fails + assertThat(handler.getReceivedMessages().entrySet(), hasSize(1)); + Message handledMessage = handler.getReceivedMessages().entrySet().iterator().next().getValue(); + assertNotNull(handledMessage); + assertEquals(testMessagePayload, handledMessage.getPayload()); + assertThat(handler.getInvocationCount(), equalTo(consumerProperties.getMaxAttempts())); + + dlqConsumerBinding.unbind(); + consumerBinding.unbind(); + + + // on the second attempt the message is not redelivered because the DLQ is set + QueueChannel successfulInputChannel = new QueueChannel(); + consumerBinding = binder.bindConsumer("retryTest." + uniqueBindingId + ".0", "testGroup", + successfulInputChannel, consumerProperties); + String testMessage2Payload = "test." + UUID.randomUUID().toString(); + Message testMessage2 = MessageBuilder.withPayload(testMessage2Payload).build(); + moduleOutputChannel.send(testMessage2); + + Message receivedMessage = receive(successfulInputChannel); + assertEquals(testMessage2Payload, receivedMessage.getPayload()); + + consumerBinding.unbind(); + producerBinding.unbind(); + } + @Test(expected = IllegalArgumentException.class) public void testValidateKafkaTopicName() { KafkaMessageChannelBinder.validateTopicName("foo:bar"); @@ -454,9 +566,7 @@ public class KafkaBinderTests extends context.refresh(); binder.setApplicationContext(context); binder.afterPropertiesSet(); - DirectChannel output = new DirectChannel(); - String testTopicName = UUID.randomUUID().toString(); ExtendedProducerProperties properties = createProducerProperties(); properties.getExtension().setSync(true); @@ -507,7 +617,6 @@ public class KafkaBinderTests extends @Test public void testAutoConfigureTopicsDisabledSucceedsIfTopicExisting() throws Exception { - String testTopicName = "existing" + System.currentTimeMillis(); AdminUtils.createTopic(kafkaTestSupport.getZkClient(), testTopicName, 5, 1, new Properties()); KafkaBinderConfigurationProperties configurationProperties = createConfigurationProperties(); @@ -525,7 +634,6 @@ public class KafkaBinderTests extends @Test public void testAutoAddPartitionsDisabledFailsIfTopicUnderpartitioned() throws Exception { - String testTopicName = "existing" + System.currentTimeMillis(); AdminUtils.createTopic(kafkaTestSupport.getZkClient(), testTopicName, 1, 1, new Properties()); KafkaBinderConfigurationProperties configurationProperties = createConfigurationProperties(); @@ -664,17 +772,40 @@ public class KafkaBinderTests extends private int invocationCount; + private final LinkedHashMap> receivedMessages = new LinkedHashMap<>(); + + private final CountDownLatch latch; + + private FailingInvocationCountingMessageHandler(int latchSize) { + latch = new CountDownLatch(latchSize); + } + private FailingInvocationCountingMessageHandler() { + this(1); } @Override public void handleMessage(Message message) throws MessagingException { invocationCount++; + Long offset = message.getHeaders().get(KafkaHeaders.OFFSET, Long.class); + // using the offset as key allows to ensure that we don't store duplicate messages on retry + if (!receivedMessages.containsKey(offset)) { + receivedMessages.put(offset, message); + latch.countDown(); + } throw new RuntimeException(); } + public LinkedHashMap> getReceivedMessages() { + return receivedMessages; + } + public int getInvocationCount() { return invocationCount; } + + public CountDownLatch getLatch() { + return latch; + } } } diff --git a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc index 695a09497..7f1cafff7 100644 --- a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc +++ b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc @@ -1136,6 +1136,13 @@ autoCommitOffset:: If set to `false`, an `Acknowledgment` header will be available in the message headers for late acknowledgment. + Default: `true`. +autoCommitOnError:: + Effective only if `autoCommitOffset` is set to `true`. +If set to `false` it suppresses auto-commits for messages that result in errors, and will commit only for successful messages, allows a stream to automatically replay from the last successfully processed message, in case of persistent failures. +If set to `true`, it will always auto-commit (if auto-commit is enabled). +If not set (default), it effectively has the same value as `enableDlq`, auto-committing erroneous messages if they are sent to a DLQ, and not committing them otherwise. ++ +Default: not set. resetOffsets:: Whether to reset offsets on the consumer to the value provided by `startOffset`. +