From b4c7c36229f8f50f46656d091410fc25a4106eb3 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 1 Sep 2021 13:29:36 -0400 Subject: [PATCH] GH-1135: Disable container retries when no DLQ set (#1136) * GH-1135: Disable container retries when no DLQ set Disable default container retries when binding retries are enabled (maxAttempts > 1) and no DLQ set. Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1135 * Address PR review comments * Addressing PR review --- .../kafka/KafkaMessageChannelBinder.java | 6 ++- .../stream/binder/kafka/KafkaBinderTests.java | 45 +++++++++++++++++++ 2 files changed, 50 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 8a7ec27ff..3155f9543 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -107,6 +107,7 @@ import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; import org.springframework.kafka.listener.ConsumerProperties; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.DefaultAfterRollbackProcessor; +import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaders; @@ -733,14 +734,17 @@ public class KafkaMessageChannelBinder extends kafkaMessageDrivenChannelAdapter.setApplicationContext(applicationContext); ErrorInfrastructure errorInfrastructure = registerErrorInfrastructure(destination, consumerGroup, extendedConsumerProperties); + if (!extendedConsumerProperties.isBatchMode() && extendedConsumerProperties.getMaxAttempts() > 1 && transMan == null) { - kafkaMessageDrivenChannelAdapter .setRetryTemplate(buildRetryTemplate(extendedConsumerProperties)); kafkaMessageDrivenChannelAdapter .setRecoveryCallback(errorInfrastructure.getRecoverer()); + if (!extendedConsumerProperties.getExtension().isEnableDlq()) { + messageListenerContainer.setCommonErrorHandler(new DefaultErrorHandler(new FixedBackOff(0L, 0L))); + } } else if (!extendedConsumerProperties.isBatchMode() && transMan != null) { messageListenerContainer.setAfterRollbackProcessor(new DefaultAfterRollbackProcessor<>( 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 350e1471c..e6ec667ba 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 @@ -1298,6 +1298,51 @@ public class KafkaBinderTests extends producerBinding.unbind(); } + @Test + @SuppressWarnings("unchecked") + public void testRetriesWithoutDlq() throws Exception { + Binder binder = getBinder(); + ExtendedProducerProperties producerProperties = createProducerProperties(); + BindingProperties producerBindingProperties = createProducerBindingProperties( + producerProperties); + + DirectChannel moduleOutputChannel = createBindableChannel("output", + producerBindingProperties); + + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + consumerProperties.setMaxAttempts(2); + consumerProperties.setBackOffInitialInterval(100); + consumerProperties.setBackOffMaxInterval(150); + + DirectChannel moduleInputChannel = createBindableChannel("input", + createConsumerBindingProperties(consumerProperties)); + + FailingInvocationCountingMessageHandler handler = new FailingInvocationCountingMessageHandler(); + moduleInputChannel.subscribe(handler); + 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(); + Message testMessage = MessageBuilder + .withPayload(testMessagePayload.getBytes()).build(); + moduleOutputChannel.send(testMessage); + + Thread.sleep(3000); + + // Since we don't have a DLQ, assert that we are invoking the handler exactly the same number of times + // as set in consumerProperties.maxAttempt and not the default set by Spring Kafka (10 times). + assertThat(handler.getInvocationCount()) + .isEqualTo(consumerProperties.getMaxAttempts()); + binderBindUnbindLatency(); + consumerBinding.unbind(); + producerBinding.unbind(); + } + //See https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/870 for motivation for this test. @Test @SuppressWarnings("unchecked")