diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 4a20cc9fd..6e2b86429 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -213,11 +213,11 @@ public class KafkaMessageChannelBinder extends private final TransactionTemplate transactionTemplate; - private final KafkaBindingRebalanceListener rebalanceListener; + private KafkaBindingRebalanceListener rebalanceListener; - private final DlqPartitionFunction dlqPartitionFunction; + private DlqPartitionFunction dlqPartitionFunction; - private final DlqDestinationResolver dlqDestinationResolver; + private DlqDestinationResolver dlqDestinationResolver; private final Map ackModeInfo = new ConcurrentHashMap<>(); @@ -309,6 +309,18 @@ public class KafkaMessageChannelBinder extends this.clientFactoryCustomizer = customizer; } + public void setRebalanceListener(KafkaBindingRebalanceListener rebalanceListener) { + this.rebalanceListener = rebalanceListener; + } + + public void setDlqPartitionFunction(DlqPartitionFunction dlqPartitionFunction) { + this.dlqPartitionFunction = dlqPartitionFunction; + } + + public void setDlqDestinationResolver(DlqDestinationResolver dlqDestinationResolver) { + this.dlqDestinationResolver = dlqDestinationResolver; + } + Map getTopicsInUse() { return this.topicsInUse; }