From d822129d6715a637010de1107154fbbd25daff6f Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 6 Sep 2023 20:36:35 -0400 Subject: [PATCH] DltAwareProcessor enhancements Cleaning up the code for the custom DltAwareProcessor --- .../pom.xml | 5 + .../kafka/streams/DltAwareProcessor.java | 58 ++++++++- ...Context.java => DltPublishingContext.java} | 7 +- ...StreamsBinderSupportAutoConfiguration.java | 4 +- .../integration/DltAwareProcessorTests.java | 120 ++++++++++++++++++ 5 files changed, 183 insertions(+), 11 deletions(-) rename binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/{DltSenderContext.java => DltPublishingContext.java} (82%) create mode 100644 binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DltAwareProcessorTests.java diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/pom.xml b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/pom.xml index c162ed379..c6b27579a 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/pom.xml +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/pom.xml @@ -76,6 +76,11 @@ test + + org.springframework.cloud + spring-cloud-stream-binder-kafka + test + diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java index 37a7c2977..5cdbbbd1a 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltAwareProcessor.java @@ -30,44 +30,88 @@ import org.springframework.util.Assert; import org.springframework.util.StringUtils; /** + * Custom {@link Processor} implementation that is capable of sending a record + * to a DLT if the processing fails. * * @author Soby Chacko * @since 4.1.0 */ public class DltAwareProcessor implements Processor { + /** + * Delegate {@link BiFunction} that is responsible for processing the data. + */ private final BiFunction> delegateFunction; + /** + * Event time for the forwarded downstream record. + */ private final Supplier recordTimeSupplier; + /** + * DLT destination. + */ private String dltDestination; - private DltSenderContext dltSenderContext; + /** + * {@link DltPublishingContext} used for DLT publishing needs. + */ + private DltPublishingContext dltPublishingContext; + /** + * A {@link BiConsumer} that does the recovery of a failed record. + */ private BiConsumer, Exception> processorRecordRecoverer; + /** + * {@link ProcessorContext} used in the processor. + */ private ProcessorContext context; + /** + * + * @param delegateFunction {@link BiFunction} to process the data + * @param dltDestination DLT destination + * @param dltPublishingContext {@link DltPublishingContext} + */ public DltAwareProcessor(BiFunction> delegateFunction, String dltDestination, - DltSenderContext dltSenderContext) { - this(delegateFunction, dltDestination, dltSenderContext, System::currentTimeMillis); + DltPublishingContext dltPublishingContext) { + this(delegateFunction, dltDestination, dltPublishingContext, System::currentTimeMillis); } + /** + * + * @param delegateFunction {@link BiFunction} to process the data + * @param dltDestination DLT destination + * @param dltPublishingContext {@link DltPublishingContext} + * @param recordTimeSupplier Supplier for downstream record timestamp + */ public DltAwareProcessor(BiFunction> delegateFunction, String dltDestination, - DltSenderContext dltSenderContext, Supplier recordTimeSupplier) { + DltPublishingContext dltPublishingContext, Supplier recordTimeSupplier) { this.delegateFunction = delegateFunction; this.recordTimeSupplier = recordTimeSupplier; Assert.isTrue(StringUtils.hasText(dltDestination), "DLT Destination topic must be provided."); this.dltDestination = dltDestination; - Assert.notNull(dltSenderContext, "DltSenderContext cannot be null"); - this.dltSenderContext = dltSenderContext; + Assert.notNull(dltPublishingContext, "DltSenderContext cannot be null"); + this.dltPublishingContext = dltPublishingContext; } + /** + * + * @param delegateFunction {@link BiFunction} to process the data + * @param processorRecordRecoverer {@link BiConsumer} that recovers failed records + */ public DltAwareProcessor(BiFunction> delegateFunction, BiConsumer, Exception> processorRecordRecoverer) { this(delegateFunction, processorRecordRecoverer, System::currentTimeMillis); } + /** + * + * @param delegateFunction {@link BiFunction} to process the data + * @param processorRecordRecoverer {@link BiConsumer} that recovers failed records + * @param recordTimeSupplier Supplier for downstream record timestamp + */ public DltAwareProcessor(BiFunction> delegateFunction, BiConsumer, Exception> processorRecordRecoverer, Supplier recordTimeSupplier) { this.delegateFunction = delegateFunction; @@ -104,7 +148,7 @@ public class DltAwareProcessor implements Processor, Exception> defaultProcessorRecordRecoverer() { return (r, e) -> { - StreamBridge streamBridge = this.dltSenderContext.getStreamBridge(); + StreamBridge streamBridge = this.dltPublishingContext.getStreamBridge(); if (streamBridge != null) { streamBridge.send(this.dltDestination, r.value()); } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltSenderContext.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltPublishingContext.java similarity index 82% rename from binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltSenderContext.java rename to binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltPublishingContext.java index 6bd89ad92..641449eb9 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltSenderContext.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/DltPublishingContext.java @@ -24,15 +24,18 @@ import org.springframework.context.ApplicationContextAware; import org.springframework.context.ConfigurableApplicationContext; /** + * The DltPublishingContext is meant to be used along with {@link DltAwareProcessor} + * when publishing failed record to a DLT. DltPublishingContext is particularly used + * for accessing framework beans such as {@link StreamBridge}. + * * @author Soby Chacko */ -public class DltSenderContext implements ApplicationContextAware, InitializingBean { +public class DltPublishingContext implements ApplicationContextAware, InitializingBean { private ConfigurableApplicationContext applicationContext; private StreamBridge streamBridge; - @Override public void afterPropertiesSet() { this.streamBridge = applicationContext.getBean(StreamBridge.class); diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index 4f7e491d3..430e2e7c9 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -394,8 +394,8 @@ public class KafkaStreamsBinderSupportAutoConfiguration { @Bean @ConditionalOnMissingBean - public DltSenderContext dltSenderContext() { - return new DltSenderContext(); + public DltPublishingContext dltSenderContext() { + return new DltPublishingContext(); } @Configuration(proxyBeanMethods = false) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DltAwareProcessorTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DltAwareProcessorTests.java new file mode 100644 index 000000000..744cf9f9f --- /dev/null +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/DltAwareProcessorTests.java @@ -0,0 +1,120 @@ +/* + * Copyright 2023-2023 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.integration; + +import java.util.Map; + +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.kstream.KStream; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.cloud.stream.binder.kafka.streams.DltAwareProcessor; +import org.springframework.cloud.stream.binder.kafka.streams.DltPublishingContext; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.kafka.core.DefaultKafkaConsumerFactory; +import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.test.EmbeddedKafkaBroker; +import org.springframework.kafka.test.condition.EmbeddedKafkaCondition; +import org.springframework.kafka.test.context.EmbeddedKafka; +import org.springframework.kafka.test.utils.KafkaTestUtils; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * @author Soby Chacko + */ +@EmbeddedKafka(topics = "hello-dlt-1") +public class DltAwareProcessorTests { + + private static final EmbeddedKafkaBroker embeddedKafka = EmbeddedKafkaCondition.getBroker(); + + private static Consumer consumer; + + @BeforeAll + public static void setUp() { + Map consumerProps = KafkaTestUtils.consumerProps("group", "false", + embeddedKafka); + consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( + consumerProps); + consumer = cf.createConsumer(); + embeddedKafka.consumeFromEmbeddedTopics(consumer, "hello-dlt-1"); + } + + @AfterAll + public static void tearDown() { + consumer.close(); + } + + @Test + void testDltAwareProcessor() { + SpringApplication app = new SpringApplication( + DltAwareProcessorTests.PublishToDltOnErrorApplication.class); + app.setWebApplicationType(WebApplicationType.NONE); + + try (ConfigurableApplicationContext context = app.run("--server.port=0", + "--spring.cloud.stream.bindings.errorStream-in-0.destination=error-stream-in", + "--spring.cloud.stream.kafka.streams.bindings.errorStream-in-0.consumer.application-id=test-error-stream-app-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.kafka.streams.binder.brokers=" + + embeddedKafka.getBrokersAsString())) { + receiveAndValidate("error-stream-in", "hello-dlt-1"); + } + } + + private void receiveAndValidate(String in, String out) { + Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( + senderProps); + try { + KafkaTemplate template = new KafkaTemplate<>(pf, true); + template.setDefaultTopic(in); + template.sendDefault("foobar"); + ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, + out); + assertThat(cr.value().contains("foobar")).isTrue(); + } + finally { + pf.destroy(); + } + } + + @EnableAutoConfiguration + static class PublishToDltOnErrorApplication { + + @Bean + public java.util.function.Consumer> errorStream(DltPublishingContext dltSenderContext) { + return input -> input + .process(() -> new DltAwareProcessor<>((k, v) -> { + throw new RuntimeException("error"); + }, "hello-dlt-1", dltSenderContext)); + } + } +}