From 7bc90c10a229e065d6e8b9cc0938c2ac3e359520 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 2 Jun 2021 16:12:57 -0400 Subject: [PATCH] GH-1084: Add txCommitRecovered Property Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1084 Update copyrights --- docs/src/main/asciidoc/overview.adoc | 5 ++++ .../properties/KafkaConsumerProperties.java | 15 +++++++++++- .../kafka/KafkaMessageChannelBinder.java | 4 +++- .../stream/binder/kafka/KafkaBinderTests.java | 23 +++++++++++++++---- 4 files changed, 41 insertions(+), 6 deletions(-) diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index 26e028bb5..9f26856bb 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -316,6 +316,11 @@ Usually needed if you want to synchronize another transaction with the Kafka tra To achieve exactly once consumption and production of records, the consumer and producer bindings must all be configured with the same transaction manager. + Default: none. +txCommitRecovered:: +When using a transactional binder, the offset of a recovered record (e.g. when retries are exhausted and the record is sent to a dead letter topic) will be committed via a new transaction, by default. +Setting this property to `false` suppresses committing the offset of recovered record. ++ +Default: true. [[reset-offsets]] ==== Resetting Offsets 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 9b1087227..ae9ef2826 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2019 the original author or authors. + * Copyright 2016-2021 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. @@ -205,6 +205,11 @@ public class KafkaConsumerProperties { */ private String transactionManager; + /** + * Set to false to NOT commit the offset of a successfully recovered recovered in the after rollback processor. + */ + private boolean txCommitRecovered = true; + /** * @return if each record needs to be acknowledged. * @@ -516,4 +521,12 @@ public class KafkaConsumerProperties { this.transactionManager = transactionManager; } + public boolean isTxCommitRecovered() { + return this.txCommitRecovered; + } + + public void setTxCommitRecovered(boolean txCommitRecovered) { + this.txCommitRecovered = txCommitRecovered; + } + } 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 a42ef4680..9c7b5b2aa 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 @@ -766,7 +766,9 @@ public class KafkaMessageChannelBinder extends throw e; } } - }, createBackOff(extendedConsumerProperties))); + }, createBackOff(extendedConsumerProperties), + new KafkaTemplate<>(transMan.getProducerFactory()), + extendedConsumerProperties.getExtension().isTxCommitRecovered())); } else { kafkaMessageDrivenChannelAdapter.setErrorChannel(errorInfrastructure.getErrorChannel()); 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 ad3ba1782..73bd1baa0 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2019 the original author or authors. + * Copyright 2016-2021 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. @@ -121,7 +121,6 @@ import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ContainerProperties; -import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaders; @@ -1049,10 +1048,16 @@ public class KafkaBinderTests extends Binding consumerBinding = binder.bindConsumer(consumerDest, "testGroup", moduleInputChannel, consumerProperties); - MessageListenerContainer container = TestUtils.getPropertyValue(consumerBinding, - "lifecycle.messageListenerContainer", MessageListenerContainer.class); + AbstractMessageListenerContainer container = TestUtils.getPropertyValue(consumerBinding, + "lifecycle.messageListenerContainer", AbstractMessageListenerContainer.class); assertThat(container.getContainerProperties().getTopicPartitionsToAssign().length) .isEqualTo(4); // 2 topics 2 partitions each + if (transactional) { + assertThat(TestUtils.getPropertyValue(container.getAfterRollbackProcessor(), "kafkaTemplate")).isNotNull(); + assertThat( + TestUtils.getPropertyValue(container.getAfterRollbackProcessor(), "commitRecovered", Boolean.class)) + .isTrue(); + } String dlqTopic = useDlqDestResolver ? "foo.dlq" : "error.dlqTest." + uniqueBindingId + ".0.testGroup"; try (AdminClient admin = AdminClient.create(Collections.singletonMap(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, @@ -1074,6 +1079,7 @@ public class KafkaBinderTests extends ExtendedConsumerProperties dlqConsumerProperties = createConsumerProperties(); dlqConsumerProperties.setMaxAttempts(1); dlqConsumerProperties.setHeaderMode(headerMode); + dlqConsumerProperties.getExtension().setTxCommitRecovered(false); ApplicationContext context = TestUtils.getPropertyValue(binder.getBinder(), "applicationContext", ApplicationContext.class); @@ -1098,6 +1104,15 @@ public class KafkaBinderTests extends dlqTopic, null, dlqChannel, dlqConsumerProperties); binderBindUnbindLatency(); + if (transactional) { + assertThat(TestUtils.getPropertyValue(dlqConsumerBinding, + "lifecycle.messageListenerContainer.afterRollbackProcessor.kafkaTemplate")).isNotNull(); + assertThat( + TestUtils.getPropertyValue(dlqConsumerBinding, + "lifecycle.messageListenerContainer.afterRollbackProcessor.commitRecovered", Boolean.class)) + .isFalse(); + } + String testMessagePayload = "test." + UUID.randomUUID().toString(); Message testMessage = MessageBuilder .withPayload(testMessagePayload.getBytes())