From 6bb4da89ce106a184b2ced896b2f11bfe8daaca7 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 24 Jul 2023 15:02:10 -0400 Subject: [PATCH] DLTAwareProcessor (Kafka Streams binder) changes - Optional BiConsumer for processor record recoverer - Optional Supplier for downstream record timestamp --- .../kafka/streams/DltAwareProcessor.java | 58 +++++++++++++++---- .../kafka/streams/DltSenderContext.java | 1 + ...StreamsBinderSupportAutoConfiguration.java | 6 ++ 3 files changed, 55 insertions(+), 10 deletions(-) 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 d5704aa47..b19348628 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 @@ -16,7 +16,9 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.function.BiConsumer; import java.util.function.BiFunction; +import java.util.function.Supplier; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.processor.api.Processor; @@ -24,6 +26,8 @@ import org.apache.kafka.streams.processor.api.ProcessorContext; import org.apache.kafka.streams.processor.api.Record; import org.springframework.cloud.stream.function.StreamBridge; +import org.springframework.util.Assert; +import org.springframework.util.StringUtils; /** * @@ -34,18 +38,44 @@ public class DltAwareProcessor implements Processor> delegateFunction; + private final Supplier recordTimeSupplier; + + private String dltDestination; + + private DltSenderContext dltSenderContext; + + private BiConsumer, Exception> processorRecordRecoverer; + private ProcessorContext context; - private final String dltDestination; + public DltAwareProcessor(BiFunction> delegateFunction, String dltDestination, + DltSenderContext dltSenderContext) { + this(delegateFunction, dltDestination, dltSenderContext, System::currentTimeMillis); + } - private final DltSenderContext dltSenderContext; - - public DltAwareProcessor(BiFunction> businessLogic, String dltDestination, DltSenderContext dltSenderContext) { - this.delegateFunction = businessLogic; + public DltAwareProcessor(BiFunction> delegateFunction, String dltDestination, + DltSenderContext dltSenderContext, 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; } + public DltAwareProcessor(BiFunction> delegateFunction, + BiConsumer, Exception> processorRecordRecoverer) { + this(delegateFunction, System::currentTimeMillis, processorRecordRecoverer); + } + + public DltAwareProcessor(BiFunction> delegateFunction, + Supplier recordTimeSupplier, BiConsumer, Exception> processorRecordRecoverer) { + this.delegateFunction = delegateFunction; + this.recordTimeSupplier = recordTimeSupplier; + Assert.notNull(processorRecordRecoverer, "You must provide a valid processor recoverer"); + this.processorRecordRecoverer = processorRecordRecoverer; + } + @Override public void init(ProcessorContext context) { Processor.super.init(context); @@ -56,15 +86,14 @@ public class DltAwareProcessor implements Processor record) { try { KeyValue keyValue = this.delegateFunction.apply(record.key(), record.value()); - //TODO: What should be the timestamp? - Record downstreamRecord = new Record<>(keyValue.key, keyValue.value, System.currentTimeMillis(), record.headers()); + Record downstreamRecord = new Record<>(keyValue.key, keyValue.value, recordTimeSupplier.get(), record.headers()); this.context.forward(downstreamRecord); } catch (Exception exception) { - StreamBridge streamBridge = this.dltSenderContext.getStreamBridge(); - if (streamBridge != null) { - streamBridge.send(dltDestination, record.value()); + if (this.processorRecordRecoverer == null) { + this.processorRecordRecoverer = defaultProcessorRecordRecoverer(); } + this.processorRecordRecoverer.accept(record, exception); } } @@ -73,4 +102,13 @@ public class DltAwareProcessor implements Processor, Exception> defaultProcessorRecordRecoverer() { + return (r, e) -> { + StreamBridge streamBridge = this.dltSenderContext.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/DltSenderContext.java index aa7c0f8cc..6bd89ad92 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/DltSenderContext.java @@ -32,6 +32,7 @@ public class DltSenderContext implements ApplicationContextAware, InitializingBe 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 4eb1ff562..4f7e491d3 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 @@ -392,6 +392,12 @@ public class KafkaStreamsBinderSupportAutoConfiguration { return new EncodingDecodingBindAdviceHandler(); } + @Bean + @ConditionalOnMissingBean + public DltSenderContext dltSenderContext() { + return new DltSenderContext(); + } + @Configuration(proxyBeanMethods = false) @ConditionalOnMissingBean(value = KafkaStreamsBinderMetrics.class, name = "outerContext") @ConditionalOnClass(name = "io.micrometer.core.instrument.MeterRegistry")