From 77e4087871407f56bdf49b493807f60051ddbd9f Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 5 Aug 2019 13:58:48 -0400 Subject: [PATCH] Handle Deserialization errors when there is a key (#711) * Handle Deserialization errors when there is a key Handle DLQ sending on deserialization errors when there is a key in the record. Resolves #635 * Adding keySerde information in the functional bindings --- ...fkaStreamsBindingInformationCatalogue.java | 20 +++++++++++++++++++ .../KafkaStreamsFunctionProcessor.java | 2 ++ ...KafkaStreamsMessageConversionDelegate.java | 10 +++++++++- ...StreamListenerSetupMethodOrchestrator.java | 1 + ...serializtionErrorHandlerByBinderTests.java | 14 +++++-------- 5 files changed, 37 insertions(+), 10 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java index 2918a3095..0115c8454 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBindingInformationCatalogue.java @@ -16,11 +16,13 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; +import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.KStream; @@ -49,6 +51,8 @@ class KafkaStreamsBindingInformationCatalogue { private ResolvableType outboundKStreamResolvable; + private final Map, Serde> keySerdeInfo = new HashMap<>(); + /** * For a given bounded {@link KStream}, retrieve it's corresponding destination on the * broker. @@ -132,4 +136,20 @@ class KafkaStreamsBindingInformationCatalogue { ResolvableType getOutboundKStreamResolvable() { return outboundKStreamResolvable; } + + /** + * Adding a mapping for KStream target to its corresponding KeySerde. + * This is used for sending to DLQ when deserialization fails. See {@link KafkaStreamsMessageConversionDelegate} + * for details. + * + * @param kStreamTarget target KStream + * @param keySerde Serde used for the key + */ + void addKeySerde(KStream kStreamTarget, Serde keySerde) { + this.keySerdeInfo.put(kStreamTarget, keySerde); + } + + Serde getKeySerde(KStream kStreamTarget) { + return this.keySerdeInfo.get(kStreamTarget); + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index 82627f235..cae881178 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -304,6 +304,8 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro (KStreamBoundElementFactory.KStreamWrapper) targetBean; //wrap the proxy created during the initial target type binding with real object (KStream) kStreamWrapper.wrap((KStream) stream); + + this.kafkaStreamsBindingInformationCatalogue.addKeySerde((KStream) kStreamWrapper, keySerde); this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); if (KStream.class.isAssignableFrom(stringResolvableTypeMap.get(input).getRawClass())) { diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java index 82a7618a6..052ca3870 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java @@ -25,6 +25,8 @@ import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.header.Header; import org.apache.kafka.common.header.Headers; import org.apache.kafka.common.header.internals.RecordHeader; +import org.apache.kafka.common.serialization.Serde; +import org.apache.kafka.common.serialization.Serializer; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.processor.Processor; @@ -282,8 +284,14 @@ public class KafkaStreamsMessageConversionDelegate { String destination = this.context.topic(); if (o2 instanceof Message) { Message message = (Message) o2; + + // We need to convert the key to a byte[] before sending to DLQ. + Serde keySerde = kstreamBindingInformationCatalogue.getKeySerde(bindingTarget); + Serializer keySerializer = keySerde.serializer(); + byte[] keyBytes = keySerializer.serialize(null, o); + KafkaStreamsMessageConversionDelegate.this.sendToDlqAndContinue - .sendToDlq(destination, (byte[]) o, + .sendToDlq(destination, keyBytes, (byte[]) message.getPayload(), this.context.partition()); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index ac97119f3..ea8671fa9 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -273,6 +273,7 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator extends AbstractKafkaStr // wrap the proxy created during the initial target type binding // with real object (KStream) kStreamWrapper.wrap((KStream) stream); + this.kafkaStreamsBindingInformationCatalogue.addKeySerde((KStream) kStreamWrapper, keySerde); this.kafkaStreamsBindingInformationCatalogue .addStreamBuilderFactory(streamsBuilderFactoryBean); for (StreamListenerParameterAdapter streamListenerParameterAdapter : adapters) { 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 5a101ecc2..bfa4990be 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 @@ -85,7 +85,7 @@ public abstract class DeserializtionErrorHandlerByBinderTests { System.setProperty("server.port", "0"); System.setProperty("spring.jmx.enabled", "false"); - Map consumerProps = KafkaTestUtils.consumerProps("foob", "false", + Map consumerProps = KafkaTestUtils.consumerProps("kafka-streams-dlq-tests", "false", embeddedKafka); // consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, // Deserializer.class.getName()); @@ -108,11 +108,9 @@ public abstract class DeserializtionErrorHandlerByBinderTests { "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$StringSerde", + + "=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.output.producer.keySerde" -// + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id" + "=deserializationByBinderAndDlqTests", @@ -122,13 +120,13 @@ public abstract class DeserializtionErrorHandlerByBinderTests { @Test @SuppressWarnings("unchecked") - public void test() throws Exception { + public void test() { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps); KafkaTemplate template = new KafkaTemplate<>(pf, true); template.setDefaultTopic("foos"); - template.sendDefault("hello"); + template.sendDefault(7, "hello"); Map consumerProps = KafkaTestUtils.consumerProps("foobar", "false", embeddedKafka); @@ -160,8 +158,6 @@ public abstract class DeserializtionErrorHandlerByBinderTests { + "=org.apache.kafka.common.serialization.Serdes$StringSerde", "spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", -// "spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde" -// + "=org.apache.kafka.common.serialization.Serdes$IntegerSerde", "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", "spring.cloud.stream.kafka.streams.bindings.input.consumer.application-id" + "=deserializationByBinderAndDlqTestsWithMultipleInputs", @@ -171,7 +167,7 @@ public abstract class DeserializtionErrorHandlerByBinderTests { @Test @SuppressWarnings("unchecked") - public void test() throws Exception { + public void test() { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( senderProps);