Learn Labs
14. Stream Processing

14.9 Consolidated reference

The design patterns and when each applies

PatternState?Windowed?When
Single-event (map/filter)nonenoper-event transform/route
Local state aggregationLOCALusuallyGROUP-BY aggregates; keys partition cleanly
Multiphase / repartitionLOCAL×2maybeGLOBAL results (top-N)
Stream-table joinLOCAL (CDC-fed)NOENRICH events from a table (fact ⋈ dimension)
Table-table joinLOCAL×2NOcurrent state ⋈ current state; equi- or foreign-key
Stream-stream joinLOCALYEScorrelate two real streams within a time window
Out-of-sequence handlingLOCALYESIoT, mobile, flaky networks
Reprocessing——new version, or bug fix
Interactive queriesLOCAL—read the state store direct

Kafka Streams operational checklist

Correctness

  • APPLICATION_ID unique 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 _beta on 2.6+/2.5+)

Performance / recovery

  • segment.bytes ≈ 100 MB and low min.compaction.lag.ms on all Streams internal topics ← recovery speed
  • num.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_OPTIMIZATION tested (with build(props)) and results validated

Testing

  • TopologyTestDriver unit 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-partitions balanced across instances
  • commit-latency-avg baselined

How this chapter closes the loop on the whole book

Ch. 1log, retention, compactionCh. 3keyed partitioningCh. 4consumer group ownershipCh. 6compacted changelog topicsCh. 7reliability / offset disciplineCh. 8transactionsLOCAL STATE IS CORRECTowns its keys · state is replayableCh. 14 — Kafka Streamssix existing properties, composed

“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
Figure 14.9.3How this chapter closes the loop on the whole book

On this page