diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java index cbe74a63f..b112c91bb 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -247,6 +247,10 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application if (!ObjectUtils.isEmpty(multiBinderKafkaStreamsBinderConfigurationProperties.getConfiguration())) { streamConfiguration.putAll(multiBinderKafkaStreamsBinderConfigurationProperties.getConfiguration()); } + if (!streamConfiguration.containsKey(StreamsConfig.REPLICATION_FACTOR_CONFIG)) { + streamConfiguration.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, + (int) multiBinderKafkaStreamsBinderConfigurationProperties.getReplicationFactor()); + } } //this is only used primarily for StreamListener based processors. Although in theory, functions can use it, diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index 7dd93f5e6..c29ecd851 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -270,6 +270,10 @@ public class KafkaStreamsBinderSupportAutoConfiguration { if (!ObjectUtils.isEmpty(configProperties.getConfiguration())) { properties.putAll(configProperties.getConfiguration()); } + if (!properties.containsKey(StreamsConfig.REPLICATION_FACTOR_CONFIG)) { + properties.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, + (int) configProperties.getReplicationFactor()); + } return properties.entrySet().stream().collect( Collectors.toMap((e) -> String.valueOf(e.getKey()), Map.Entry::getValue)); }