From bf6a227f3275182f270bb7def810d4580aea77b1 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Sat, 28 Sep 2019 16:03:41 -0400 Subject: [PATCH] Function level binding properties If there are multiple functions in a Kafka Streams application, and if they want to have a separate set of configuration for each, then it should be able to set that at the function level. For e.g. spring.cloud.stream.kafka.streams.binder.functions.... Resolves #757 --- .../AbstractKafkaStreamsBinderProcessor.java | 28 ++++++++--- ...aStreamsBinderConfigurationProperties.java | 49 +++++++++++++++---- ...sBinderWordCountBranchesFunctionTests.java | 4 +- .../KafkaStreamsFunctionStateStoreTests.java | 2 +- .../MultipleFunctionsInSameAppTests.java | 22 +++++++-- 5 files changed, 79 insertions(+), 26 deletions(-) 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 94d414aac..bb2e81bf0 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 @@ -157,6 +157,25 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application Map streamConfigGlobalProperties = applicationContext .getBean("streamConfigGlobalProperties", Map.class); + if (kafkaStreamsBinderConfigurationProperties != null) { + final Map functionConfigMap = kafkaStreamsBinderConfigurationProperties.getFunctions(); + if (!CollectionUtils.isEmpty(functionConfigMap)) { + final KafkaStreamsBinderConfigurationProperties.Functions functionConfig = functionConfigMap.get(beanNamePostPrefix); + final Map functionSpecificConfig = functionConfig.getConfiguration(); + if (!CollectionUtils.isEmpty(functionSpecificConfig)) { + streamConfigGlobalProperties.putAll(functionSpecificConfig); + } + + String applicationId = functionConfig.getApplicationId(); + if (!StringUtils.isEmpty(applicationId)) { + streamConfigGlobalProperties.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId); + } + } + } + + //this is only used primarily for StreamListener based processors. Although in theory, functions can use it, + //it is ideal for functions to use the approach used in the above if statement by using a property like + //spring.cloud.stream.kafka.streams.binder.functions.process.configuration.num.threads (assuming that process is the function name). KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties .getExtendedConsumerProperties(inboundName); streamConfigGlobalProperties @@ -165,17 +184,12 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application String bindingLevelApplicationId = extendedConsumerProperties.getApplicationId(); // override application.id if set at the individual binding level. // We provide this for backward compatibility with StreamListener based processors. - // For function based processors see the next else if conditional block + // For function based processors see the approach used above + // (i.e. use a property like spring.cloud.stream.kafka.streams.binder.functions.process.applicationId). if (StringUtils.hasText(bindingLevelApplicationId)) { streamConfigGlobalProperties.put(StreamsConfig.APPLICATION_ID_CONFIG, bindingLevelApplicationId); } - else if (kafkaStreamsBinderConfigurationProperties != null && !CollectionUtils.isEmpty(kafkaStreamsBinderConfigurationProperties.getFunctions())) { - String applicationId = kafkaStreamsBinderConfigurationProperties.getFunctions().get(beanNamePostPrefix + ".applicationId"); - if (!StringUtils.isEmpty(applicationId)) { - streamConfigGlobalProperties.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId); - } - } //If the application id is not set by any mechanism, then generate it. streamConfigGlobalProperties.computeIfAbsent(StreamsConfig.APPLICATION_ID_CONFIG, diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java index 48a618b4e..caa6dd7f2 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsBinderConfigurationProperties.java @@ -57,10 +57,18 @@ public class KafkaStreamsBinderConfigurationProperties private String applicationId; - private Map functions = new HashMap<>(); - private StateStoreRetry stateStoreRetry = new StateStoreRetry(); + private Map functions = new HashMap<>(); + + public Map getFunctions() { + return functions; + } + + public void setFunctions(Map functions) { + this.functions = functions; + } + public StateStoreRetry getStateStoreRetry() { return stateStoreRetry; } @@ -69,14 +77,6 @@ public class KafkaStreamsBinderConfigurationProperties this.stateStoreRetry = stateStoreRetry; } - public Map getFunctions() { - return functions; - } - - public void setFunctions(Map functions) { - this.functions = functions; - } - public String getApplicationId() { return this.applicationId; } @@ -125,4 +125,33 @@ public class KafkaStreamsBinderConfigurationProperties } } + public static class Functions { + + /** + * Function specific application id. + */ + private String applicationId; + + /** + * Funcion specific configuraiton to use. + */ + private Map configuration; + + public String getApplicationId() { + return applicationId; + } + + public void setApplicationId(String applicationId) { + this.applicationId = applicationId; + } + + public Map getConfiguration() { + return configuration; + } + + public void setConfiguration(Map configuration) { + this.configuration = configuration; + } + } + } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java index 23c7be3a1..ffc8c4bb3 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBinderWordCountBranchesFunctionTests.java @@ -91,9 +91,7 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { "=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.timeWindow.length=5000", - "--spring.cloud.stream.kafka.streams.timeWindow.advanceBy=0", - "--spring.cloud.stream.kafka.streams.bindings.input.consumer.applicationId" + + "--spring.cloud.stream.kafka.streams.binder.applicationId" + "=KafkaStreamsBinderWordCountBranchesFunctionTests-abc", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString()); try { diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java index aa61bc78c..2c7f2fc35 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsFunctionStateStoreTests.java @@ -59,7 +59,7 @@ public class KafkaStreamsFunctionStateStoreTests { try (ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.input.destination=words", - "--spring.cloud.stream.kafka.streams.default.consumer.application-id=testKafkaStreamsFuncionWithMultipleStateStores", + "--spring.cloud.stream.kafka.streams.binder.application-id=testKafkaStreamsFuncionWithMultipleStateStores", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + "=org.apache.kafka.common.serialization.Serdes$StringSerde", diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/MultipleFunctionsInSameAppTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/MultipleFunctionsInSameAppTests.java index d50332f4e..fee4cd9fc 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/MultipleFunctionsInSameAppTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/MultipleFunctionsInSameAppTests.java @@ -17,6 +17,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; import java.util.Map; +import java.util.Properties; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.function.BiConsumer; @@ -36,6 +37,7 @@ import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; +import org.springframework.kafka.config.StreamsBuilderFactoryBean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -78,7 +80,7 @@ public class MultipleFunctionsInSameAppTests { SpringApplication app = new SpringApplication(MultipleFunctionsInSameApp.class); app.setWebApplicationType(WebApplicationType.NONE); - try (ConfigurableApplicationContext ignored = app.run( + try (ConfigurableApplicationContext context = app.run( "--server.port=0", "--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.process-in-0.destination=purchases", @@ -89,12 +91,22 @@ public class MultipleFunctionsInSameAppTests { "--spring.cloud.stream.kafka.streams.binder.functions.analyze.applicationId=analyze-id-0", "--spring.cloud.stream.kafka.streams.binder.functions.process.applicationId=process-id-0", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", - "--spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde" + - "=org.apache.kafka.common.serialization.Serdes$StringSerde", + "--spring.cloud.stream.kafka.streams.binder.functions.process.configuration.client.id=process-client", + "--spring.cloud.stream.kafka.streams.binder.functions.analyze.configuration.client.id=analyze-client", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { receiveAndValidate("purchases", "coffee", "electronics"); + + StreamsBuilderFactoryBean processStreamsBuilderFactoryBean = context + .getBean("&stream-builder-process", StreamsBuilderFactoryBean.class); + + StreamsBuilderFactoryBean analyzeStreamsBuilderFactoryBean = context + .getBean("&stream-builder-analyze", StreamsBuilderFactoryBean.class); + + final Properties processStreamsConfiguration = processStreamsBuilderFactoryBean.getStreamsConfiguration(); + final Properties analyzeStreamsConfiguration = analyzeStreamsBuilderFactoryBean.getStreamsConfiguration(); + + assertThat(processStreamsConfiguration.getProperty("client.id")).isEqualTo("process-client"); + assertThat(analyzeStreamsConfiguration.getProperty("client.id")).isEqualTo("analyze-client"); } }