Designing Data-Intensive Applications
Ch. 11

Processing Streams

Stream processors apply transformations to unbounded data — windowing time is the hard part.

Stream processing applies computations to data in motion. Unlike batch, the input is unbounded. Key challenges: handling out-of-order events, defining time windows, and maintaining state across failures.

In practice

Flink powers real-time fraud detection at banks — joining a payment stream against a Redis feature store within millisecond windows. Kafka Streams embeds processing in JVM microservices without a separate cluster. ksqlDB lets analysts write SQL over Kafka topics. Materialize maintains incrementally updated SQL views — like a live data warehouse fed by streams.

Uber at scale

Surge pricing aggregates ride demand in 5-minute tumbling windows over a geo-partitioned event stream — event time, not processing time, defines when a window closes.

typescript — Flink tumbling window aggregation
// Uber surge pricing — Flink tumbling window on ride stream
stream
  .keyBy((event) => event.geoHash)
  .window(TumblingEventTimeWindows.of(Time.minutes(5)))
  .aggregate({
    add: (acc, ride) => ({ count: acc.count + 1, fare: acc.fare + ride.fare }),
    getResult: (acc) => acc,
  })
  .filter((stats) => stats.count > demandThreshold)
  .map((stats) => ({ geoHash: stats.key, surgeMultiplier: 1.5 }));
Key Takeaways
  • Uses: notifications, search indexing, metrics, fraud detection, recommendations.
  • Event time (when it happened) vs processing time (when observed) diverge.
  • Windowing groups events by time ranges for aggregation.
  • Stream joins combine two streams or a stream with a table.
  • Fault tolerance via checkpointing and replay from durable logs.
  • Apache Flink, Kafka Streams, ksqlDB, and Materialize process streams in production.
FlinkKafka StreamsksqlDBMaterializestream processingwindowingcheckpoint