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