diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
index 0a4d2561d..2cd6e5958 100644
--- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
+++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java
@@ -216,6 +216,13 @@ public class KafkaConsumerProperties {
*/
private String commonErrorHandlerBeanName;
+ /**
+ * When using the reactive binder, automatically commit offsets when records received
+ * by the poll have been processed.
+ * @since 4.0.2
+ */
+ private boolean reactiveAutoCommit;
+
/**
* @return if each record needs to be acknowledged.
*
@@ -542,4 +549,14 @@ public class KafkaConsumerProperties {
public void setCommonErrorHandlerBeanName(String commonErrorHandlerBeanName) {
this.commonErrorHandlerBeanName = commonErrorHandlerBeanName;
}
+
+ public boolean isReactiveAutoCommit() {
+ return this.reactiveAutoCommit;
+ }
+
+ public void setReactiveAutoCommit(boolean reactiveAutoCommit) {
+ this.reactiveAutoCommit = reactiveAutoCommit;
+ }
+
+
}
diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/pom.xml b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/pom.xml
index 4695f0a60..b108204f1 100644
--- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/pom.xml
+++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/pom.xml
@@ -48,7 +48,7 @@
io.projectreactor.kafka
reactor-kafka
- 1.3.8
+ 1.3.16
org.springframework.boot
diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java
index 736ef26a3..695200787 100644
--- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java
+++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java
@@ -25,12 +25,14 @@ import java.util.concurrent.atomic.AtomicReference;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
+import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Sinks;
import reactor.kafka.receiver.KafkaReceiver;
import reactor.kafka.receiver.ReceiverOptions;
+import reactor.kafka.receiver.ReceiverRecord;
import reactor.kafka.sender.KafkaSender;
import reactor.kafka.sender.SenderOptions;
import reactor.kafka.sender.SenderRecord;
@@ -56,6 +58,9 @@ import org.springframework.context.Lifecycle;
import org.springframework.integration.core.MessageProducer;
import org.springframework.integration.endpoint.MessageProducerSupport;
import org.springframework.integration.handler.AbstractMessageHandler;
+import org.springframework.integration.support.MessageBuilder;
+import org.springframework.kafka.support.KafkaHeaders;
+import org.springframework.kafka.support.converter.KafkaMessageHeaders;
import org.springframework.kafka.support.converter.MessageConverter;
import org.springframework.kafka.support.converter.MessagingMessageConverter;
import org.springframework.kafka.support.converter.RecordMessageConverter;
@@ -201,11 +206,24 @@ public class ReactorKafkaBinder
protected void doStart() {
List>> fluxes = new ArrayList<>();
int concurrency = properties.getConcurrency();
+ boolean autoCommit = properties.getExtension().isReactiveAutoCommit();
for (int i = 0; i < concurrency; i++) {
- fluxes.add(this.receivers.get(i)
- .receive()
- .map(record -> (Message