From 4a4ce6975a32bd6dd40175a7f7f0fa80ca861486 Mon Sep 17 00:00:00 2001 From: Steven PG Date: Wed, 18 Oct 2023 19:10:04 -0400 Subject: [PATCH] Add DltAwareProcessor trace logging Since exception is not propogated, adding a trace log allows an optional way for a developer to utilize the defaultRecoverer and still do some basic review of a given exception. --- .../stream/binder/kafka/streams/DltAwareProcessor.java | 6 ++++++ 1 file changed, 6 insertions(+) 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 3c4c9816c..b9ed5c780 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 @@ -19,6 +19,8 @@ package org.springframework.cloud.stream.binder.kafka.streams; import java.util.function.BiConsumer; import java.util.function.Function; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.apache.kafka.streams.processor.api.Record; import org.springframework.cloud.stream.function.StreamBridge; @@ -32,10 +34,13 @@ import org.springframework.util.StringUtils; * Custom {@link RecordRecoverableProcessor} that is capable of sending the failed record to a DLT during recovery. * * @author Soby Chacko + * @author Steven Gantz * @sinc 4.1.0 */ public class DltAwareProcessor extends RecordRecoverableProcessor { + private static final Log LOG = LogFactory.getLog(DltAwareProcessor.class); + /** * DLT destination. */ @@ -74,6 +79,7 @@ public class DltAwareProcessor extends RecordRecoverablePr if (streamBridge != null) { Message message = MessageBuilder.withPayload(r.value()) .setHeader(KafkaHeaders.KEY, r.key()).build(); + DltAwareProcessor.LOG.trace("Recovered from Exception: ", e); streamBridge.send(this.dltDestination, message); } };