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));
+ }
+ }
+}