diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java index 4debb72c8..1a9993db8 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java @@ -36,6 +36,7 @@ import org.springframework.util.StringUtils; * @author Marius Bogoevici * @author Soby Chacko * @author Gary Russell + * @author Rafal Zukowski */ @ConfigurationProperties(prefix = "spring.cloud.stream.kafka.binder") public class KafkaBinderConfigurationProperties { @@ -83,7 +84,7 @@ public class KafkaBinderConfigurationProperties { */ private int zkConnectionTimeout = 10000; - private int requiredAcks = 1; + private String requiredAcks = "1"; private int replicationFactor = 1; @@ -219,11 +220,15 @@ public class KafkaBinderConfigurationProperties { this.maxWait = maxWait; } - public int getRequiredAcks() { + public String getRequiredAcks() { return this.requiredAcks; } public void setRequiredAcks(int requiredAcks) { + this.requiredAcks = String.valueOf(requiredAcks); + } + + public void setRequiredAcks(String requiredAcks) { this.requiredAcks = requiredAcks; } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsDlqDispatch.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsDlqDispatch.java index fe7c489ce..4c7db8b75 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsDlqDispatch.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsDlqDispatch.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; /** * @author Soby Chacko + * @author Rafal Zukowski */ import java.util.HashMap; import java.util.Map; @@ -103,7 +104,7 @@ class KafkaStreamsDlqDispatch { Map props = new HashMap<>(); props.put(ProducerConfig.RETRIES_CONFIG, 0); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); - props.put(ProducerConfig.ACKS_CONFIG, String.valueOf(configurationProperties.getRequiredAcks())); + props.put(ProducerConfig.ACKS_CONFIG, configurationProperties.getRequiredAcks()); if (!ObjectUtils.isEmpty(configurationProperties.getProducerConfiguration())) { props.putAll(configurationProperties.getProducerConfiguration()); }