Update Kafka Streams branching docs/tests
Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1133
This commit is contained in:
@@ -277,7 +277,7 @@ public Function<KTable<String, String>, KStream<String, String>> 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<Object, String>, KStream<?, WordCount>[]> process() {
|
||||
Predicate<Object, WordCount> isFrench = (k, v) -> v.word.equals("french");
|
||||
Predicate<Object, WordCount> 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<String, KStream<Object, WordCount>> 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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user