14. Stream Processing
14.9 Consolidated reference
The design patterns and when each applies
| Pattern | State? | Windowed? | When |
|---|---|---|---|
| Single-event (map/filter) | none | no | per-event transform/route |
| Local state aggregation | LOCAL | usually | GROUP-BY aggregates; keys partition cleanly |
| Multiphase / repartition | LOCAL×2 | maybe | GLOBAL results (top-N) |
| Stream-table join | LOCAL (CDC-fed) | NO | ENRICH events from a table (fact ⋈ dimension) |
| Table-table join | LOCAL×2 | NO | current state ⋈ current state; equi- or foreign-key |
| Stream-stream join | LOCAL | YES | correlate two real streams within a time window |
| Out-of-sequence handling | LOCAL | YES | IoT, mobile, flaky networks |
| Reprocessing | — | — | new version, or bug fix |
| Interactive queries | LOCAL | — | read the state store direct |
Kafka Streams operational checklist
Correctness
APPLICATION_IDunique per app per cluster- Serdes provided for input, output, and intermediate objects
- Event time used (
TimestampExtractor); event time in the payload if the record-creation time isn't the real event time - One time zone across the whole pipeline
- Window size / advance interval / grace period all chosen deliberately
- Join topics: same partition count, partitioned by the join key
- Join windows express causality (
.before(0s)where appropriate) processing.guarantee = exactly_once(or_betaon 2.6+/2.5+)
Performance / recovery
segment.bytes≈ 100 MB and lowmin.compaction.lag.mson all Streams internal topics ← recovery speednum.standby.replicas> 0 for stateful apps ← recovery speed- Compaction verified as actually running (Ch. 13 §6 loggers)
- Threads per instance and instance count ≤ partition count
TOPOLOGY_OPTIMIZATIONtested (withbuild(props)) and results validated
Testing
TopologyTestDriverunit tests- Testcontainers integration tests (caching behavior is not covered by the test driver)
- Reprocessing plan = “run a second app as a new consumer group”, not “reset offsets and state”
Monitoring (Ch. 13)
- Consumer-group lag via Burrow (a Streams app is a consumer group)
sync-rate≈ 0 (rebalance storms)assigned-partitionsbalanced across instancescommit-latency-avgbaselined
How this chapter closes the loop on the whole book
“Kafka Streams is not a new system. It is these six existing Kafka properties composed into a programming model.”
- · the shuffle is a topic
- · the state store's durability is a compacted topic
- · task failover is consumer group rebalancing
- · exactly-once is Kafka transactions
- · the cluster is just more copies of your app