14. Stream Processing
14.8 What actually breaks in production — Ch. 14 consolidated
| # | Symptom | Root cause | Fix |
|---|---|---|---|
| 1 | Aggregates reset to zero after a restart | State kept in local variables — "when the application is stopped or crashes, the state is lost, WHICH CHANGES THE RESULTS" | Use a proper state store (Kafka Streams Materialized + changelog topic) |
| 2 | Windowed results are wrong after a producer outage | A producer returned with hours of backlog for already-closed windows | Configure a grace period; use event time, not processing time |
| 3 | The same event produces different results in different runs | Used processing time — "it can even differ for TWO THREADS IN THE SAME APPLICATION" | Event time via TimestampExtractor |
| 4 | Timestamps are wrong for CDC-sourced events | Kafka's auto timestamp is record-creation time, not the original event time | "Add the event time as a FIELD IN THE RECORD ITSELF so BOTH timestamps are available" |
| 5 | Window results are nonsensical across regions | Mixed time zones | Standardize the whole pipeline on one time zone; store the zone in the record if you must mix |
| 6 | Aggregation is slow, spilling to disk | Local state exceeds available memory | Partition into smaller substreams; "spilling to disk has SIGNIFICANT PERFORMANCE IMPACT" |
| 7 | Database overwhelmed by enrichment lookups | Per-record external lookup: 5–15 ms each, and "stream systems handle 100K–500K events/s but the DB only ~10K/s" | Stream-table join over a CDC-maintained local copy |
| 8 | Enrichment uses stale data | Hand-rolled cache with a refresh interval — "too often → hammering the DB; too long → stale" | CDC-driven cache updates |
| 9 | Local aggregation can't compute a global result (top 10) | "all the top 10 stocks could be in partitions assigned to OTHER instances" | Multiphase: local aggregate → single-partition summary topic → single instance |
| 10 | A Kafka Streams join refuses to run / produces nothing | "Kafka Streams REQUIRES that all topics that participate in a join have THE SAME NUMBER OF PARTITIONS and be PARTITIONED BASED ON THE JOIN KEY" | Repartition or recreate topics to match |
| 11 | A stream-stream join matches unrelated events | Window too wide, or symmetric when causality is one-directional | JoinWindows.of(1s).before(0s) — express causality as window asymmetry |
| 12 | Late events silently dropped | No grace period configured | Define the reconciliation period; remember "the longer the windows stay available, the MORE MEMORY is required" |
| 13 | Late corrections never reach consumers | Output topic not compacted | "Those are USUALLY COMPACTED TOPICS" — a new result for the same key replaces the old |
| 14 | Reprocessing corrupted results / lost data | Reset offsets and state in place | "OUR RECOMMENDATION IS TO USE THE FIRST METHOD" — run a second app as a new consumer group and switch clients over |
| 15 | groupByKey() doesn't group anything | It doesn't — "despite its name, this operation DOES NOT DO ANY GROUPING. Rather, it ensures the stream is PARTITIONED based on the record key" | Understand it as a repartition assertion |
| 16 | Serialization failure on an intermediate result | Provided Serdes for input/output but not for the aggregation result object | "Remember to provide a Serde for EVERY object you want to store in Kafka — input, output, AND, IN SOME CASES, INTERMEDIATE RESULTS" |
| 17 | Windowed output can't be deserialized | Windowed Serde missing the window size — "deserialization REQUIRES the window size, because ONLY THE START TIME is stored" | WindowedSerdes.timeWindowedSerdeFrom(class, windowSize) |
| 18 | Two Kafka Streams apps collide on internal topics/stores | Duplicate APPLICATION_ID_CONFIG — it "names the internal local stores AND THE TOPICS related to them" | "THIS NAME MUST BE UNIQUE for each Kafka Streams application working with the same Kafka cluster" |
| 19 | Enabled topology optimization; nothing changed | Called build() without the props | "If you only call build() without passing the config, OPTIMIZATION IS STILL DISABLED" |
| 20 | Optimization changed the results | Optimizations rewrite the physical plan | "TEST with and without... and VALIDATE THAT THE RESULTS ARE IDENTICAL in various known scenarios" |
| 21 | Unit tests pass but production fails | TopologyTestDriver "does NOT simulate Kafka Streams CACHING behavior... there are ENTIRE CLASSES OF ERRORS IT WILL NOT DETECT" | Add integration tests — Testcontainers preferred |
| 22 | Adding threads doesn't increase throughput | Task count == partition count; you've saturated it | More partitions (mind Ch. 3 §9.4); more instances only helps up to partition count |
| 23 | After an instance fails, a subset of keys goes stale for minutes | Task must replay the changelog to warm the state store | segment.bytes = 100 MB (not 1 GB) + low min.compaction.lag.ms; and standby replicas |
| 24 | Recovery replays far more than expected | The active segment is never compacted — with 1 GB segments, up to 1 GB/partition is uncompacted | Same fix as #23 |
| 25 | Changelog topics grow without bound | Compaction not working on internal topics (Ch. 13 §6: silently halted cleaner threads) | Kafka uses compaction "to make sure they don't grow endlessly and that re-creating the state is ALWAYS FEASIBLE" — verify it's actually running |
| 26 | Duplicate contributions to aggregates after a failure | At-least-once processing | processing.guarantee=exactly_once (or exactly_once_beta on 2.6+/2.5+ brokers) |
| 27 | Rebalances pause the whole app | Eager rebalancing | Kafka Streams inherits cooperative rebalancing and static group membership (Ch. 4) |
| 28 | Built a stream processing app for pure data movement | Wrong tool | "Reconsider whether you want... a SIMPLER INGEST-FOCUSED SYSTEM LIKE KAFKA CONNECT" |
| 29 | Sub-millisecond SLA not met by a streaming app | Wrong paradigm | "REQUEST-RESPONSE PATTERNS ARE OFTEN BETTER SUITED"; if you must stream, avoid microbatch frameworks |