From 7b8f0dcab77965ae389e717cc010c4b4bcaedde1 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 13 Nov 2019 09:41:02 -0500 Subject: [PATCH] Kafka Streams - DLQ control per consumer binding (#801) * Kafka Streams - DLQ control per consumer binding Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/800 * Fine-grained DLQ control and deserialization exception handlers per input binding * Deprecate KafkaStreamsBinderConfigurationProperties.SerdeError in preference to the new enum `KafkaStreamsBinderConfigurationProperties.DeserializationExceptionHandler` based properties * Add tests, modifying docs * Addressing PR review comments --- docs/src/main/asciidoc/kafka-streams.adoc | 43 +++++++++++---- .../AbstractKafkaStreamsBinderProcessor.java | 24 +++++++++ .../DeserializationExceptionHandler.java | 43 +++++++++++++++ ...StreamsBinderSupportAutoConfiguration.java | 6 +-- .../streams/KafkaStreamsBinderUtils.java | 9 +++- ...aStreamsBinderConfigurationProperties.java | 43 +++++++++++---- .../KafkaStreamsConsumerProperties.java | 14 +++++ ...serializationErrorHandlerByKafkaTests.java | 41 ++++++++++++++- ...serializtionErrorHandlerByBinderTests.java | 52 +++++++++++++++++-- 9 files changed, 249 insertions(+), 26 deletions(-) create mode 100644 spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationExceptionHandler.java diff --git a/docs/src/main/asciidoc/kafka-streams.adoc b/docs/src/main/asciidoc/kafka-streams.adoc index 616de99da..cf6ef1ce0 100644 --- a/docs/src/main/asciidoc/kafka-streams.adoc +++ b/docs/src/main/asciidoc/kafka-streams.adoc @@ -801,20 +801,20 @@ For details on this support, please see https://cwiki.apache.org/confluence/disp Out of the box, Apache Kafka Streams provides two kinds of deserialization exception handlers - `LogAndContinueExceptionHandler` and `LogAndFailExceptionHandler`. As the name indicates, the former will log the error and continue processing the next records and the latter will log the error and fail. `LogAndFailExceptionHandler` is the default deserialization exception handler. -=== Handling Deserialization Exceptions in the Binder +==== Handling Deserialization Exceptions in the Binder Kafka Streams binder allows to specify the deserialization exception handlers above using the following property. [source] ---- -spring.cloud.stream.kafka.streams.binder.serdeError: logAndContinue +spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler: logAndContinue ---- or [source] ---- -spring.cloud.stream.kafka.streams.binder.serdeError: logAndFail +spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler: logAndFail ---- In addition to the above two deserialization exception handlers, the binder also provides a third one for sending the erroneous records (poison pills) to a DLQ (dead letter queue) topic. @@ -822,10 +822,10 @@ Here is how you enable this DLQ exception handler. [source] ---- -spring.cloud.stream.kafka.streams.binder.serdeError: sendToDlq +spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler: sendToDlq ---- -When the above property is set, all the deserialization error records are automatically sent to the DLQ topic. +When the above property is set, all the records in deserialization error are automatically sent to the DLQ topic. You can set the topic name where the DLQ messages are published as below. @@ -834,11 +834,35 @@ You can set the topic name where the DLQ messages are published as below. spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer.dlqName: custom-dlq (Change the binding name accordingly) ---- -If this is set, then the error records are sent to the topic `custom-dlq`. If this is not set, then it will create a DLQ -topic with the name `error..`. -For instance, if your binding's destination topic is `inputTopic` and the applicatioin ID is `process-applicationId`, then the default DLQ topic is `error.inputTopic.process-applicationId`. +If this is set, then the error records are sent to the topic `custom-dlq`. +If this is not set, then it will create a DLQ topic with the name `error..`. +For instance, if your binding's destination topic is `inputTopic` and the application ID is `process-applicationId`, then the default DLQ topic is `error.inputTopic.process-applicationId`. It is always recommended to explicitly create a DLQ topic for each input binding if it is your intention to enable DLQ. +==== DLQ per input consumer binding + +The property `spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler` is applicable for the entire application. +This implies that if there are multiple functions or `StreamListener` methods in the same application, this property is applied to all of them. +However, if you have multiple processors or multiple input bindings within a single processor, then you can use the finer-grained DLQ control that the binder provides per input consumer binding. + +If you have the following processor, + +``` +@Bean +public BiFunction, KTable, KStream> process() { +... +} +``` + +and you only want to enable DLQ on the first input binding and logAndSkip on the second binding, then you can do so on the consumer as below. + +`spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer.deserializationExceptionHandler: sendToDlq` +`spring.cloud.stream.kafka.streams.bindings.process-in-1.consumer.deserializationExceptionHandler: logAndSkip` + +Setting deserialization exception handlers this way has a higher precedence than setting at the binder level. + +==== DLQ partitioning + By default, records are published to the Dead-Letter topic using the same partition as the original record. This means the Dead-Letter topic must have at least as many partitions as the original record. @@ -859,9 +883,10 @@ public DlqPartitionFunction partitionFunction() { NOTE: If you set a consumer binding's `dlqPartitions` property to 1 (and the binder's `minPartitionCount` is equal to `1`), there is no need to supply a `DlqPartitionFunction`; the framework will always use partition 0. If you set a consumer binding's `dlqPartitions` property to a value greater than `1` (or the binder's `minPartitionCount` is greater than `1`), you **must** provide a `DlqPartitionFunction` bean, even if the partition count is the same as the original topic's. + A couple of things to keep in mind when using the exception handling feature in Kafka Streams binder. -* The property `spring.cloud.stream.kafka.streams.binder.serdeError` is applicable for the entire application. This implies +* The property `spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler` is applicable for the entire application. This implies that if there are multiple functions or `StreamListener` methods in the same application, this property is applied to all of them. * The exception handling for deserialization works consistently with native deserialization and framework provided message conversion. diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java index 1e18a7879..d4d736d17 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -27,6 +27,8 @@ import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.Topology; +import org.apache.kafka.streams.errors.LogAndContinueExceptionHandler; +import org.apache.kafka.streams.errors.LogAndFailExceptionHandler; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; @@ -55,6 +57,7 @@ import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.config.StreamsBuilderFactoryBeanCustomizer; import org.springframework.kafka.core.CleanupConfig; +import org.springframework.kafka.streams.RecoveringDeserializationExceptionHandler; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.CollectionUtils; @@ -212,6 +215,27 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application concurrency); } + // Override deserialization exception handlers per binding + final DeserializationExceptionHandler deserializationExceptionHandler = + extendedConsumerProperties.getDeserializationExceptionHandler(); + if (deserializationExceptionHandler == DeserializationExceptionHandler.logAndFail) { + streamConfigGlobalProperties.put( + StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, + LogAndFailExceptionHandler.class); + } + else if (deserializationExceptionHandler == DeserializationExceptionHandler.logAndContinue) { + streamConfigGlobalProperties.put( + StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, + LogAndContinueExceptionHandler.class); + } + else if (deserializationExceptionHandler == DeserializationExceptionHandler.sendToDlq) { + streamConfigGlobalProperties.put( + StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, + RecoveringDeserializationExceptionHandler.class); + streamConfigGlobalProperties.put(RecoveringDeserializationExceptionHandler.KSTREAM_DESERIALIZATION_RECOVERER, + applicationContext.getBean(SendToDlqAndContinue.class)); + } + KafkaStreamsConfiguration kafkaStreamsConfiguration = new KafkaStreamsConfiguration(streamConfigGlobalProperties); StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.cleanupConfig == null diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationExceptionHandler.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationExceptionHandler.java new file mode 100644 index 000000000..33192c642 --- /dev/null +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationExceptionHandler.java @@ -0,0 +1,43 @@ +/* + * Copyright 2019-2019 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.kafka.streams; + +/** + * Enumeration for various {@link org.apache.kafka.streams.errors.DeserializationExceptionHandler} types. + * + * @author Soby Chacko + * @since 3.0.0 + */ +public enum DeserializationExceptionHandler { + + /** + * Deserialization error handler with log and continue. + * See {@link org.apache.kafka.streams.errors.LogAndContinueExceptionHandler} + */ + logAndContinue, + /** + * Deserialization error handler with log and fail. + * See {@link org.apache.kafka.streams.errors.LogAndFailExceptionHandler} + */ + logAndFail, + /** + * Deserialization error handler with DLQ send. + * See {@link org.springframework.kafka.streams.RecoveringDeserializationExceptionHandler} + */ + sendToDlq + +} diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index fb2fb80ac..21cc3ba91 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -221,19 +221,19 @@ public class KafkaStreamsBinderSupportAutoConfiguration { Serdes.ByteArraySerde.class.getName()); if (configProperties - .getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.logAndContinue) { + .getDeserializationExceptionHandler() == DeserializationExceptionHandler.logAndContinue) { properties.put( StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, LogAndContinueExceptionHandler.class); } else if (configProperties - .getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.logAndFail) { + .getDeserializationExceptionHandler() == DeserializationExceptionHandler.logAndFail) { properties.put( StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, LogAndFailExceptionHandler.class); } else if (configProperties - .getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.sendToDlq) { + .getDeserializationExceptionHandler() == DeserializationExceptionHandler.sendToDlq) { properties.put( StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, RecoveringDeserializationExceptionHandler.class); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java index 84898ae01..c24bd8352 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderUtils.java @@ -67,8 +67,15 @@ final class KafkaStreamsBinderUtils { ExtendedConsumerProperties extendedConsumerProperties = (ExtendedConsumerProperties) properties; + if (binderConfigurationProperties - .getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.sendToDlq) { + .getDeserializationExceptionHandler() == DeserializationExceptionHandler.sendToDlq) { + extendedConsumerProperties.getExtension().setEnableDlq(true); + } + // check for deserialization handler at the consumer binding, as that takes precedence. + final DeserializationExceptionHandler deserializationExceptionHandler = + properties.getExtension().getDeserializationExceptionHandler(); + if (deserializationExceptionHandler == DeserializationExceptionHandler.sendToDlq) { extendedConsumerProperties.getExtension().setEnableDlq(true); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java index caa6dd7f2..9aeb4c7d0 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java @@ -21,6 +21,7 @@ import java.util.Map; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.streams.DeserializationExceptionHandler; /** * Kafka Streams binder configuration properties. @@ -37,7 +38,10 @@ public class KafkaStreamsBinderConfigurationProperties /** * Enumeration for various Serde errors. + * + * @deprecated in favor of {@link DeserializationExceptionHandler}. */ + @Deprecated public enum SerdeError { /** @@ -61,6 +65,16 @@ public class KafkaStreamsBinderConfigurationProperties private Map functions = new HashMap<>(); + private KafkaStreamsBinderConfigurationProperties.SerdeError serdeError; + + /** + * {@link org.apache.kafka.streams.errors.DeserializationExceptionHandler} to use when + * there is a deserialization exception. This handler will be applied against all input bindings + * unless overridden at the consumer binding. + */ + private DeserializationExceptionHandler deserializationExceptionHandler; + + public Map getFunctions() { return functions; } @@ -85,21 +99,32 @@ public class KafkaStreamsBinderConfigurationProperties this.applicationId = applicationId; } - /** - * {@link org.apache.kafka.streams.errors.DeserializationExceptionHandler} to use when - * there is a Serde error. - * {@link KafkaStreamsBinderConfigurationProperties.SerdeError} values are used to - * provide the exception handler on consumer binding. - */ - private KafkaStreamsBinderConfigurationProperties.SerdeError serdeError; - + @Deprecated public KafkaStreamsBinderConfigurationProperties.SerdeError getSerdeError() { return this.serdeError; } + @Deprecated public void setSerdeError( KafkaStreamsBinderConfigurationProperties.SerdeError serdeError) { - this.serdeError = serdeError; + this.serdeError = serdeError; + if (serdeError == SerdeError.logAndContinue) { + this.deserializationExceptionHandler = DeserializationExceptionHandler.logAndContinue; + } + else if (serdeError == SerdeError.logAndFail) { + this.deserializationExceptionHandler = DeserializationExceptionHandler.logAndFail; + } + else if (serdeError == SerdeError.sendToDlq) { + this.deserializationExceptionHandler = DeserializationExceptionHandler.sendToDlq; + } + } + + public DeserializationExceptionHandler getDeserializationExceptionHandler() { + return deserializationExceptionHandler; + } + + public void setDeserializationExceptionHandler(DeserializationExceptionHandler deserializationExceptionHandler) { + this.deserializationExceptionHandler = deserializationExceptionHandler; } public static class StateStoreRetry { diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java index debfd169c..d0ba77836 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsConsumerProperties.java @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.properties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; +import org.springframework.cloud.stream.binder.kafka.streams.DeserializationExceptionHandler; /** * Extended properties for Kafka Streams consumer. @@ -43,6 +44,11 @@ public class KafkaStreamsConsumerProperties extends KafkaConsumerProperties { */ private String materializedAs; + /** + * Per input binding deserialization handler. + */ + private DeserializationExceptionHandler deserializationExceptionHandler; + /** * {@link org.apache.kafka.streams.processor.TimestampExtractor} bean name to use for this consumer. */ @@ -87,4 +93,12 @@ public class KafkaStreamsConsumerProperties extends KafkaConsumerProperties { public void setTimestampExtractorBeanName(String timestampExtractorBeanName) { this.timestampExtractorBeanName = timestampExtractorBeanName; } + + public DeserializationExceptionHandler getDeserializationExceptionHandler() { + return deserializationExceptionHandler; + } + + public void setDeserializationExceptionHandler(DeserializationExceptionHandler deserializationExceptionHandler) { + this.deserializationExceptionHandler = deserializationExceptionHandler; + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java index bdbfe8247..703f0f31c 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializationErrorHandlerByKafkaTests.java @@ -111,7 +111,7 @@ public abstract class DeserializationErrorHandlerByKafkaTests { @SpringBootTest(properties = { "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=deser-kafka-dlq", "spring.cloud.stream.bindings.input.group=group", - "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", + "spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler=sendToDlq", "spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=" + "org.apache.kafka.common.serialization.Serdes$IntegerSerde" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) public static class DeserializationByKafkaAndDlqTests @@ -147,6 +147,45 @@ public abstract class DeserializationErrorHandlerByKafkaTests { } + @SpringBootTest(properties = { + "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=deser-kafka-dlq", + "spring.cloud.stream.bindings.input.group=group", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.deserializationExceptionHandler=sendToDlq", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.valueSerde=" + + "org.apache.kafka.common.serialization.Serdes$IntegerSerde" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) + public static class DeserializationByKafkaAndDlqPerBindingTests + extends DeserializationErrorHandlerByKafkaTests { + + @Test + public void test() { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( + senderProps); + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic("DeserializationErrorHandlerByKafkaTests-In"); + template.sendDefault(1, null, "foobar"); + + Map consumerProps = KafkaTestUtils.consumerProps("foobar", + "false", embeddedKafka); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( + consumerProps); + Consumer consumer1 = cf.createConsumer(); + embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, "error.DeserializationErrorHandlerByKafkaTests-In.group"); + + ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer1, + "error.DeserializationErrorHandlerByKafkaTests-In.group"); + assertThat(cr.value()).isEqualTo("foobar"); + assertThat(cr.partition()).isEqualTo(0); // custom partition function + + // Ensuring that the deserialization was indeed done by Kafka natively + verify(conversionDelegate, never()).deserializeOnInbound(any(Class.class), + any(KStream.class)); + verify(conversionDelegate, never()).serializeOnOutbound(any(KStream.class)); + } + + } + @SpringBootTest(properties = { "spring.cloud.stream.bindings.input.destination=word1,word2", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id=deser-kafka-dlq-multi-input", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java index cb81eeb55..9e55fefb3 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DeserializtionErrorHandlerByBinderTests.java @@ -22,9 +22,9 @@ import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.streams.KeyValue; +import org.apache.kafka.streams.kstream.Grouped; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.kstream.Serialized; import org.apache.kafka.streams.kstream.TimeWindows; import org.junit.AfterClass; import org.junit.BeforeClass; @@ -110,7 +110,7 @@ public abstract class DeserializtionErrorHandlerByBinderTests { + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", "spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", + "spring.cloud.stream.kafka.streams.binder.deserializationExceptionHandler=sendToDlq", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id" + "=deserializationByBinderAndDlqTests", "spring.cloud.stream.kafka.streams.bindings.input.consumer.dlqPartitions=1", @@ -145,7 +145,53 @@ public abstract class DeserializtionErrorHandlerByBinderTests { verify(conversionDelegate).deserializeOnInbound(any(Class.class), any(KStream.class)); } + } + @SpringBootTest(properties = { + "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=false", + "spring.cloud.stream.bindings.output.producer.useNativeEncoding=false", + "spring.cloud.stream.bindings.input.destination=foos", + "spring.cloud.stream.bindings.output.destination=counts-id", + "spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", + "spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", + "spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + + "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.deserializationExceptionHandler=sendToDlq", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id" + + "=deserializationByBinderAndDlqTests", + "spring.cloud.stream.kafka.streams.bindings.input.consumer.dlqPartitions=1", + "spring.cloud.stream.bindings.input.group=foobar-group" }, webEnvironment = SpringBootTest.WebEnvironment.NONE) + public static class DeserializationByBinderAndDlqSetOnConsumerBindingTests + extends DeserializtionErrorHandlerByBinderTests { + + @Test + public void test() { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( + senderProps); + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic("foos"); + template.sendDefault(1, 7, "hello"); + + Map consumerProps = KafkaTestUtils.consumerProps("foobar", + "false", embeddedKafka); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( + consumerProps); + Consumer consumer1 = cf.createConsumer(); + embeddedKafka.consumeFromAnEmbeddedTopic(consumer1, + "error.foos.foobar-group"); + + ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer1, + "error.foos.foobar-group"); + assertThat(cr.value()).isEqualTo("hello"); + assertThat(cr.partition()).isEqualTo(0); + + // Ensuring that the deserialization was indeed done by the binder + verify(conversionDelegate).deserializeOnInbound(any(Class.class), + any(KStream.class)); + } } @SpringBootTest(properties = { @@ -211,7 +257,7 @@ public abstract class DeserializtionErrorHandlerByBinderTests { public KStream process(KStream input) { return input.filter((key, product) -> product.getId() == 123) .map((key, value) -> new KeyValue<>(value, value)) - .groupByKey(Serialized.with(new JsonSerde<>(Product.class), + .groupByKey(Grouped.with(new JsonSerde<>(Product.class), new JsonSerde<>(Product.class))) .windowedBy(TimeWindows.of(5000)) .count(Materialized.as("id-count-store-x")).toStream()