From c7fa1ce275d3546cc9f9e309c2dc302c1ca58afd Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 8 Oct 2019 06:04:54 -0500 Subject: [PATCH] Fix how stream function properties for bindings are used --- .../kafka/streams/KafkaStreamsFunctionProcessor.java | 2 +- .../function/KafkaStreamsBindableProxyFactory.java | 4 ++-- .../KafkaStreamsBinderWordCountBranchesFunctionTests.java | 8 +++++--- .../function/StreamToGlobalKTableFunctionTests.java | 8 ++++++-- 4 files changed, 14 insertions(+), 8 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index 4807f3d78..4875cc828 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -260,7 +260,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro } private List getOutputBindings(String functionName, int outputs) { - List outputBindings = this.streamFunctionProperties.getOutputBindings().get(functionName); + List outputBindings = this.streamFunctionProperties.getOutputBindings(functionName); List outputBindingNames = new ArrayList<>(); if (!CollectionUtils.isEmpty(outputBindings)) { outputBindingNames.addAll(outputBindings); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java index e5c5fa1c6..0886485c7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java @@ -121,7 +121,7 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto // if the type is array, we need to do a late binding as we don't know the number of // output bindings at this point in the flow. - List outputBindings = streamFunctionProperties.getOutputBindings().get(this.functionName); + List outputBindings = streamFunctionProperties.getOutputBindings(this.functionName); String outputBinding = null; if (!CollectionUtils.isEmpty(outputBindings)) { @@ -168,7 +168,7 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto */ private List buildInputBindings() { List inputs = new ArrayList<>(); - List inputBindings = streamFunctionProperties.getInputBindings().get(this.functionName); + List inputBindings = streamFunctionProperties.getInputBindings(this.functionName); if (!CollectionUtils.isEmpty(inputBindings)) { inputs.addAll(inputBindings); return inputs; 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 ffc8c4bb3..8897a14af 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 @@ -79,13 +79,15 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { app.setWebApplicationType(WebApplicationType.NONE); ConfigurableApplicationContext context = app.run("--server.port=0", - "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.inputBindings.process=input", - "--spring.cloud.stream.function.outputBindings.process=output1,output2,output3", + "--spring.cloud.stream.function.bindings.process-in-0=input", "--spring.cloud.stream.bindings.input.destination=words", + "--spring.cloud.stream.function.bindings.process-out-0=output1", "--spring.cloud.stream.bindings.output1.destination=counts", + "--spring.cloud.stream.function.bindings.process-out-1=output2", "--spring.cloud.stream.bindings.output2.destination=foo", + "--spring.cloud.stream.function.bindings.process-out-2=output3", "--spring.cloud.stream.bindings.output3.destination=bar", + "--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/StreamToGlobalKTableFunctionTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java index 355be2893..7af1eddf6 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/function/StreamToGlobalKTableFunctionTests.java @@ -67,12 +67,16 @@ public class StreamToGlobalKTableFunctionTests { app.setWebApplicationType(WebApplicationType.NONE); try (ConfigurableApplicationContext ignored = app.run("--server.port=0", "--spring.jmx.enabled=false", - "--spring.cloud.stream.function.inputBindings.process=order,customer,product", - "--spring.cloud.stream.function.outputBindings.process=enriched-order", + "--spring.cloud.stream.function.definition=process", + "--spring.cloud.stream.function.bindings.process-in-0=order", + "--spring.cloud.stream.function.bindings.process-in-1=customer", + "--spring.cloud.stream.function.bindings.process-in-2=product", + "--spring.cloud.stream.function.bindings.process-out-0=enriched-order", "--spring.cloud.stream.bindings.order.destination=orders", "--spring.cloud.stream.bindings.customer.destination=customers", "--spring.cloud.stream.bindings.product.destination=products", "--spring.cloud.stream.bindings.enriched-order.destination=enriched-order", + "--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" +