diff --git a/docs/src/main/asciidoc/kafka-streams.adoc b/docs/src/main/asciidoc/kafka-streams.adoc index 8ba285d92..5fc0bc419 100644 --- a/docs/src/main/asciidoc/kafka-streams.adoc +++ b/docs/src/main/asciidoc/kafka-streams.adoc @@ -277,7 +277,7 @@ public Function, KStream> bar() { ===== Multiple Output Bindings -Kafka Streams allows to write outbound data into multiple topics. This feature is known as branching in Kafka Streams. +Kafka Streams allows writing outbound data into multiple topics. This feature is known as branching in Kafka Streams. When using multiple output bindings, you need to provide an array of KStream (`KStream[]`) as the outbound return type. Here is an example: @@ -291,21 +291,30 @@ public Function, KStream[]> process() { Predicate isFrench = (k, v) -> v.word.equals("french"); Predicate isSpanish = (k, v) -> v.word.equals("spanish"); - return input -> input - .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) - .groupBy((key, value) -> value) - .windowedBy(TimeWindows.of(5000)) - .count(Materialized.as("WordCounts-branch")) - .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); + return input -> { + final Map> stringKStreamMap = input + .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) + .groupBy((key, value) -> value) + .windowedBy(TimeWindows.of(Duration.ofSeconds(5))) + .count(Materialized.as("WordCounts-branch")) + .toStream() + .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, + new Date(key.window().start()), new Date(key.window().end())))) + .split() + .branch(isEnglish) + .branch(isFrench) + .branch(isSpanish) + .noDefaultBranch(); + + return stringKStreamMap.values().toArray(new KStream[0]); + }; } ---- The programming model remains the same, however the outbound parameterized type is `KStream[]`. -The default output binding names are `process-out-0`, `process-out-1`, `process-out-2` respectively. -The reason why the binder generates three output bindings is because it detects the length of the returned `KStream` array. +The default output binding names are `process-out-0`, `process-out-1`, `process-out-2` respectively for the function above. +The reason why the binder generates three output bindings is because it detects the length of the returned `KStream` array as three. +Note that in this example, we provide a `noDefaultBranch()`; if we have used `defaultBranch()` instead, that would have required an extra output binding, essentially returning a `KStream` array of length four. ===== Summary of Function based Programming Styles for Kafka Streams 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 8897a14af..1f1d443ef 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 @@ -1,5 +1,5 @@ /* - * Copyright 2019-2019 the original author or authors. + * Copyright 2019-2021 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. @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams.function; +import java.time.Duration; import java.util.Arrays; import java.util.Date; import java.util.Map; @@ -48,6 +49,9 @@ import org.springframework.kafka.test.utils.KafkaTestUtils; import static org.assertj.core.api.Assertions.assertThat; +/** + * @author Soby Chacko + */ public class KafkaStreamsBinderWordCountBranchesFunctionTests { @ClassRule @@ -179,22 +183,30 @@ public class KafkaStreamsBinderWordCountBranchesFunctionTests { public static class WordCountProcessorApplication { @Bean - @SuppressWarnings("unchecked") + @SuppressWarnings({"unchecked"}) public Function, KStream[]> process() { 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 -> input - .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) - .groupBy((key, value) -> value) - .windowedBy(TimeWindows.of(5000)) - .count(Materialized.as("WordCounts-branch")) - .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); + return input -> { + final Map> stringKStreamMap = input + .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) + .groupBy((key, value) -> value) + .windowedBy(TimeWindows.of(Duration.ofSeconds(5))) + .count(Materialized.as("WordCounts-branch")) + .toStream() + .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, + new Date(key.window().start()), new Date(key.window().end())))) + .split() + .branch(isEnglish) + .branch(isFrench) + .branch(isSpanish) + .noDefaultBranch(); + + return stringKStreamMap.values().toArray(new KStream[0]); + }; } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/WordCountMultipleBranchesIntegrationTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/WordCountMultipleBranchesIntegrationTests.java deleted file mode 100644 index baa38547d..000000000 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/WordCountMultipleBranchesIntegrationTests.java +++ /dev/null @@ -1,236 +0,0 @@ -/* - * Copyright 2017-2019 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 - * - * https://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.kafka.streams.integration; - -import java.time.Duration; -import java.util.Arrays; -import java.util.Date; -import java.util.Map; - -import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.consumer.ConsumerConfig; -import org.apache.kafka.clients.consumer.ConsumerRecord; -import org.apache.kafka.streams.KeyValue; -import org.apache.kafka.streams.kstream.KStream; -import org.apache.kafka.streams.kstream.Materialized; -import org.apache.kafka.streams.kstream.Predicate; -import org.apache.kafka.streams.kstream.TimeWindows; -import org.junit.AfterClass; -import org.junit.BeforeClass; -import org.junit.ClassRule; -import org.junit.Test; - -import org.springframework.boot.SpringApplication; -import org.springframework.boot.WebApplicationType; -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -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.context.ConfigurableApplicationContext; -import org.springframework.kafka.core.DefaultKafkaConsumerFactory; -import org.springframework.kafka.core.DefaultKafkaProducerFactory; -import org.springframework.kafka.core.KafkaTemplate; -import org.springframework.kafka.test.EmbeddedKafkaBroker; -import org.springframework.kafka.test.rule.EmbeddedKafkaRule; -import org.springframework.kafka.test.utils.KafkaTestUtils; -import org.springframework.messaging.handler.annotation.SendTo; - -import static org.assertj.core.api.Assertions.assertThat; - -/** - * @author Marius Bogoevici - * @author Soby Chacko - * @author Gary Russell - */ -public class WordCountMultipleBranchesIntegrationTests { - - @ClassRule - public static EmbeddedKafkaRule embeddedKafkaRule = new EmbeddedKafkaRule(1, true, - "counts", "foo", "bar"); - - private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule - .getEmbeddedKafka(); - - private static Consumer consumer; - - @BeforeClass - public static void setUp() throws Exception { - Map consumerProps = KafkaTestUtils.consumerProps("groupx", - "false", embeddedKafka); - consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>( - consumerProps); - consumer = cf.createConsumer(); - embeddedKafka.consumeFromEmbeddedTopics(consumer, "counts", "foo", "bar"); - } - - @AfterClass - public static void tearDown() { - consumer.close(); - } - - @Test - public void testKstreamWordCountWithStringInputAndPojoOuput() throws Exception { - SpringApplication app = new SpringApplication( - WordCountProcessorApplication.class); - app.setWebApplicationType(WebApplicationType.NONE); - - ConfigurableApplicationContext context = app.run("--server.port=0", - "--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input.destination=words", - "--spring.cloud.stream.bindings.output1.destination=counts", - "--spring.cloud.stream.bindings.output2.destination=foo", - "--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", - "--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" - + "=WordCountMultipleBranchesIntegrationTests-abc", - "--spring.cloud.stream.kafka.streams.binder.brokers=" - + embeddedKafka.getBrokersAsString()); - try { - receiveAndValidate(context); - } - finally { - context.close(); - } - } - - private void receiveAndValidate(ConfigurableApplicationContext context) - throws Exception { - Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); - DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>( - senderProps); - KafkaTemplate template = new KafkaTemplate<>(pf, true); - template.setDefaultTopic("words"); - template.sendDefault("english"); - ConsumerRecord cr = KafkaTestUtils.getSingleRecord(consumer, - "counts"); - assertThat(cr.value().contains("\"word\":\"english\",\"count\":1")).isTrue(); - - template.sendDefault("french"); - template.sendDefault("french"); - cr = KafkaTestUtils.getSingleRecord(consumer, "foo"); - assertThat(cr.value().contains("\"word\":\"french\",\"count\":2")).isTrue(); - - template.sendDefault("spanish"); - template.sendDefault("spanish"); - template.sendDefault("spanish"); - cr = KafkaTestUtils.getSingleRecord(consumer, "bar"); - assertThat(cr.value().contains("\"word\":\"spanish\",\"count\":3")).isTrue(); - } - - @EnableBinding(KStreamProcessorX.class) - @EnableAutoConfiguration - public static class WordCountProcessorApplication { - - @StreamListener("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.of(Duration.ofSeconds(5))) - .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())))) - .branch(isEnglish, isFrench, isSpanish); - } - - } - - interface KStreamProcessorX { - - @Input("input") - KStream input(); - - @Output("output1") - KStream output1(); - - @Output("output2") - KStream output2(); - - @Output("output3") - KStream output3(); - - } - - static class WordCount { - - private String word; - - private long count; - - private Date start; - - private Date end; - - WordCount(String word, long count, Date start, Date end) { - this.word = word; - this.count = count; - this.start = start; - this.end = end; - } - - public String getWord() { - return word; - } - - public void setWord(String word) { - this.word = word; - } - - public long getCount() { - return count; - } - - public void setCount(long count) { - this.count = count; - } - - public Date getStart() { - return start; - } - - public void setStart(Date start) { - this.start = start; - } - - public Date getEnd() { - return end; - } - - public void setEnd(Date end) { - this.end = end; - } - - } - -}