From 387d435e31d99c469a023700e52fd766b3e673b2 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 4 Nov 2021 14:47:05 -0400 Subject: [PATCH] GH-1994: Option to Strip Exception Headers Resolves https://github.com/spring-projects/spring-kafka/issues/1994 Add an option to strip previous exception headers when republishing a dead letter record. * Change default to true for 2.8. * Fix typo in doc. --- .../src/main/asciidoc/kafka.adoc | 18 +++++ .../src/main/asciidoc/retrytopic.adoc | 27 +++++++ .../src/main/asciidoc/whats-new.adoc | 7 ++ .../DeadLetterPublishingRecoverer.java | 53 +++++++++--- .../DeadLetterPublishingRecovererFactory.java | 2 +- .../retrytopic/RetryTopicBootstrapper.java | 2 +- .../RetryTopicInternalBeanNames.java | 11 ++- .../DeadLetterPublishingRecovererTests.java | 80 ++++++++++++++----- .../RetryTopicBootstrapperTests.java | 4 +- .../RetryTopicInternalBeanNamesTests.java | 25 ++++-- 10 files changed, 187 insertions(+), 42 deletions(-) diff --git a/spring-kafka-docs/src/main/asciidoc/kafka.adoc b/spring-kafka-docs/src/main/asciidoc/kafka.adoc index 2851a555..7eae911e 100644 --- a/spring-kafka-docs/src/main/asciidoc/kafka.adoc +++ b/spring-kafka-docs/src/main/asciidoc/kafka.adoc @@ -5532,6 +5532,24 @@ Starting with version 2.7, the recoverer checks that the partition selected by t If the partition is not present, the partition in the `ProducerRecord` is set to `null`, allowing the `KafkaProducer` to select the partition. You can disable this check by setting the `verifyPartition` property to `false`. +[[dlpr-headers]] +===== Managing Dead Letter Record Headers + +Referring to <> above, the `DeadLetterPublishingRecoverer` has two properties used to manage headers when those headers already exist (such as when reprocessing a dead letter record that failed, including when using <>). + +* `appendOriginalHeaders` (default `true`) +* `stripPreviousExceptionHeaders` (default `true` since version 2.8) + +Apache Kafka supports multiple headers with the same name; to obtain the "latest" value, you can use `headers.lastHeader(headerName)`; to get an iterator over multiple headers, use `headers.headers(headerName).iterator()`. + +When repeatedly republishing a failed record, these headers can grow (and eventually cause publication to fail due to a `RecordTooLargeException`); this is especially true for the exception headers and particularly for the stack trace headers. + +The reason for the two properties is because, while you might want to retain only the last exception information, you might want to retain the history of which topic(s) the record passed through for each failure. + +`appendOriginalHeaders` is applied to all headers named `*ORIGINAL*` while `stripPreviousExceptionHeaders` is applied to all headers named `*EXCEPTION*`. + +Also see <>. + [[exp-backoff]] ===== `ExponentialBackOffWithMaxRetries` Implementation diff --git a/spring-kafka-docs/src/main/asciidoc/retrytopic.adoc b/spring-kafka-docs/src/main/asciidoc/retrytopic.adoc index 160fca80..a439756b 100644 --- a/spring-kafka-docs/src/main/asciidoc/retrytopic.adoc +++ b/spring-kafka-docs/src/main/asciidoc/retrytopic.adoc @@ -357,6 +357,33 @@ public RetryTopicConfiguration myOtherRetryTopic(KafkaTemplate NOTE: By default the topics are autocreated with one partition and a replication factor of one. +[[retry-headers]] +===== Failure Header Management + +When considering how to manage failure headers (original headers and exception headers), the framework delegates to the `DeadLetterPublishingRecover` to decide whether to append or replace the headers. + +By default, it explicitly sets `appendOriginalHeaders` to `false` and leaves `stripPreviousExceptionHeaders` to the default used by the `DeadLetterPublishingRecover`. + +This means that, currently, records published to multiple retry topics may grow to large size, especially when the stack trace is large. + +See <> for more information. + +To reconfigure the framework to use different settings for these properties, replace standard `DeadLetterPublishingRecovererFactory` bean by adding a `recovererCustomizer`: + +==== +[source, java] +---- +@Bean(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME) +DeadLetterPublishingRecovererFactory factory(DestinationTopicResolver resolver) { + DeadLetterPublishingRecovererFactory factory = new DeadLetterPublishingRecovererFactory(resolver); + factory.setDeadLetterPublishingRecovererCustomizer(dlpr -> { + dlpr.appendOriginalHeaders(true); + dlpr.setStripPreviousExceptionHeaders(false); + }); + return factory; +} +---- +==== ==== Topic Naming diff --git a/spring-kafka-docs/src/main/asciidoc/whats-new.adoc b/spring-kafka-docs/src/main/asciidoc/whats-new.adoc index a7a39d83..d0ee028a 100644 --- a/spring-kafka-docs/src/main/asciidoc/whats-new.adoc +++ b/spring-kafka-docs/src/main/asciidoc/whats-new.adoc @@ -66,3 +66,10 @@ The `interceptBeforeTx` container property is now `true` by default. The `DelegatingByTopicSerializer` and `DelegatingByTopicDeserializer` are now provided. See <> for more information. + +[[x28-dlpr]] +==== `DeadLetterPublishingRecover` Changes + +The property `stripPreviousExceptionHeaders` is now `true` by default. + +See <> for more information. diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/DeadLetterPublishingRecoverer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/DeadLetterPublishingRecoverer.java index a0b2732d..399a50c2 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/DeadLetterPublishingRecoverer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/DeadLetterPublishingRecoverer.java @@ -93,7 +93,7 @@ public class DeadLetterPublishingRecoverer extends ExceptionClassifier implement private Duration waitForSendResultTimeout = Duration.ofSeconds(THIRTY); - private boolean replaceOriginalHeaders = true; + private boolean appendOriginalHeaders = true; private boolean failIfSendResultIsError = true; @@ -101,6 +101,8 @@ public class DeadLetterPublishingRecoverer extends ExceptionClassifier implement private long timeoutBuffer = Duration.ofSeconds(FIVE).toMillis(); + private boolean stripPreviousExceptionHeaders = true; + /** * Create an instance with the provided template and a default destination resolving * function that returns a TopicPartition based on the original topic (appended with ".DLT") @@ -246,13 +248,27 @@ public class DeadLetterPublishingRecoverer extends ExceptionClassifier implement } /** - * Set to false if you don't want to replace the dead letter original headers if - * they are already present. + * Set to false if you don't want to append the current "original" headers (topic, + * partition etc.) if they are already present. When false, only the first "original" + * headers are retained. * @param replaceOriginalHeaders set to false not to replace. * @since 2.7 + * @deprecated in favor of {@link #setAppendOriginalHeaders(boolean)}. */ + @Deprecated public void setReplaceOriginalHeaders(boolean replaceOriginalHeaders) { - this.replaceOriginalHeaders = replaceOriginalHeaders; + this.appendOriginalHeaders = replaceOriginalHeaders; + } + + /** + * Set to false if you don't want to append the current "original" headers (topic, + * partition etc.) if they are already present. When false, only the first "original" + * headers are retained. + * @param appendOriginalHeaders set to false not to replace. + * @since 2.7.9 + */ + public void setAppendOriginalHeaders(boolean appendOriginalHeaders) { + this.appendOriginalHeaders = appendOriginalHeaders; } /** @@ -299,6 +315,18 @@ public class DeadLetterPublishingRecoverer extends ExceptionClassifier implement this.timeoutBuffer = buffer; } + /** + * Set to false to retain previous exception headers as well as headers for the + * current exception. Default is true, which means only the current headers are + * retained; setting it to false this can cause a growth in record size when a record + * is republished many times. + * @param stripPreviousExceptionHeaders false to retain all. + * @since 2.7.9 + */ + public void setStripPreviousExceptionHeaders(boolean stripPreviousExceptionHeaders) { + this.stripPreviousExceptionHeaders = stripPreviousExceptionHeaders; + } + @SuppressWarnings("unchecked") @Override public void accept(ConsumerRecord record, @Nullable Consumer consumer, Exception exception) { @@ -563,32 +591,39 @@ public class DeadLetterPublishingRecoverer extends ExceptionClassifier implement } private void maybeAddHeader(Headers kafkaHeaders, String header, byte[] value) { - if (this.replaceOriginalHeaders || kafkaHeaders.lastHeader(header) == null) { + if (this.appendOriginalHeaders || kafkaHeaders.lastHeader(header) == null) { kafkaHeaders.add(header, value); } } void addExceptionInfoHeaders(Headers kafkaHeaders, Exception exception, boolean isKey) { - kafkaHeaders.add(new RecordHeader(isKey ? this.headerNames.exceptionInfo.keyExceptionFqcn + appendOrReplace(kafkaHeaders, new RecordHeader(isKey ? this.headerNames.exceptionInfo.keyExceptionFqcn : this.headerNames.exceptionInfo.exceptionFqcn, exception.getClass().getName().getBytes(StandardCharsets.UTF_8))); if (!isKey && exception.getCause() != null) { - kafkaHeaders.add(new RecordHeader(this.headerNames.exceptionInfo.exceptionCauseFqcn, + appendOrReplace(kafkaHeaders, new RecordHeader(this.headerNames.exceptionInfo.exceptionCauseFqcn, exception.getCause().getClass().getName().getBytes(StandardCharsets.UTF_8))); } String message = exception.getMessage(); if (message != null) { - kafkaHeaders.add(new RecordHeader(isKey + appendOrReplace(kafkaHeaders, new RecordHeader(isKey ? this.headerNames.exceptionInfo.keyExceptionMessage : this.headerNames.exceptionInfo.exceptionMessage, exception.getMessage().getBytes(StandardCharsets.UTF_8))); } - kafkaHeaders.add(new RecordHeader(isKey + appendOrReplace(kafkaHeaders, new RecordHeader(isKey ? this.headerNames.exceptionInfo.keyExceptionStacktrace : this.headerNames.exceptionInfo.exceptionStacktrace, getStackTraceAsString(exception).getBytes(StandardCharsets.UTF_8))); } + private void appendOrReplace(Headers headers, RecordHeader header) { + if (this.stripPreviousExceptionHeaders) { + headers.remove(header.key()); + } + headers.add(header); + } + private String getStackTraceAsString(Throwable cause) { StringWriter stringWriter = new StringWriter(); PrintWriter printWriter = new PrintWriter(stringWriter, true); diff --git a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/DeadLetterPublishingRecovererFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/DeadLetterPublishingRecovererFactory.java index ad11e11f..31c038cf 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/DeadLetterPublishingRecovererFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/DeadLetterPublishingRecovererFactory.java @@ -137,7 +137,7 @@ public class DeadLetterPublishingRecovererFactory { recoverer.setHeadersFunction((consumerRecord, e) -> addHeaders(consumerRecord, e, getAttempts(consumerRecord))); recoverer.setFailIfSendResultIsError(true); - recoverer.setReplaceOriginalHeaders(false); + recoverer.setAppendOriginalHeaders(false); recoverer.setThrowIfNoDestinationReturned(false); this.recovererCustomizer.accept(recoverer); this.fatalExceptions.forEach(ex -> recoverer.addNotRetryableExceptions(ex)); diff --git a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapper.java b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapper.java index a9cce40c..0e47332f 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapper.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapper.java @@ -81,7 +81,7 @@ public class RetryTopicBootstrapper { DefaultDestinationTopicProcessor.class); registerIfNotContains(RetryTopicInternalBeanNames.LISTENER_CONTAINER_FACTORY_CONFIGURER_NAME, ListenerContainerFactoryConfigurer.class); - registerIfNotContains(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME, + registerIfNotContains(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME, DeadLetterPublishingRecovererFactory.class); registerIfNotContains(RetryTopicInternalBeanNames.RETRY_TOPIC_CONFIGURER, RetryTopicConfigurer.class); registerIfNotContains(RetryTopicInternalBeanNames.DESTINATION_TOPIC_CONTAINER_NAME, diff --git a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNames.java b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNames.java index 8200417c..ef56912a 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNames.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNames.java @@ -52,7 +52,16 @@ public abstract class RetryTopicInternalBeanNames { /** * {@link DeadLetterPublishingRecovererFactory} bean name. */ - public static final String DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME = "internalDeadLetterPublishingRecovererProvider"; + public static final String DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME = + "internalDeadLetterPublishingRecovererProvider"; + + /** + * {@link DeadLetterPublishingRecovererFactory} bean name. + * @deprecated in favor of {@link #DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME} + */ + @Deprecated + public static final String DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME = + DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME; /** * {@link DestinationTopicContainer} bean name. diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/DeadLetterPublishingRecovererTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/DeadLetterPublishingRecovererTests.java index be766a70..21650771 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/DeadLetterPublishingRecovererTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/DeadLetterPublishingRecovererTests.java @@ -38,6 +38,7 @@ import java.io.UncheckedIOException; import java.time.Duration; import java.util.Collections; import java.util.HashMap; +import java.util.Iterator; import java.util.LinkedHashMap; import java.util.Map; import java.util.Optional; @@ -331,15 +332,15 @@ public class DeadLetterPublishingRecovererTests { @SuppressWarnings({"unchecked", "rawtypes"}) @Test - void dontReplaceOriginalHeaders() { + void dontAppendOriginalHeaders() { KafkaOperations template = mock(KafkaOperations.class); ListenableFuture future = mock(ListenableFuture.class); given(template.send(any(ProducerRecord.class))).willReturn(future); ConsumerRecord record = new ConsumerRecord<>("foo", 0, 0L, 1234L, TimestampType.CREATE_TIME, 123, 123, "bar", null, new RecordHeaders(), Optional.empty()); DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template); - recoverer.setReplaceOriginalHeaders(false); - recoverer.accept(record, new RuntimeException()); + recoverer.setAppendOriginalHeaders(false); + recoverer.accept(record, new RuntimeException(new IllegalStateException())); ArgumentCaptor producerRecordCaptor = ArgumentCaptor.forClass(ProducerRecord.class); then(template).should(times(1)).send(producerRecordCaptor.capture()); Headers headers = producerRecordCaptor.getValue().headers(); @@ -348,32 +349,50 @@ public class DeadLetterPublishingRecovererTests { Header originalOffsetHeader = headers.lastHeader(KafkaHeaders.DLT_ORIGINAL_OFFSET); Header originalTimestampHeader = headers.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP); Header originalTimestampType = headers.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP_TYPE); + Header firstExceptionType = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_FQCN); + Header firstExceptionCauseType = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_CAUSE_FQCN); + Header firstExceptionMessage = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_MESSAGE); + Header firstExceptionStackTrace = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_STACKTRACE); ConsumerRecord anotherRecord = new ConsumerRecord<>("bar", 1, 12L, 4321L, TimestampType.LOG_APPEND_TIME, 321, 321, "bar", null, new RecordHeaders(), Optional.empty()); headers.forEach(header -> anotherRecord.headers().add(header)); - recoverer.accept(anotherRecord, new RuntimeException()); + recoverer.accept(anotherRecord, new RuntimeException(new IllegalStateException())); ArgumentCaptor anotherProducerRecordCaptor = ArgumentCaptor.forClass(ProducerRecord.class); - then(template).should(times(2)).send(producerRecordCaptor.capture()); - Headers anotherHeaders = producerRecordCaptor.getAllValues().get(1).headers(); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TOPIC)).isEqualTo(originalTopicHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_PARTITION)).isEqualTo(originalPartitionHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_OFFSET)).isEqualTo(originalOffsetHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP)).isEqualTo(originalTimestampHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP_TYPE)).isEqualTo(originalTimestampType); + then(template).should(times(2)).send(anotherProducerRecordCaptor.capture()); + Headers anotherHeaders = anotherProducerRecordCaptor.getAllValues().get(1).headers(); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TOPIC)).isSameAs(originalTopicHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_PARTITION)).isSameAs(originalPartitionHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_OFFSET)).isSameAs(originalOffsetHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP)).isSameAs(originalTimestampHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP_TYPE)) + .isSameAs(originalTimestampType); + Iterator
originalTopics = anotherHeaders.headers(KafkaHeaders.DLT_ORIGINAL_TOPIC).iterator(); + assertThat(originalTopics.next()).isSameAs(originalTopicHeader); + assertThat(originalTopics.hasNext()).isFalse(); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_FQCN)).isNotSameAs(firstExceptionType); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_CAUSE_FQCN)) + .isNotSameAs(firstExceptionCauseType); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_MESSAGE)).isNotSameAs(firstExceptionMessage); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_STACKTRACE)) + .isNotSameAs(firstExceptionStackTrace); + Iterator
exceptionHeaders = anotherHeaders.headers(KafkaHeaders.DLT_EXCEPTION_FQCN).iterator(); + assertThat(exceptionHeaders.next()).isNotSameAs(firstExceptionType); + assertThat(exceptionHeaders.hasNext()).isFalse(); } @SuppressWarnings({"unchecked", "rawtypes"}) @Test - void replaceOriginalHeaders() { + void appendOriginalHeaders() { KafkaOperations template = mock(KafkaOperations.class); ListenableFuture future = mock(ListenableFuture.class); given(template.send(any(ProducerRecord.class))).willReturn(future); ConsumerRecord record = new ConsumerRecord<>("foo", 0, 0L, 1234L, TimestampType.CREATE_TIME, 123, 123, "bar", null, new RecordHeaders(), Optional.empty()); DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template); - recoverer.setReplaceOriginalHeaders(true); - recoverer.accept(record, new RuntimeException()); + recoverer.setAppendOriginalHeaders(true); + recoverer.setStripPreviousExceptionHeaders(false); + recoverer.accept(record, new RuntimeException(new IllegalStateException())); ArgumentCaptor producerRecordCaptor = ArgumentCaptor.forClass(ProducerRecord.class); then(template).should(times(1)).send(producerRecordCaptor.capture()); Headers headers = producerRecordCaptor.getValue().headers(); @@ -382,19 +401,40 @@ public class DeadLetterPublishingRecovererTests { Header originalOffsetHeader = headers.lastHeader(KafkaHeaders.DLT_ORIGINAL_OFFSET); Header originalTimestampHeader = headers.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP); Header originalTimestampType = headers.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP_TYPE); + Header firstExceptionType = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_FQCN); + Header firstExceptionCauseType = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_CAUSE_FQCN); + Header firstExceptionMessage = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_MESSAGE); + Header firstExceptionStackTrace = headers.lastHeader(KafkaHeaders.DLT_EXCEPTION_STACKTRACE); ConsumerRecord anotherRecord = new ConsumerRecord<>("bar", 1, 12L, 4321L, TimestampType.LOG_APPEND_TIME, 321, 321, "bar", null, new RecordHeaders(), Optional.empty()); headers.forEach(header -> anotherRecord.headers().add(header)); - recoverer.accept(anotherRecord, new RuntimeException()); + recoverer.accept(anotherRecord, new RuntimeException(new IllegalStateException())); ArgumentCaptor anotherProducerRecordCaptor = ArgumentCaptor.forClass(ProducerRecord.class); then(template).should(times(2)).send(anotherProducerRecordCaptor.capture()); Headers anotherHeaders = anotherProducerRecordCaptor.getAllValues().get(1).headers(); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TOPIC)).isNotEqualTo(originalTopicHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_PARTITION)).isNotEqualTo(originalPartitionHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_OFFSET)).isNotEqualTo(originalOffsetHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP)).isNotEqualTo(originalTimestampHeader); - assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP_TYPE)).isNotEqualTo(originalTimestampType); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TOPIC)).isNotSameAs(originalTopicHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_PARTITION)) + .isNotSameAs(originalPartitionHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_OFFSET)).isNotSameAs(originalOffsetHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP)) + .isNotSameAs(originalTimestampHeader); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_ORIGINAL_TIMESTAMP_TYPE)) + .isNotSameAs(originalTimestampType); + Iterator
originalTopics = anotherHeaders.headers(KafkaHeaders.DLT_ORIGINAL_TOPIC).iterator(); + assertThat(originalTopics.next()).isSameAs(originalTopicHeader); + assertThat(originalTopics.next()).isNotNull(); + assertThat(originalTopics.hasNext()).isFalse(); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_FQCN)).isNotSameAs(firstExceptionType); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_CAUSE_FQCN)) + .isNotSameAs(firstExceptionCauseType); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_MESSAGE)).isNotSameAs(firstExceptionMessage); + assertThat(anotherHeaders.lastHeader(KafkaHeaders.DLT_EXCEPTION_STACKTRACE)) + .isNotSameAs(firstExceptionStackTrace); + Iterator
exceptionHeaders = anotherHeaders.headers(KafkaHeaders.DLT_EXCEPTION_FQCN).iterator(); + assertThat(exceptionHeaders.next()).isSameAs(firstExceptionType); + assertThat(exceptionHeaders.next()).isNotNull(); + assertThat(exceptionHeaders.hasNext()).isFalse(); } @SuppressWarnings({ "unchecked", "rawtypes" }) diff --git a/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapperTests.java b/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapperTests.java index 877ee662..1ab6bd27 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapperTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicBootstrapperTests.java @@ -115,7 +115,7 @@ class RetryTopicBootstrapperTests { .registerBeanDefinition(RetryTopicInternalBeanNames.LISTENER_CONTAINER_FACTORY_CONFIGURER_NAME, new RootBeanDefinition(ListenerContainerFactoryConfigurer.class)); then(this.applicationContext).should(times(1)) - .registerBeanDefinition(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME, + .registerBeanDefinition(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME, new RootBeanDefinition(DeadLetterPublishingRecovererFactory.class)); then(this.applicationContext).should(times(1)) .registerBeanDefinition(RetryTopicInternalBeanNames.RETRY_TOPIC_CONFIGURER, @@ -158,7 +158,7 @@ class RetryTopicBootstrapperTests { .registerBeanDefinition(RetryTopicInternalBeanNames.LISTENER_CONTAINER_FACTORY_RESOLVER_NAME, new RootBeanDefinition(ListenerContainerFactoryResolver.class)); then(this.applicationContext).should(times(0)) - .registerBeanDefinition(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME, + .registerBeanDefinition(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME, new RootBeanDefinition(DeadLetterPublishingRecovererFactory.class)); then(this.applicationContext).should(times(0)) .registerBeanDefinition(RetryTopicInternalBeanNames.RETRY_TOPIC_CONFIGURER, diff --git a/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNamesTests.java b/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNamesTests.java index 4775c37b..be854f32 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNamesTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/retrytopic/RetryTopicInternalBeanNamesTests.java @@ -22,6 +22,7 @@ import org.junit.jupiter.api.Test; /** * @author Tomaz Fernandes + * @author Gary Russell * @since 2.7 */ class RetryTopicInternalBeanNamesTests { @@ -47,14 +48,22 @@ class RetryTopicInternalBeanNamesTests { @Test public void assertRetryTopicInternalBeanNamesConstants() { new RetryTopicInternalBeanNames() { }; // for coverage - assertThat(RetryTopicInternalBeanNames.DESTINATION_TOPIC_PROCESSOR_NAME).isEqualTo(DESTINATION_TOPIC_PROCESSOR_NAME); - assertThat(RetryTopicInternalBeanNames.KAFKA_CONSUMER_BACKOFF_MANAGER).isEqualTo(KAFKA_CONSUMER_BACKOFF_MANAGER); + assertThat(RetryTopicInternalBeanNames.DESTINATION_TOPIC_PROCESSOR_NAME) + .isEqualTo(DESTINATION_TOPIC_PROCESSOR_NAME); + assertThat(RetryTopicInternalBeanNames.KAFKA_CONSUMER_BACKOFF_MANAGER) + .isEqualTo(KAFKA_CONSUMER_BACKOFF_MANAGER); assertThat(RetryTopicInternalBeanNames.RETRY_TOPIC_CONFIGURER).isEqualTo(RETRY_TOPIC_CONFIGURER); - assertThat(RetryTopicInternalBeanNames.LISTENER_CONTAINER_FACTORY_RESOLVER_NAME).isEqualTo(LISTENER_CONTAINER_FACTORY_RESOLVER_NAME); - assertThat(RetryTopicInternalBeanNames.LISTENER_CONTAINER_FACTORY_CONFIGURER_NAME).isEqualTo(LISTENER_CONTAINER_FACTORY_CONFIGURER_NAME); - assertThat(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME).isEqualTo(DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME); - assertThat(RetryTopicInternalBeanNames.DESTINATION_TOPIC_CONTAINER_NAME).isEqualTo(DESTINATION_TOPIC_CONTAINER_NAME); - assertThat(RetryTopicInternalBeanNames.DEFAULT_LISTENER_FACTORY_BEAN_NAME).isEqualTo(DEFAULT_LISTENER_FACTORY_BEAN_NAME); - assertThat(RetryTopicInternalBeanNames.DEFAULT_KAFKA_TEMPLATE_BEAN_NAME).isEqualTo(DEFAULT_KAFKA_TEMPLATE_BEAN_NAME); + assertThat(RetryTopicInternalBeanNames.LISTENER_CONTAINER_FACTORY_RESOLVER_NAME) + .isEqualTo(LISTENER_CONTAINER_FACTORY_RESOLVER_NAME); + assertThat(RetryTopicInternalBeanNames.LISTENER_CONTAINER_FACTORY_CONFIGURER_NAME) + .isEqualTo(LISTENER_CONTAINER_FACTORY_CONFIGURER_NAME); + assertThat(RetryTopicInternalBeanNames.DEAD_LETTER_PUBLISHING_RECOVERER_FACTORY_BEAN_NAME) + .isEqualTo(DEAD_LETTER_PUBLISHING_RECOVERER_PROVIDER_NAME); + assertThat(RetryTopicInternalBeanNames.DESTINATION_TOPIC_CONTAINER_NAME) + .isEqualTo(DESTINATION_TOPIC_CONTAINER_NAME); + assertThat(RetryTopicInternalBeanNames.DEFAULT_LISTENER_FACTORY_BEAN_NAME) + .isEqualTo(DEFAULT_LISTENER_FACTORY_BEAN_NAME); + assertThat(RetryTopicInternalBeanNames.DEFAULT_KAFKA_TEMPLATE_BEAN_NAME) + .isEqualTo(DEFAULT_KAFKA_TEMPLATE_BEAN_NAME); } }