From 3a5aed61c9744bec4776745db1720e0f3f34df44 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 10 Jan 2018 14:19:09 -0500 Subject: [PATCH] Redesigning the branching support in KStream binder Instead of relying on a property based approach, use proper multiple output bindings to support the branching feature. This is accomplished through a combination of using `SendTo` annotation with multiple outputs and overriding the default StreamListener method setup orchestration in the binder. --- .../src/main/asciidoc/overview.adoc | 88 +++++---- .../stream/binder/kstream/KStreamBinder.java | 58 +----- ...StreamListenerSetupMethodOrchestrator.java | 183 ++++++++++++++++++ .../config/KStreamBinderConfiguration.java | 7 - ...KStreamBinderSupportAutoConfiguration.java | 10 + .../config/KStreamProducerProperties.java | 9 - ...CountMultipleBranchesIntegrationTests.java | 48 +++-- 7 files changed, 276 insertions(+), 127 deletions(-) create mode 100644 spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamListenerSetupMethodOrchestrator.java diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc index c246ab446..4911311dc 100644 --- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc @@ -608,26 +608,16 @@ Keys will not get converted, but if the Serdes are different for keys from what Kafka Streams allow outbound data to be split into multiple topics based on some predicates. Spring Cloud Stream Kafka Streams binder provides support for this feature without losing the overall programming model exposed through `StreamListener` in the end user application. You write the application in the usual way as demonstrated above in the word count example. -The actual splitting and branching into multiple topics are done by the framework behind the scenes. -When using the branching feature, you are required to do two things. -First, you need to provide the following property that specifies the extra branches (topics) in the order. -The first topic will always be the one specified through the main outbound destination. - -`spring.cloud.stream.kstream.bindings.output.producer.additionalBranches=foo,bar` - -If your main output destination is foobar provided through `spring.cloud.stream.bindings.output.destination=foobar`, then your 3 output branches (topics) are foobar, foo and bar. - -Second, you need to provide a `Bean` in your application context, that returns a `org.apache.kafka.streams.kstream.Predicate[]`. -The presence of this bean is the trigger to the framework to perform branching into multiple topics. -The individual Predicates in this bean, should match with the output branches in the same order. -Each Predicate should get its own output branch, otherwise, it fails. -If you provide more output branches than there are Predicates, that is fine, but the number of branches cannot be less than the Predicates. +When using the branching feature, you are required to do a few things. +First, you need to make sure that your return type is `KStream[]` instead of a regular `KStream`. +Then you need to use the `SendTo` annotation containing the output bindings in the order (example below). +For each of these output bindings, you need to configure destination, content-type etc. as required by any other standard Spring Cloud Stream application Here is an example: [source] ---- -@EnableBinding(KStreamProcessor.class) +@EnableBinding(KStreamProcessorWithBranches.class) @EnableAutoConfiguration public static class WordCountProcessorApplication { @@ -635,25 +625,37 @@ public static class WordCountProcessorApplication { private TimeWindows timeWindows; @StreamListener("input") - @SendTo("output") - public KStream process(KStream input) { - return input - .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) - .groupBy((key, value) -> value) - .windowedBy(timeWindows) - .count(Materialized.as("WordCounts-1")) - .toStream() - .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))); + @SendTo({"output1","output2","output3}) + public KStream[] process(KStream input) { + + Predicate isEnglish = (k, v) -> v.word.equals("english"); + Predicate isFrench = (k, v) -> v.word.equals("french"); + Predicate isSpanish = (k, v) -> v.word.equals("spanish"); + + return input + .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) + .groupBy((key, value) -> value) + .windowedBy(timeWindows) + .count(Materialized.as("WordCounts-1")) + .toStream() + .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))) + .branch(isEnglish, isFrench, isSpanish); } - @Bean - public Predicate[] predicates() { - Predicate isEnglish = (k, v) -> v.word.equals("english"); - Predicate isFrench = (k, v) -> v.word.equals("french"); - Predicate isSpanish = (k, v) -> v.word.equals("spanish"); - return new Predicate[] {isEnglish, isFrench, isSpanish}; - } + interface KStreamProcessorWithBranches { + @Input("input") + KStream input(); + + @Output("output1") + KStream output1(); + + @Output("output2") + KStream output2(); + + @Output("output3") + KStream output3(); + } } ---- @@ -661,16 +663,25 @@ Then in the properties: [source] ---- -spring.cloud.stream.bindings.output.contentType: application/json +spring.cloud.stream.bindings.output1.contentType: application/json +spring.cloud.stream.bindings.output2.contentType: application/json +spring.cloud.stream.bindings.output3.contentType: application/json spring.cloud.stream.kstream.binder.configuration.commit.interval.ms: 1000 spring.cloud.stream.kstream.binder.configuration: key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde -spring.cloud.stream.bindings.output: - destination: foobar +spring.cloud.stream.bindings.output1: + destination: foo + producer: + headerMode: raw +spring.cloud.stream.bindings.output2: + destination: bar + producer: + headerMode: raw +spring.cloud.stream.bindings.output3: + destination: fox producer: headerMode: raw -spring.cloud.stream.kstream.bindings.output.producer.additionalBranches: foo,bar spring.cloud.stream.bindings.input: destination: words consumer: @@ -710,13 +721,6 @@ spring.cloud.stream.kstream.bindings.output.producer.keySerde=org.apache.kafka.c spring.cloud.stream.kstream.bindings.output.producer.valueSerde=org.apache.kafka.common.serialization.Serdes$LongSerde ---- -Additional output branches: - -[source] ----- -spring.cloud.stream.kstream.bindings.output.producer.additionalBranches (comma separated values) ----- - TimeWindow properties: [source] diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java index 4cbb680e0..943d2ee5d 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBinder.java @@ -24,7 +24,6 @@ import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KeyValueMapper; -import org.apache.kafka.streams.kstream.Predicate; import org.apache.kafka.streams.kstream.Produced; import org.springframework.cloud.stream.binder.AbstractBinder; @@ -40,7 +39,6 @@ import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProv import org.springframework.cloud.stream.binder.kstream.config.KStreamConsumerProperties; import org.springframework.cloud.stream.binder.kstream.config.KStreamExtendedBindingProperties; import org.springframework.cloud.stream.binder.kstream.config.KStreamProducerProperties; -import org.springframework.util.Assert; import org.springframework.util.StringUtils; /** @@ -59,8 +57,6 @@ public class KStreamBinder extends private final KafkaBinderConfigurationProperties binderConfigurationProperties; - private Predicate[] predicates; - private final MessageConversionDelegate messageConversionDelegate; public KStreamBinder(KafkaBinderConfigurationProperties binderConfigurationProperties, @@ -93,26 +89,10 @@ public class KStreamBinder extends new KafkaProducerProperties()); this.kafkaTopicProvisioner.provisionProducerDestination(name, extendedProducerProperties); - String[] branches = new String[]{}; - if (predicates != null && predicates.length > 0) { - String additionalBranches = properties.getExtension().getAdditionalBranches(); - if (!StringUtils.hasText(additionalBranches)) { - Assert.isTrue(predicates.length == 1, "More than 1 predicate bean found, but no additional output branches"); - } - else { - branches = StringUtils.commaDelimitedListToStringArray(additionalBranches); - Assert.isTrue(branches.length + 1 >= predicates.length, - "Number of output topics and org.apache.kafka.streams.kstream.Predicate[] beans don't match"); - for (String branch : branches) { - this.kafkaTopicProvisioner.provisionProducerDestination(branch, extendedProducerProperties); - } - } - } - Serde keySerde = getKeySerde(properties); Serde valueSerde = getValueSerde(properties); - to(properties.isUseNativeEncoding(), name, outboundBindTarget, (Serde) keySerde, (Serde) valueSerde, branches); + to(properties.isUseNativeEncoding(), name, outboundBindTarget, (Serde) keySerde, (Serde) valueSerde); return new DefaultBinding<>(name, null, outboundBindTarget, null); } @@ -173,42 +153,19 @@ public class KStreamBinder extends @SuppressWarnings("unchecked") private void to(boolean isNativeEncoding, String name, KStream outboundBindTarget, - Serde keySerde, Serde valueSerde, String[] branches) { + Serde keySerde, Serde valueSerde) { KeyValueMapper> keyValueMapper = null; if (!isNativeEncoding) { keyValueMapper = messageConversionDelegate.outboundKeyValueMapper(name); } - if (predicates != null && predicates.length > 0) { - KStream[] toBranches = outboundBindTarget.branch(predicates); - String[] topics = getOutputTopicsInProperOrder(name, branches); - for (int i = 0; i < toBranches.length; i++) { - if (!isNativeEncoding) { - toBranches[i].map(keyValueMapper).to(topics[i], Produced.with(keySerde, valueSerde)); - } - else { - toBranches[i].to(topics[i], Produced.with(keySerde, valueSerde)); - } - } - } else { - if (!isNativeEncoding) { + if (!isNativeEncoding) { outboundBindTarget.map(keyValueMapper).to(name, Produced.with(keySerde, valueSerde)); } - else { - outboundBindTarget.to(name, Produced.with(keySerde, valueSerde)); - } + else { + outboundBindTarget.to(name, Produced.with(keySerde, valueSerde)); } } - private static String[] getOutputTopicsInProperOrder(String name, String[] branches) { - String[] topics = new String[branches.length + 1]; - topics[0] = name; - int j = 1; - for (String branch : branches) { - topics[j++] = branch; - } - return topics; - } - @Override public KStreamConsumerProperties getExtendedConsumerProperties(String channelName) { return this.kStreamExtendedBindingProperties.getExtendedConsumerProperties(channelName); @@ -218,9 +175,4 @@ public class KStreamBinder extends public KStreamProducerProperties getExtendedProducerProperties(String channelName) { return this.kStreamExtendedBindingProperties.getExtendedProducerProperties(channelName); } - - public void setPredicates(Predicate[] predicates) { - this.predicates = predicates; - } - } diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamListenerSetupMethodOrchestrator.java new file mode 100644 index 000000000..eb2a68388 --- /dev/null +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamListenerSetupMethodOrchestrator.java @@ -0,0 +1,183 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.kstream; + +import java.lang.reflect.Method; +import java.util.Collection; + +import org.apache.kafka.streams.kstream.KStream; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.BeanInitializationException; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binding.StreamListenerErrorMessages; +import org.springframework.cloud.stream.binding.StreamListenerParameterAdapter; +import org.springframework.cloud.stream.binding.StreamListenerResultAdapter; +import org.springframework.cloud.stream.binding.StreamListenerSetupMethodOrchestrator; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.core.MethodParameter; +import org.springframework.core.annotation.AnnotationUtils; +import org.springframework.messaging.handler.annotation.SendTo; +import org.springframework.util.Assert; +import org.springframework.util.ObjectUtils; +import org.springframework.util.StringUtils; + +/** + * Kafka Streams specific implementation for {@link StreamListenerSetupMethodOrchestrator} + * that overrides the default mechanisms for invoking StreamListener adapters. + * + * @author Soby Chacko + */ +public class KStreamListenerSetupMethodOrchestrator implements StreamListenerSetupMethodOrchestrator, ApplicationContextAware { + + private ConfigurableApplicationContext applicationContext; + + private StreamListenerParameterAdapter streamListenerParameterAdapter; + private Collection streamListenerResultAdapters; + + public KStreamListenerSetupMethodOrchestrator(StreamListenerParameterAdapter streamListenerParameterAdapter, + Collection streamListenerResultAdapters) { + this.streamListenerParameterAdapter = streamListenerParameterAdapter; + this.streamListenerResultAdapters = streamListenerResultAdapters; + } + + @Override + public boolean supports(Method method) { + return methodParameterSuppports(method) && methodReturnTypeSuppports(method); + } + + private boolean methodReturnTypeSuppports(Method method) { + Class returnType = method.getReturnType(); + if (returnType.equals(KStream.class) || + (returnType.isArray() && returnType.getComponentType().equals(KStream.class))) { + return true; + } + return false; + } + + private boolean methodParameterSuppports(Method method) { + MethodParameter methodParameter = MethodParameter.forExecutable(method, 0); + Class parameterType = methodParameter.getParameterType(); + return parameterType.equals(KStream.class); + } + + @Override + @SuppressWarnings({"rawtypes", "unchecked"}) + public void orchestrateStreamListenerSetupMethod(StreamListener streamListener, Method method, Object bean) { + String[] methodAnnotatedOutboundNames = getOutboundBindingTargetNames(method); + validateStreamListenerMethod(streamListener, method, methodAnnotatedOutboundNames); + + String methodAnnotatedInboundName = streamListener.value(); + Object[] adaptedInboundArguments = adaptAndRetrieveInboundArguments(method, methodAnnotatedInboundName, + this.applicationContext, + this.streamListenerParameterAdapter); + + try { + Object result = method.invoke(bean, adaptedInboundArguments); + + if (result.getClass().isArray()) { + Assert.isTrue(methodAnnotatedOutboundNames.length == ((Object[]) result).length, "Big error"); + } else { + Assert.isTrue(methodAnnotatedOutboundNames.length == 1, "Big error"); + } + if (result.getClass().isArray()) { + Object[] outboundKStreams = (Object[]) result; + int i = 0; + for (Object outboundKStream : outboundKStreams) { + Object targetBean = this.applicationContext.getBean(methodAnnotatedOutboundNames[i++]); + for (StreamListenerResultAdapter streamListenerResultAdapter : streamListenerResultAdapters) { + if (streamListenerResultAdapter.supports(outboundKStream.getClass(), targetBean.getClass())) { + streamListenerResultAdapter.adapt(outboundKStream, targetBean); + break; + } + } + } + } + else { + Object targetBean = this.applicationContext.getBean(methodAnnotatedOutboundNames[0]); + for (StreamListenerResultAdapter streamListenerResultAdapter : streamListenerResultAdapters) { + if (streamListenerResultAdapter.supports(result.getClass(), targetBean.getClass())) { + streamListenerResultAdapter.adapt(result, targetBean); + break; + } + } + } + } + catch (Exception e) { + throw new BeanInitializationException("Cannot setup StreamListener for " + method, e); + } + } + + @Override + public final void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + this.applicationContext = (ConfigurableApplicationContext) applicationContext; + } + + private void validateStreamListenerMethod(StreamListener streamListener, Method method, String[] methodAnnotatedOutboundNames) { + String methodAnnotatedInboundName = streamListener.value(); + for (String s : methodAnnotatedOutboundNames) { + if (StringUtils.hasText(s)) { + Assert.isTrue(isDeclarativeOutput(method, s), "Method must be declarative"); + } + } + if (StringUtils.hasText(methodAnnotatedInboundName)) { + int methodArgumentsLength = method.getParameterTypes().length; + + for (int parameterIndex = 0; parameterIndex < methodArgumentsLength; parameterIndex++) { + MethodParameter methodParameter = MethodParameter.forExecutable(method, parameterIndex); + Assert.isTrue(isDeclarativeInput(methodAnnotatedInboundName, methodParameter), "Method must be declarative"); + } + } + } + + @SuppressWarnings("unchecked") + private boolean isDeclarativeOutput(Method m, String targetBeanName) { + boolean declarative; + Class returnType = m.getReturnType(); + if (returnType.isArray()){ + Class targetBeanClass = this.applicationContext.getType(targetBeanName); + declarative = this.streamListenerResultAdapters.stream() + .anyMatch(slpa -> slpa.supports(returnType.getComponentType(), targetBeanClass)); + return declarative; + } + Class targetBeanClass = this.applicationContext.getType(targetBeanName); + declarative = this.streamListenerResultAdapters.stream() + .anyMatch(slpa -> slpa.supports(returnType, targetBeanClass)); + return declarative; + } + + @SuppressWarnings("unchecked") + private boolean isDeclarativeInput(String targetBeanName, MethodParameter methodParameter) { + if (!methodParameter.getParameterType().isAssignableFrom(Object.class) && this.applicationContext.containsBean(targetBeanName)) { + Class targetBeanClass = this.applicationContext.getType(targetBeanName); + return this.streamListenerParameterAdapter.supports(targetBeanClass, methodParameter); + } + return false; + } + + private static String[] getOutboundBindingTargetNames(Method method) { + SendTo sendTo = AnnotationUtils.findAnnotation(method, SendTo.class); + if (sendTo != null) { + Assert.isTrue(!ObjectUtils.isEmpty(sendTo.value()), StreamListenerErrorMessages.ATLEAST_ONE_OUTPUT); + Assert.isTrue(sendTo.value().length >= 1, "At least one outbound destination need to be provided."); + return sendTo.value(); + } + return null; + } +} diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java index 4db42d63c..a0ab4c5d6 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderConfiguration.java @@ -19,7 +19,6 @@ package org.springframework.cloud.stream.binder.kstream.config; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.streams.StreamsConfig; -import org.apache.kafka.streams.kstream.Predicate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; @@ -45,9 +44,6 @@ public class KStreamBinderConfiguration { @Autowired private KafkaProperties kafkaProperties; - @Autowired(required = false) - private Predicate[] predicates; - @Bean public KafkaTopicProvisioner provisioningProvider(KafkaBinderConfigurationProperties binderConfigurationProperties) { return new KafkaTopicProvisioner(binderConfigurationProperties, kafkaProperties); @@ -60,9 +56,6 @@ public class KStreamBinderConfiguration { MessageConversionDelegate messageConversionDelegate) { KStreamBinder kStreamBinder = new KStreamBinder(binderConfigurationProperties, kafkaTopicProvisioner, kStreamExtendedBindingProperties, streamsConfig, messageConversionDelegate); - if (predicates != null) { - kStreamBinder.setPredicates(predicates); - } return kStreamBinder; } diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderSupportAutoConfiguration.java index beb545c9f..3a4826c50 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamBinderSupportAutoConfiguration.java @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kstream.config; +import java.util.Collection; import java.util.Properties; import org.apache.kafka.common.serialization.Serdes; @@ -29,8 +30,10 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kstream.KStreamBoundElementFactory; import org.springframework.cloud.stream.binder.kstream.KStreamListenerParameterAdapter; +import org.springframework.cloud.stream.binder.kstream.KStreamListenerSetupMethodOrchestrator; import org.springframework.cloud.stream.binder.kstream.KStreamStreamListenerResultAdapter; import org.springframework.cloud.stream.binder.kstream.MessageConversionDelegate; +import org.springframework.cloud.stream.binding.StreamListenerResultAdapter; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.context.annotation.Bean; @@ -90,6 +93,13 @@ public class KStreamBinderSupportAutoConfiguration { return new KStreamListenerParameterAdapter(messageConversionDelegate); } + @Bean + public KStreamListenerSetupMethodOrchestrator kStreamListenerSetupMethodOrchestrator( + KStreamListenerParameterAdapter kafkaStreamListenerParameterAdapter, + Collection streamListenerResultAdapters){ + return new KStreamListenerSetupMethodOrchestrator(kafkaStreamListenerParameterAdapter, streamListenerResultAdapters); + } + @Bean public MessageConversionDelegate messageConversionDelegate(BindingServiceProperties bindingServiceProperties, CompositeMessageConverterFactory compositeMessageConverterFactory) { diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamProducerProperties.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamProducerProperties.java index 051dc43ad..4946ff850 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamProducerProperties.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/config/KStreamProducerProperties.java @@ -22,13 +22,4 @@ package org.springframework.cloud.stream.binder.kstream.config; */ public class KStreamProducerProperties extends KStreamCommonProperties { - private String additionalBranches; - - public String getAdditionalBranches() { - return additionalBranches; - } - - public void setAdditionalBranches(String additionalBranches) { - this.additionalBranches = additionalBranches; - } } diff --git a/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/WordCountMultipleBranchesIntegrationTests.java b/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/WordCountMultipleBranchesIntegrationTests.java index a33a6752a..397e04bce 100644 --- a/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/WordCountMultipleBranchesIntegrationTests.java +++ b/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/WordCountMultipleBranchesIntegrationTests.java @@ -38,11 +38,11 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.Input; +import org.springframework.cloud.stream.annotation.Output; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.kstream.annotations.KStreamProcessor; import org.springframework.cloud.stream.binder.kstream.config.KStreamApplicationSupportProperties; import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -85,12 +85,15 @@ public class WordCountMultipleBranchesIntegrationTests { ConfigurableApplicationContext context = app.run("--server.port=0", "--spring.cloud.stream.bindings.input.destination=words", - "--spring.cloud.stream.bindings.output.destination=counts", - "--spring.cloud.stream.bindings.output.contentType=application/json", + "--spring.cloud.stream.bindings.output1.destination=counts", + "--spring.cloud.stream.bindings.output1.contentType=application/json", + "--spring.cloud.stream.bindings.output2.destination=foo", + "--spring.cloud.stream.bindings.output2.contentType=application/json", + "--spring.cloud.stream.bindings.output3.destination=bar", + "--spring.cloud.stream.bindings.output3.contentType=application/json", "--spring.cloud.stream.kstream.binder.configuration.commit.interval.ms=1000", "--spring.cloud.stream.kstream.binder.configuration.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde", "--spring.cloud.stream.kstream.binder.configuration.value.serde=org.apache.kafka.common.serialization.Serdes$StringSerde", - "--spring.cloud.stream.kstream.bindings.output.producer.additionalBranches=foo,bar", "--spring.cloud.stream.bindings.output.producer.headerMode=raw", "--spring.cloud.stream.bindings.input.consumer.headerMode=raw", "--spring.cloud.stream.kstream.timeWindow.length=5000", @@ -125,7 +128,7 @@ public class WordCountMultipleBranchesIntegrationTests { assertThat(cr.value().contains("\"word\":\"spanish\",\"count\":3")).isTrue(); } - @EnableBinding(KStreamProcessor.class) + @EnableBinding(KStreamProcessorX.class) @EnableAutoConfiguration @EnableConfigurationProperties(KStreamApplicationSupportProperties.class) public static class WordCountProcessorApplication { @@ -134,25 +137,38 @@ public class WordCountMultipleBranchesIntegrationTests { private TimeWindows timeWindows; @StreamListener("input") - @SendTo("output") - public KStream process(KStream input) { + @SendTo({"output1","output2","output3"}) + @SuppressWarnings("unchecked") + public KStream[] process(KStream input) { + + Predicate isEnglish = (k, v) -> v.word.equals("english"); + Predicate isFrench = (k, v) -> v.word.equals("french"); + Predicate isSpanish = (k, v) -> v.word.equals("spanish"); + return input .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) .groupBy((key, value) -> value) .windowedBy(timeWindows) .count(Materialized.as("WordCounts-multi")) .toStream() - .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))); + .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))) + .branch(isEnglish, isFrench, isSpanish); } + } - @Bean - public Predicate[] predicates() { - Predicate isEnglish = (k, v) -> v.word.equals("english"); - Predicate isFrench = (k, v) -> v.word.equals("french"); - Predicate isSpanish = (k, v) -> v.word.equals("spanish"); - return new Predicate[] {isEnglish, isFrench, isSpanish}; - } + interface KStreamProcessorX { + @Input("input") + KStream input(); + + @Output("output1") + KStream output1(); + + @Output("output2") + KStream output2(); + + @Output("output3") + KStream output3(); } static class WordCount {