|
1 |
| -== [[KStreamBranch]] KStreamBranch |
| 1 | +== [[KStreamBranch]] KStreamBranch -- ProcessorSupplier of KStreamBranchProcessors |
2 | 2 |
|
3 |
| -`KStreamBranch` is...FIXME |
| 3 | +`KStreamBranch` is a custom <<kafka-streams-ProcessorSupplier.adoc#, ProcessorSupplier>> of <<KStreamBranchProcessor, KStreamBranchProcessors>> for <<kafka-streams-KStream.adoc#branch, KStream.branch>> operator. |
| 4 | + |
| 5 | +[source, scala] |
| 6 | +---- |
| 7 | +// Scala API for Kafka Streams |
| 8 | +import org.apache.kafka.streams.scala._ |
| 9 | +import ImplicitConversions._ |
| 10 | +import Serdes._ |
| 11 | +
|
| 12 | +val builder = new StreamsBuilder |
| 13 | +def alwaysTrue(k: String, v: String) = true |
| 14 | +builder |
| 15 | + .stream[String, String]("input") |
| 16 | + .branch(alwaysTrue) |
| 17 | +val topology = builder.build |
| 18 | +scala> println(topology.describe) |
| 19 | +Topologies: |
| 20 | + Sub-topology: 0 |
| 21 | + Source: KSTREAM-SOURCE-0000000000 (topics: [input]) |
| 22 | + --> KSTREAM-BRANCH-0000000001 |
| 23 | + Processor: KSTREAM-BRANCH-0000000001 (stores: []) |
| 24 | + --> KSTREAM-BRANCHCHILD-0000000002 |
| 25 | + <-- KSTREAM-SOURCE-0000000000 |
| 26 | + Processor: KSTREAM-BRANCHCHILD-0000000002 (stores: []) |
| 27 | + --> none |
| 28 | + <-- KSTREAM-BRANCH-0000000001 |
| 29 | +---- |
| 30 | + |
| 31 | +[[creating-instance]] |
| 32 | +`KStreamBranch` takes the following to be created: |
| 33 | + |
| 34 | +* [[predicates]] `Predicate<K, V>[]` |
| 35 | +* [[childNodes]] Child nodes |
| 36 | +
|
| 37 | +`KStreamBranch` is <<creating-instance, created>> exclusively when `KStreamImpl` is requested to <<kafka-streams-internals-KStreamImpl.adoc#branch, branch>>. |
| 38 | + |
| 39 | +NOTE: `KStreamImpl` is the default <<kafka-streams-KStream.adoc#, KStream>>. |
| 40 | + |
| 41 | +[[get]] |
| 42 | +When <<kafka-streams-ProcessorSupplier.adoc#get, requested for a Processor>>, `KStreamBranch` gives a new <<KStreamBranchProcessor, KStreamBranchProcessor>>. |
| 43 | + |
| 44 | +=== [[KStreamBranchProcessor]] KStreamBranchProcessor |
| 45 | + |
| 46 | +`KStreamBranchProcessor` is a custom <<kafka-streams-Processor.adoc#, record processor>> (indirectly as <<kafka-streams-AbstractProcessor.adoc#, AbstractProcessor>>) that allows for <<process, forwarding a record to exactly one of the child processors>> (_branching on them_). |
| 47 | + |
| 48 | +[[process]] |
| 49 | +When requested to <<kafka-streams-Processor.adoc#process, process a record>>, `KStreamBranchProcessor` walks over the <<predicates, predicates>> and requests each and every predicate to test the record. When `true`, `process` requests the `ProcessorContext` to <<kafka-streams-ProcessorContext.adoc#forward, forward the record>> to a corresponding <<kafka-streams-To.adoc#child, child processor>>. |
| 50 | + |
| 51 | +NOTE: `process` requests the predicates until positive is found or finishes without forwarding a record. |
0 commit comments