Skip to content

Commit 74ecee4

Browse files
KStreamBranch -- ProcessorSupplier of KStreamBranchProcessors
1 parent 381c5c0 commit 74ecee4

4 files changed

+56
-8
lines changed

kafka-streams-AbstractProcessor.adoc

+3-3
Original file line numberDiff line numberDiff line change
@@ -15,8 +15,8 @@
1515
| KStreamAggregateProcessor
1616
| [[KStreamAggregateProcessor]]
1717

18-
| KStreamBranchProcessor
19-
| [[KStreamBranchProcessor]]
18+
| <<kafka-streams-internals-KStreamBranch.adoc#KStreamBranchProcessor, KStreamBranchProcessor>>
19+
| [[KStreamBranchProcessor]] Forwards a record to exactly one of the child processors
2020

2121
| <<kafka-streams-internals-KStreamFilter.adoc#KStreamFilterProcessor, KStreamFilterProcessor>>
2222
| [[KStreamFilterProcessor]]
@@ -46,7 +46,7 @@
4646
| [[KStreamPassThroughProcessor]]
4747

4848
| <<kafka-streams-internals-KStreamPeek.adoc#KStreamPeekProcessor, KStreamPeekProcessor>>
49-
| [[KStreamPeekProcessor]]
49+
| [[KStreamPeekProcessor]] Executes an action with records (_peeks at them_)
5050

5151
| KStreamPrintProcessor
5252
| [[KStreamPrintProcessor]]

kafka-streams-ProcessorSupplier.adoc

+2-2
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ Processor<K, V> get()
2424
| [[KStreamAggProcessorSupplier]]
2525

2626
| <<kafka-streams-internals-KStreamBranch.adoc#, KStreamBranch>>
27-
| [[KStreamBranch]]
27+
| [[KStreamBranch]] Represents <<kafka-streams-KStream.adoc#branch, KStream.branch>> operator
2828

2929
| <<kafka-streams-internals-KStreamFilter.adoc#, KStreamFilter>>
3030
| [[KStreamFilter]]
@@ -57,7 +57,7 @@ Processor<K, V> get()
5757
| [[KStreamPassThrough]]
5858

5959
| <<kafka-streams-internals-KStreamPeek.adoc#, KStreamPeek>>
60-
| [[KStreamPeek]] <<kafka-streams-internals-KStreamPeek.adoc#KStreamPeekProcessor, KStreamPeekProcessors>> for <<kafka-streams-KStream.adoc#foreach, KStream.foreach>> and <<kafka-streams-KStream.adoc#peek, KStream.peek>> operators
60+
| [[KStreamPeek]] Represents <<kafka-streams-KStream.adoc#foreach, KStream.foreach>> and <<kafka-streams-KStream.adoc#peek, KStream.peek>> operators
6161

6262
| KStreamPrint
6363
| [[KStreamPrint]]
+50-2
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,51 @@
1-
== [[KStreamBranch]] KStreamBranch
1+
== [[KStreamBranch]] KStreamBranch -- ProcessorSupplier of KStreamBranchProcessors
22

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.

kafka-streams-internals-KStreamPeek.adoc

+1-1
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ When <<kafka-streams-ProcessorSupplier.adoc#get, requested for a Processor>>, `K
4343

4444
=== [[KStreamPeekProcessor]] KStreamPeekProcessor
4545

46-
`KStreamPeekProcessor` is a custom <<kafka-streams-Processor.adoc#, record processor>> (indirectly as <<kafka-streams-AbstractProcessor.adoc#, AbstractProcessor>>) that allows for <<process, executing an action with records>> (_peek at them_).
46+
`KStreamPeekProcessor` is a custom <<kafka-streams-Processor.adoc#, record processor>> (indirectly as <<kafka-streams-AbstractProcessor.adoc#, AbstractProcessor>>) that allows for <<process, executing an action with records>> (_peeks at them_).
4747

4848
[[process]]
4949
When requested to <<kafka-streams-Processor.adoc#process, process a record>>, `KStreamPeekProcessor` executes the <<action, ForeachAction>> with the record. If the <<forwardDownStream, forwardDownStream>> flag is enabled, `KStreamPeekProcessor` requests the `ProcessorContext` to <<kafka-streams-ProcessorContext.adoc#forward, forward the record downstreams>>.

0 commit comments

Comments
 (0)