From 13693e8e66faf11a30e5e4aa8d93cc79ede38791 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 1 May 2018 12:02:36 -0400 Subject: [PATCH] Kafka Streams DLQ related changes DLQ handling needs to be adjusted in kafka streams binder due to the multiplexing of input topics. This commit changes it accordingly in KStream and KTable binders. Add tests to verify. --- .../binder/kafka/streams/KStreamBinder.java | 47 +++++++++++-------- .../binder/kafka/streams/KTableBinder.java | 27 +++++++---- ...KafkaStreamsMessageConversionDelegate.java | 2 +- ...serializationErrorHandlerByKafkaTests.java | 44 ++++++++++++++++- ...serializtionErrorHandlerByBinderTests.java | 47 ++++++++++++++++++- 5 files changed, 135 insertions(+), 32 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java index 28a26c8f6..6d458af72 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java @@ -42,7 +42,7 @@ import org.springframework.util.StringUtils; /** * {@link org.springframework.cloud.stream.binder.Binder} implementation for {@link KStream}. * This implemenation extends from the {@link AbstractBinder} directly. - * + *

* Provides both producer and consumer bindings for the bound KStream. * * @author Marius Bogoevici @@ -67,10 +67,10 @@ class KStreamBinder extends private final KeyValueSerdeResolver keyValueSerdeResolver; KStreamBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, - KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, - KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, - KeyValueSerdeResolver keyValueSerdeResolver) { + KafkaTopicProvisioner kafkaTopicProvisioner, + KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, + KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, + KeyValueSerdeResolver keyValueSerdeResolver) { this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; this.kafkaStreamsMessageConversionDelegate = kafkaStreamsMessageConversionDelegate; @@ -92,24 +92,34 @@ class KStreamBinder extends if (!StringUtils.hasText(group)) { group = binderConfigurationProperties.getApplicationId(); } + String[] inputTopics = StringUtils.commaDelimitedListToStringArray(name); for (String inputTopic : inputTopics) { this.kafkaTopicProvisioner.provisionConsumerDestination(inputTopic, group, extendedConsumerProperties); } - StreamsConfig streamsConfig = this.KafkaStreamsBindingInformationCatalogue.getStreamsConfig(inputTarget); - if (extendedConsumerProperties.getExtension().isEnableDlq()) { - String dlqName = StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName()) ? - "error." + name + "." + group : extendedConsumerProperties.getExtension().getDlqName(); - KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = new KafkaStreamsDlqDispatch(dlqName, binderConfigurationProperties, - extendedConsumerProperties.getExtension()); - SendToDlqAndContinue sendToDlqAndContinue = this.getApplicationContext().getBean(SendToDlqAndContinue.class); - sendToDlqAndContinue.addKStreamDlqDispatch(name, kafkaStreamsDlqDispatch); - DeserializationExceptionHandler deserializationExceptionHandler = streamsConfig.defaultDeserializationExceptionHandler(); - if(deserializationExceptionHandler instanceof SendToDlqAndContinue) { - ((SendToDlqAndContinue)deserializationExceptionHandler).addKStreamDlqDispatch(name, kafkaStreamsDlqDispatch); + if (extendedConsumerProperties.getExtension().isEnableDlq()) { + StreamsConfig streamsConfig = this.KafkaStreamsBindingInformationCatalogue.getStreamsConfig(inputTarget); + + KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = !StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName()) ? + new KafkaStreamsDlqDispatch(extendedConsumerProperties.getExtension().getDlqName(), binderConfigurationProperties, + extendedConsumerProperties.getExtension()) : null; + for (String inputTopic : inputTopics) { + if (StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName())) { + String dlqName = "error." + inputTopic + "." + group; + kafkaStreamsDlqDispatch = new KafkaStreamsDlqDispatch(dlqName, binderConfigurationProperties, + extendedConsumerProperties.getExtension()); + } + SendToDlqAndContinue sendToDlqAndContinue = this.getApplicationContext().getBean(SendToDlqAndContinue.class); + sendToDlqAndContinue.addKStreamDlqDispatch(inputTopic, kafkaStreamsDlqDispatch); + + DeserializationExceptionHandler deserializationExceptionHandler = streamsConfig.defaultDeserializationExceptionHandler(); + if (deserializationExceptionHandler instanceof SendToDlqAndContinue) { + ((SendToDlqAndContinue) deserializationExceptionHandler).addKStreamDlqDispatch(inputTopic, kafkaStreamsDlqDispatch); + } } } + return new DefaultBinding<>(name, group, inputTarget, null); } @@ -128,13 +138,12 @@ class KStreamBinder extends @SuppressWarnings("unchecked") private void to(boolean isNativeEncoding, String name, KStream outboundBindTarget, - Serde keySerde, Serde valueSerde) { + Serde keySerde, Serde valueSerde) { if (!isNativeEncoding) { LOG.info("Native encoding is disabled for " + name + ". Outbound message conversion done by Spring Cloud Stream."); kafkaStreamsMessageConversionDelegate.serializeOnOutbound(outboundBindTarget) .to(name, Produced.with(keySerde, valueSerde)); - } - else { + } else { LOG.info("Native encoding is enabled for " + name + ". Outbound serialization done at the broker."); outboundBindTarget.to(name, Produced.with(keySerde, valueSerde)); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java index a2624bb6c..6c68358d9 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java @@ -80,17 +80,24 @@ class KTableBinder extends } if (extendedConsumerProperties.getExtension().isEnableDlq()) { - String dlqName = StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName()) ? - "error." + name + "." + group : extendedConsumerProperties.getExtension().getDlqName(); - KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = new KafkaStreamsDlqDispatch(dlqName, binderConfigurationProperties, - extendedConsumerProperties.getExtension()); - SendToDlqAndContinue sendToDlqAndContinue = this.getApplicationContext().getBean(SendToDlqAndContinue.class); - sendToDlqAndContinue.addKStreamDlqDispatch(name, kafkaStreamsDlqDispatch); - StreamsConfig streamsConfig = this.KafkaStreamsBindingInformationCatalogue.getStreamsConfig(inputTarget); - DeserializationExceptionHandler deserializationExceptionHandler = streamsConfig.defaultDeserializationExceptionHandler(); - if(deserializationExceptionHandler instanceof SendToDlqAndContinue) { - ((SendToDlqAndContinue)deserializationExceptionHandler).addKStreamDlqDispatch(name, kafkaStreamsDlqDispatch); + + KafkaStreamsDlqDispatch kafkaStreamsDlqDispatch = !StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName()) ? + new KafkaStreamsDlqDispatch(extendedConsumerProperties.getExtension().getDlqName(), binderConfigurationProperties, + extendedConsumerProperties.getExtension()) : null; + for (String inputTopic : inputTopics) { + if (StringUtils.isEmpty(extendedConsumerProperties.getExtension().getDlqName())) { + String dlqName = "error." + inputTopic + "." + group; + kafkaStreamsDlqDispatch = new KafkaStreamsDlqDispatch(dlqName, binderConfigurationProperties, + extendedConsumerProperties.getExtension()); + } + SendToDlqAndContinue sendToDlqAndContinue = this.getApplicationContext().getBean(SendToDlqAndContinue.class); + sendToDlqAndContinue.addKStreamDlqDispatch(inputTopic, kafkaStreamsDlqDispatch); + + DeserializationExceptionHandler deserializationExceptionHandler = streamsConfig.defaultDeserializationExceptionHandler(); + if (deserializationExceptionHandler instanceof SendToDlqAndContinue) { + ((SendToDlqAndContinue) deserializationExceptionHandler).addKStreamDlqDispatch(inputTopic, kafkaStreamsDlqDispatch); + } } } return new DefaultBinding<>(name, group, inputTarget, null); 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 863f0eb45..e41ed81b1 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 @@ -165,7 +165,7 @@ class KafkaStreamsMessageConversionDelegate { @Override public void process(Object o, Object o2) { if (kstreamBindingInformationCatalogue.isDlqEnabled(bindingTarget)) { - String destination = kstreamBindingInformationCatalogue.getDestination(bindingTarget); + String destination = context.topic(); if (o2 instanceof Message) { Message message = (Message) o2; sendToDlqAndContinue.sendToDlq(destination, (byte[]) o, (byte[]) message.getPayload(), context.partition()); diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationErrorHandlerByKafkaTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationErrorHandlerByKafkaTests.java index 7ab78075c..8c77d7121 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationErrorHandlerByKafkaTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializationErrorHandlerByKafkaTests.java @@ -68,7 +68,8 @@ import static org.mockito.Mockito.verify; public abstract class DeserializationErrorHandlerByKafkaTests { @ClassRule - public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, "counts", "error.words.group"); + public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, "counts", "error.words.group", + "error.word1.groupx", "error.word2.groupx"); @SpyBean KafkaStreamsMessageConversionDelegate KafkaStreamsMessageConversionDelegate; @@ -130,6 +131,47 @@ public abstract class DeserializationErrorHandlerByKafkaTests { } } + @SpringBootTest(properties = { + "spring.cloud.stream.bindings.input.consumer.useNativeDecoding=true", + "spring.cloud.stream.bindings.output.producer.useNativeEncoding=true", + "spring.cloud.stream.bindings.input.destination=word1,word2", + "spring.cloud.stream.bindings.input.group=groupx", + "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", + "spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=" + + "org.apache.kafka.common.serialization.Serdes$IntegerSerde"}, + webEnvironment= SpringBootTest.WebEnvironment.NONE + ) + public static class DeserializationByKafkaAndDlqTestsWithMultipleInputs extends DeserializationErrorHandlerByKafkaTests { + + @Test + @SuppressWarnings("unchecked") + public void test() throws Exception { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic("word1"); + template.sendDefault("foobar"); + + template.setDefaultTopic("word2"); + template.sendDefault("foobar"); + + Map consumerProps = KafkaTestUtils.consumerProps("foobarx", "false", embeddedKafka); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); + Consumer consumer1 = cf.createConsumer(); + embeddedKafka.consumeFromEmbeddedTopics(consumer1, "error.word1.groupx", "error.word2.groupx"); + + //TODO: Investigate why the ordering matters below: i.e. if we consume from error.word1.groupx first, an exception is thrown. + ConsumerRecord cr1 = KafkaTestUtils.getSingleRecord(consumer1, "error.word2.groupx"); + assertThat(cr1.value().equals("foobar")).isTrue(); + ConsumerRecord cr2 = KafkaTestUtils.getSingleRecord(consumer1, "error.word1.groupx"); + assertThat(cr2.value().equals("foobar")).isTrue(); + + //Ensuring that the deserialization was indeed done by Kafka natively + verify(KafkaStreamsMessageConversionDelegate, never()).deserializeOnInbound(any(Class.class), any(KStream.class)); + verify(KafkaStreamsMessageConversionDelegate, never()).serializeOnOutbound(any(KStream.class)); + } + } @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializtionErrorHandlerByBinderTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializtionErrorHandlerByBinderTests.java index b73fc8a59..9e8057afb 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializtionErrorHandlerByBinderTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/DeserializtionErrorHandlerByBinderTests.java @@ -62,7 +62,8 @@ import static org.mockito.Mockito.verify; public abstract class DeserializtionErrorHandlerByBinderTests { @ClassRule - public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, "counts-id", "error.foos.foobar-group"); + public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, "counts-id", "error.foos.foobar-group", + "error.foos1.fooz-group", "error.foos2.fooz-group"); @SpyBean KafkaStreamsMessageConversionDelegate KafkaStreamsMessageConversionDelegate; @@ -128,6 +129,50 @@ public abstract class DeserializtionErrorHandlerByBinderTests { } } + @SpringBootTest(properties = { + "spring.cloud.stream.bindings.input.destination=foos1,foos2", + "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", + "spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=org.apache.kafka.common.serialization.Serdes$StringSerde", + "spring.cloud.stream.bindings.output.producer.headerMode=raw", + "spring.cloud.stream.kafka.streams.bindings.output.producer.keySerde=org.apache.kafka.common.serialization.Serdes$IntegerSerde", + "spring.cloud.stream.bindings.input.consumer.headerMode=raw", + "spring.cloud.stream.kafka.streams.binder.serdeError=sendToDlq", + "spring.cloud.stream.bindings.input.group=fooz-group"}, + webEnvironment= SpringBootTest.WebEnvironment.NONE + ) + public static class DeserializationByBinderAndDlqTestsWithMultipleInputs extends DeserializtionErrorHandlerByBinderTests { + + @Test + @SuppressWarnings("unchecked") + public void test() throws Exception { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic("foos1"); + template.sendDefault("hello"); + + template.setDefaultTopic("foos2"); + template.sendDefault("hello"); + + Map consumerProps = KafkaTestUtils.consumerProps("foobar1", "false", embeddedKafka); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); + Consumer consumer1 = cf.createConsumer(); + embeddedKafka.consumeFromEmbeddedTopics(consumer1, "error.foos1.fooz-group", "error.foos2.fooz-group"); + + ConsumerRecord cr1 = KafkaTestUtils.getSingleRecord(consumer1, "error.foos1.fooz-group"); + assertThat(cr1.value().equals("hello")).isTrue(); + + ConsumerRecord cr2 = KafkaTestUtils.getSingleRecord(consumer1, "error.foos2.fooz-group"); + assertThat(cr2.value().equals("hello")).isTrue(); + + //Ensuring that the deserialization was indeed done by the binder + verify(KafkaStreamsMessageConversionDelegate).deserializeOnInbound(any(Class.class), any(KStream.class)); + } + } + @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration public static class ProductCountApplication {