Stream processing is a data processing paradigm in which computations are performed continuously on unbounded sequences of records as they arrive, rather than on static stored datasets, enabling low-latency analytics, transformations, and reactions to events within milliseconds to seconds of their occurrence. It is characterised by windowing operations, stateful operators, time-based semantics (event time versus processing time), and exactly-once or at-least-once delivery guarantees.
Content
- The theoretical foundations of stream processing trace to the Stanford STREAM project and MIT Aurora/Borealis work in the early 2000s, which formalised windowing semantics and continuous query languages. Commercial systems such as Esper (complex event processing), IBM InfoSphere Streams, and Oracle CEP followed. Apache Storm (2011, open-sourced from BackType/Twitter) was the first widely adopted open-source distributed stream processor, followed by Apache Samza (LinkedIn, 2013), Apache Spark Streaming (2012), and Apache Flink (2014).
- Modern stream processing frameworks distinguish event time (when an event occurred) from processing time (when it is processed), handling late-arriving data through watermarks and allowed lateness policies. Flink’s stateful stream processing model provides exactly-once semantics via distributed snapshots (Chandy-Lamport inspired checkpointing), enabling strong consistency guarantees at scale. The Kafka Streams and ksqlDB libraries bring stream processing directly into the Kafka ecosystem without a separate cluster.
- Industry use cases span financial services (real-time transaction monitoring, order book processing), telecommunications (network anomaly detection), e-commerce (personalisation, inventory updates), and operational intelligence (infrastructure metrics, log analysis). Confluent, Amazon Kinesis, Google Dataflow, and Azure Event Hubs provide managed cloud stream processing, reducing operational overhead. At scale, major deployments process millions of events per second with sub-100-millisecond end-to-end latency.
- By 2024–2025, the distinction between batch and stream processing has blurred further with Apache Flink’s unified batch-streaming engine and Apache Spark’s Structured Streaming adopting micro-batch approaches. GenAI workloads are driving new streaming patterns for real-time RAG, agent event processing, and model output streaming. RisingWave, Materialize, and other streaming SQL databases are gaining adoption by lowering the operational complexity of building streaming applications.