From 50ec8f0919ff4a322f238d12e42f45f4a73bf7c9 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 15 Oct 2020 13:26:07 -0400 Subject: [PATCH] Change Log etc. Replication Factor If the user has not explicitly set the `SreamsConfig.REPLICATION_FACTOR_CONFIG`, set it from the binder property. This is used for infrastructure topics (change logs and repartition topics). --- .../kafka/streams/AbstractKafkaStreamsBinderProcessor.java | 4 ++++ .../streams/KafkaStreamsBinderSupportAutoConfiguration.java | 4 ++++ 2 files changed, 8 insertions(+) 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)); }