From b4c6ec8332c73ca62b14bb8ad7b1cbbd0644e411 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 31 Mar 2022 16:38:05 -0400 Subject: [PATCH] Concurrency property issues in KStream binder In Kafka Streams binder, when using a function with camelcase names, it causes issues for parsing binding level concurrency properties. Fixing this issue. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2316 --- .../AbstractKafkaStreamsBinderProcessor.java | 8 ++-- .../MultipleFunctionsInSameAppTests.java | 42 +++++++++---------- 2 files changed, 25 insertions(+), 25 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java index 9c787b58b..a7680e960 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -371,10 +371,10 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application public Object onSuccess(ConfigurationPropertyName name, Bindable target, BindContext context, Object result) { if (!concurrencyExplicitlyProvided[0]) { - - concurrencyExplicitlyProvided[0] = name.getLastElement(ConfigurationPropertyName.Form.UNIFORM) - .equals("concurrency") && - ConfigurationPropertyName.of("spring.cloud.stream.bindings." + inboundName + ".consumer").isAncestorOf(name); + concurrencyExplicitlyProvided[0] = name.getLastElement(ConfigurationPropertyName.Form.UNIFORM).equals("concurrency") && + // name is normalized to contain only uniform elements and thus safe to call toLowerCase here. + ConfigurationPropertyName.of("spring.cloud.stream.bindings." + inboundName.toLowerCase() + ".consumer") + .isAncestorOf(name); } return result; } diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java index d9679f926..f85218999 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/MultipleFunctionsInSameAppTests.java @@ -86,26 +86,26 @@ public class MultipleFunctionsInSameAppTests { try (ConfigurableApplicationContext context = app.run( "--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.definition=process;analyze;anotherProcess;yetAnotherProcess", - "--spring.cloud.stream.bindings.process-in-0.destination=purchases", - "--spring.cloud.stream.bindings.process-out-0.destination=coffee", - "--spring.cloud.stream.bindings.process-out-1.destination=electronics", + "--spring.cloud.stream.function.definition=processItem;analyze;anotherProcess;yetAnotherProcess", + "--spring.cloud.stream.bindings.processItem-in-0.destination=purchases", + "--spring.cloud.stream.bindings.processItem-out-0.destination=coffee", + "--spring.cloud.stream.bindings.processItem-out-1.destination=electronics", "--spring.cloud.stream.bindings.analyze-in-0.destination=coffee", "--spring.cloud.stream.bindings.analyze-in-1.destination=electronics", "--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.functions.processItem.applicationId=processItem-id-0", "--spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms=1000", - "--spring.cloud.stream.bindings.process-in-0.consumer.concurrency=2", + "--spring.cloud.stream.bindings.processItem-in-0.consumer.concurrency=2", "--spring.cloud.stream.bindings.analyze-in-0.consumer.concurrency=1", "--spring.cloud.stream.kafka.streams.binder.configuration.num.stream.threads=3", - "--spring.cloud.stream.kafka.streams.binder.functions.process.configuration.client.id=process-client", + "--spring.cloud.stream.kafka.streams.binder.functions.processItem.configuration.client.id=processItem-client", "--spring.cloud.stream.kafka.streams.binder.functions.analyze.configuration.client.id=analyze-client", "--spring.cloud.stream.kafka.streams.binder.functions.anotherProcess.configuration.client.id=anotherProcess-client", "--spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString())) { receiveAndValidate("purchases", "coffee", "electronics"); StreamsBuilderFactoryBean processStreamsBuilderFactoryBean = context - .getBean("&stream-builder-process", StreamsBuilderFactoryBean.class); + .getBean("&stream-builder-processItem", StreamsBuilderFactoryBean.class); StreamsBuilderFactoryBean analyzeStreamsBuilderFactoryBean = context .getBean("&stream-builder-analyze", StreamsBuilderFactoryBean.class); @@ -117,7 +117,7 @@ public class MultipleFunctionsInSameAppTests { final Properties analyzeStreamsConfiguration = analyzeStreamsBuilderFactoryBean.getStreamsConfiguration(); final Properties anotherProcessStreamsConfiguration = anotherProcessStreamsBuilderFactoryBean.getStreamsConfiguration(); - assertThat(processStreamsConfiguration.getProperty("client.id")).isEqualTo("process-client"); + assertThat(processStreamsConfiguration.getProperty("client.id")).isEqualTo("processItem-client"); assertThat(analyzeStreamsConfiguration.getProperty("client.id")).isEqualTo("analyze-client"); Integer concurrency = (Integer) processStreamsConfiguration.get(StreamsConfig.NUM_STREAM_THREADS_CONFIG); @@ -144,14 +144,14 @@ public class MultipleFunctionsInSameAppTests { try (ConfigurableApplicationContext context = app.run( "--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.definition=process;analyze", - "--spring.cloud.stream.bindings.process-in-0.destination=purchases", - "--spring.cloud.stream.kafka.streams.bindings.process-in-0.consumer.startOffset=latest", - "--spring.cloud.stream.bindings.process-in-0.binder=kafka1", - "--spring.cloud.stream.bindings.process-out-0.destination=coffee", - "--spring.cloud.stream.bindings.process-out-0.binder=kafka1", - "--spring.cloud.stream.bindings.process-out-1.destination=electronics", - "--spring.cloud.stream.bindings.process-out-1.binder=kafka1", + "--spring.cloud.stream.function.definition=processItem;analyze", + "--spring.cloud.stream.bindings.processItem-in-0.destination=purchases", + "--spring.cloud.stream.kafka.streams.bindings.processItem-in-0.consumer.startOffset=latest", + "--spring.cloud.stream.bindings.processItem-in-0.binder=kafka1", + "--spring.cloud.stream.bindings.processItem-out-0.destination=coffee", + "--spring.cloud.stream.bindings.processItem-out-0.binder=kafka1", + "--spring.cloud.stream.bindings.processItem-out-1.destination=electronics", + "--spring.cloud.stream.bindings.processItem-out-1.binder=kafka1", "--spring.cloud.stream.bindings.analyze-in-0.destination=coffee", "--spring.cloud.stream.bindings.analyze-in-0.binder=kafka2", "--spring.cloud.stream.bindings.analyze-in-1.destination=electronics", @@ -160,7 +160,7 @@ public class MultipleFunctionsInSameAppTests { "--spring.cloud.stream.binders.kafka1.type=kstream", "--spring.cloud.stream.binders.kafka1.environment.spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.binders.kafka1.environment.spring.cloud.stream.kafka.streams.binder.applicationId=my-app-1", - "--spring.cloud.stream.binders.kafka1.environment.spring.cloud.stream.kafka.streams.binder.configuration.client.id=process-client", + "--spring.cloud.stream.binders.kafka1.environment.spring.cloud.stream.kafka.streams.binder.configuration.client.id=processItem-client", "--spring.cloud.stream.binders.kafka2.type=kstream", "--spring.cloud.stream.binders.kafka2.environment.spring.cloud.stream.kafka.streams.binder.brokers=" + embeddedKafka.getBrokersAsString(), "--spring.cloud.stream.binders.kafka2.environment.spring.cloud.stream.kafka.streams.binder.applicationId=my-app-2", @@ -170,7 +170,7 @@ public class MultipleFunctionsInSameAppTests { receiveAndValidate("purchases", "coffee", "electronics"); StreamsBuilderFactoryBean processStreamsBuilderFactoryBean = context - .getBean("&stream-builder-process", StreamsBuilderFactoryBean.class); + .getBean("&stream-builder-processItem", StreamsBuilderFactoryBean.class); StreamsBuilderFactoryBean analyzeStreamsBuilderFactoryBean = context .getBean("&stream-builder-analyze", StreamsBuilderFactoryBean.class); @@ -180,7 +180,7 @@ public class MultipleFunctionsInSameAppTests { assertThat(processStreamsConfiguration.getProperty("application.id")).isEqualTo("my-app-1"); assertThat(analyzeStreamsConfiguration.getProperty("application.id")).isEqualTo("my-app-2"); - assertThat(processStreamsConfiguration.getProperty("client.id")).isEqualTo("process-client"); + assertThat(processStreamsConfiguration.getProperty("client.id")).isEqualTo("processItem-client"); assertThat(analyzeStreamsConfiguration.getProperty("client.id")).isEqualTo("analyze-client"); Integer concurrency = (Integer) analyzeStreamsConfiguration.get(StreamsConfig.NUM_STREAM_THREADS_CONFIG); @@ -216,7 +216,7 @@ public class MultipleFunctionsInSameAppTests { public static class MultipleFunctionsInSameApp { @Bean - public Function, KStream[]> process() { + public Function, KStream[]> processItem() { return input -> input.branch( (s, p) -> p.equalsIgnoreCase("coffee"), (s, p) -> p.equalsIgnoreCase("electronics"));