Learn Labs
14. Stream Processing

14.8 What actually breaks in production — Ch. 14 consolidated

Production failure catalog
0 rows
#SymptomRoot causeFix
1Aggregates reset to zero after a restartState 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)
2Windowed results are wrong after a producer outageA producer returned with hours of backlog for already-closed windowsConfigure a grace period; use event time, not processing time
3The same event produces different results in different runsUsed processing time — "it can even differ for TWO THREADS IN THE SAME APPLICATION"Event time via TimestampExtractor
4Timestamps are wrong for CDC-sourced eventsKafka'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"
5Window results are nonsensical across regionsMixed time zonesStandardize the whole pipeline on one time zone; store the zone in the record if you must mix
6Aggregation is slow, spilling to diskLocal state exceeds available memoryPartition into smaller substreams; "spilling to disk has SIGNIFICANT PERFORMANCE IMPACT"
7Database overwhelmed by enrichment lookupsPer-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
8Enrichment uses stale dataHand-rolled cache with a refresh interval — "too often → hammering the DB; too long → stale"CDC-driven cache updates
9Local 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
10A 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
11A stream-stream join matches unrelated eventsWindow too wide, or symmetric when causality is one-directionalJoinWindows.of(1s).before(0s) — express causality as window asymmetry
12Late events silently droppedNo grace period configuredDefine the reconciliation period; remember "the longer the windows stay available, the MORE MEMORY is required"
13Late corrections never reach consumersOutput topic not compacted"Those are USUALLY COMPACTED TOPICS" — a new result for the same key replaces the old
14Reprocessing corrupted results / lost dataReset 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
15groupByKey() doesn't group anythingIt 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
16Serialization failure on an intermediate resultProvided 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"
17Windowed output can't be deserializedWindowed Serde missing the window size — "deserialization REQUIRES the window size, because ONLY THE START TIME is stored"WindowedSerdes.timeWindowedSerdeFrom(class, windowSize)
18Two Kafka Streams apps collide on internal topics/storesDuplicate 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"
19Enabled topology optimization; nothing changedCalled build() without the props"If you only call build() without passing the config, OPTIMIZATION IS STILL DISABLED"
20Optimization changed the resultsOptimizations rewrite the physical plan"TEST with and without... and VALIDATE THAT THE RESULTS ARE IDENTICAL in various known scenarios"
21Unit tests pass but production failsTopologyTestDriver "does NOT simulate Kafka Streams CACHING behavior... there are ENTIRE CLASSES OF ERRORS IT WILL NOT DETECT"Add integration tests — Testcontainers preferred
22Adding threads doesn't increase throughputTask count == partition count; you've saturated itMore partitions (mind Ch. 3 §9.4); more instances only helps up to partition count
23After an instance fails, a subset of keys goes stale for minutesTask must replay the changelog to warm the state storesegment.bytes = 100 MB (not 1 GB) + low min.compaction.lag.ms; and standby replicas
24Recovery replays far more than expectedThe active segment is never compacted — with 1 GB segments, up to 1 GB/partition is uncompactedSame fix as #23
25Changelog topics grow without boundCompaction 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
26Duplicate contributions to aggregates after a failureAt-least-once processingprocessing.guarantee=exactly_once (or exactly_once_beta on 2.6+/2.5+ brokers)
27Rebalances pause the whole appEager rebalancingKafka Streams inherits cooperative rebalancing and static group membership (Ch. 4)
28Built a stream processing app for pure data movementWrong tool"Reconsider whether you want... a SIMPLER INGEST-FOCUSED SYSTEM LIKE KAFKA CONNECT"
29Sub-millisecond SLA not met by a streaming appWrong paradigm"REQUEST-RESPONSE PATTERNS ARE OFTEN BETTER SUITED"; if you must stream, avoid microbatch frameworks