skip to content

How do you route a single KStream into multiple branches in modern Kafka Streams, and how do you recombine streams? Contrast split/branch with merge.

level: middleimportance: should knowfreq 45%

answer

  1. split().branch(p, Branched.as()).defaultBranch()
  2. First match wins; unmatched dropped sans default
  3. Returns Map<String,KStream>, keys = prefix+name
  4. Old branch()->KStream[] deprecated
  5. merge = interleave, same types, no cross-order

basics

~10 s

Use KStream.split() then chained .branch(predicate, Branched.as("name")) and optionally .defaultBranch(), which returns a Map<String, KStream> of routed sub-streams. merge(otherStream) does the opposite: it interleaves two same-typed streams into one. Both are stateless.

solid answer

~40 s

The modern routing API is `stream.split([Named.as(prefix)]).branch(predicate, Branched.as("X")).branch(...).defaultBranch()`, returning a `Map<String, KStream<K,V>>` keyed by the branch names (prefix + the Branched name). Predicates are evaluated top-to-bottom and a record goes to the FIRST matching branch only; unmatched records are dropped unless you add `defaultBranch()`. The old `KStream.branch(Predicate...)` returning a `KStream[]` array is deprecated in favor of split(). merge is the inverse: `streamA.merge(streamB)` produces one stream containing all records from both, with the same key/value types, interleaved with no ordering guarantee across the two inputs. All of split/branch/merge are stateless and key-preserving, so they never set the repartition flag. Branched can also take a Consumer/Function to handle a branch inline (e.g. send to a topic) instead of returning it in the map.

go deeper

for a junior

Know split/branch routes a stream and merge combines two streams.

for a middle

Use the modern split()/Branched API, know first-match and default-branch semantics, and merge's same-type requirement.

for a senior

Explain Branched.withConsumer inline handling, the returned-map naming, and that these are stateless/key-preserving.

for a principal

Weigh branch vs multiple filters for fan-out, and reason about ordering guarantees after merge in downstream stateful ops.

## Splitting one stream into many (branch) You often want to fan a stream into several downstream paths by content — e.g. route orders by region, or errors vs. successes. **Modern API (Kafka 2.8+):** ``` Map<String, KStream<K,V>> branches = stream.split(Named.as("route-")) .branch((k,v) -> v.region == US, Branched.as("us")) .branch((k,v) -> v.region == EU, Branched.as("eu")) .defaultBranch(Branched.as("other")); // access: branches.get("route-us"), branches.get("route-eu"), branches.get("route-other") ``` Key semantics: - Predicates are tested **in order**; a record lands in the **first** branch whose predicate returns true. It is NOT copied to multiple branches. - A record matching **no** predicate is **dropped**, unless you provide `defaultBranch()`. - The returned **Map** keys are `<prefix><branch-name>` (the `Named.as` prefix concatenated with the `Branched.as` name). - `Branched.withConsumer(ks -> ks.to("topic"))` or `Branched.withFunction(...)` lets you process a branch inline; such branches are NOT added to the returned map. **Deprecated API:** `KStream<K,V>[] arr = stream.branch(p1, p2, p3);` returned an array indexed positionally. Deprecated because positional indexing was error-prone and unnamed; prefer split()/branch(). ## If you want a record in MULTIPLE outputs branch is exclusive (first match wins). To send the same record down several independent paths, just apply multiple separate operators to the same stream (filter it several times), or use `filter` per path — don't expect branch to duplicate. ## Recombining: merge `KStream<K,V> merged = streamA.merge(streamB);` produces one stream with **all** records from both inputs. Requirements/behavior: - Both streams must share the same key and value types. - Records are interleaved; there is **no ordering guarantee** between the two sources (order within each source's partition is preserved). - merge is stateless and key-preserving — no repartition, no state store. ## Statelessness split/branch/merge all process each record independently and keep the key, so none of them sets the repartition-required flag or needs a changelog. They're pure topology-shaping operators.

  • A record matches two of your branch predicates. Which branch gets it?
    Only the FIRST one in declaration order. branch is exclusive — predicates are evaluated top-to-bottom and the record goes to the first match, never to multiple branches.
  • What happens to records that match none of the predicates?
    They are dropped, unless you add a defaultBranch() to capture them. There is no implicit catch-all.

saying these in an interview costs you the question

  • Saying branch copies a matching record into every matching branch (it's first-match-only)
  • Assuming unmatched records pass through (they're dropped without defaultBranch)
  • Claiming merge guarantees global ordering across the two streams
  • Trying to merge streams of different key/value types

context