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 {