Streaming Data
Stream processing handles data continuously as it arrives, rather than in fixed batches. This enables real-time analytics, fraud detection, live dashboards, and event-driven applications. The core tools in this space are Apache Kafka, Apache Flink, and Spark Structured Streaming.
Apache Kafka
Kafka is the de facto standard for building event streaming platforms. It acts as a distributed commit log, providing durable, scalable, and fault-tolerant message storage.
Key Concepts
- Topic - A named channel for related messages. Producers write to topics; consumers read from them.
- Partition - Each topic splits into partitions for parallelism. Messages within a partition are ordered.
- Consumer Group - A group of consumers that coordinate to read partitions in parallel.
- Offset - A unique integer identifying a message's position within a partition. Consumers track offsets for checkpointing.
- Broker - A Kafka server that stores and serves topic partitions.
Example: Producing and Consuming
# Producer
from kafka import KafkaProducer
import json
producer = KafkaProducer(
bootstrap_servers=["kafka:9092"],
value_serializer=lambda v: json.dumps(v).encode(),
)
producer.send(
"user-events",
{"user_id": 42, "event": "page_view", "url": "/pricing"},
)
producer.flush()# Consumer
from kafka import KafkaConsumer
import json
consumer = KafkaConsumer(
"user-events",
bootstrap_servers=["kafka:9092"],
group_id="analytics-consumer",
value_deserializer=lambda v: json.loads(v.decode()),
)
for message in consumer:
event = message.value
print(f"Processing: {event}")Apache Flink
Flink is a stream processing framework designed for stateful computations over unbounded data streams. It supports event-time processing, exactly-once semantics, and complex windowing.
Flink Windowing Example
-- Flink SQL: tumbling window aggregation CREATE TABLE page_views ( user_id INT, page_url STRING, view_time TIMESTAMP(3), WATERMARK FOR view_time AS view_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'page_views', 'properties.bootstrap.servers' = 'kafka:9092', 'format' = 'json' ); SELECT TUMBLE_START(view_time, INTERVAL '1' MINUTE) AS window_start, page_url, COUNT(*) AS views FROM page_views GROUP BY TUMBLE(view_time, INTERVAL '1' MINUTE), page_url;Spark Structured Streaming
Spark Structured Streaming provides a high-level API for stream processing built on the Spark SQL engine. It treats streams as unbounded tables, making it easy to express streaming queries with the same DataFrame API.
Spark Streaming Example
from pyspark.sql import SparkSession from pyspark.sql.functions import window, col spark = SparkSession.builder.appName("streaming").getOrCreate() events = ( spark .readStream .format("kafka") .option("kafka.bootstrap.servers", "kafka:9092") .option("subscribe", "clicks") .load() .selectExpr("CAST(value AS STRING) as json") .selectExpr("from_json(json, 'user_id INT, action STRING, ts TIMESTAMP') as data") .select("data.*") ) aggregated = ( events .groupBy(window(col("ts"), "5 minutes"), col("action")) .count() ) query = ( aggregated .writeStream .outputMode("complete") .format("console") .start() ) query.awaitTermination()Exactly-Once Semantics
Stream processing systems offer three delivery guarantees:
- At-most-once - Each message is processed zero or one times. Fast but can drop messages.
- At-least-once - Each message is processed one or more times. Duplicates possible but no drops.
- Exactly-once - Each message is processed exactly once. Achieved through transactional writes and idempotent sinks.
Flink and Kafka Streams both support end-to-end exactly-once semantics via two-phase commit protocols.
Event-Driven Architectures
Streaming is the backbone of event-driven architecture (EDA), where services communicate through events rather than direct API calls. Key patterns include:
- Event sourcing - State is derived from a log of events
- CQRS - Separate read and write models, synchronized via events
- Event notification - Services publish events others subscribe to