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 b19348628..37a7c2977 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 @@ -65,11 +65,11 @@ public class DltAwareProcessor implements Processor> delegateFunction, BiConsumer, Exception> processorRecordRecoverer) { - this(delegateFunction, System::currentTimeMillis, processorRecordRecoverer); + this(delegateFunction, processorRecordRecoverer, System::currentTimeMillis); } public DltAwareProcessor(BiFunction> delegateFunction, - Supplier recordTimeSupplier, BiConsumer, Exception> processorRecordRecoverer) { + BiConsumer, Exception> processorRecordRecoverer, Supplier recordTimeSupplier) { this.delegateFunction = delegateFunction; this.recordTimeSupplier = recordTimeSupplier; Assert.notNull(processorRecordRecoverer, "You must provide a valid processor recoverer");