stream processing
Stream Processing
Learn how stream processing continuously transforms events as they arrive. Understand stateless and stateful operations, filtering, mapping, aggregation, joins, windows, event time, processing time, watermarks, late events, checkpoints, delivery semantics, partitioning, backpressure, schema evolution, Kafka Streams, Apache Flink, and reliable stream-processing architecture.
Introduction
Modern applications continuously generate events such as transactions, clicks, sensor readings, enrollments, progress updates, inventory changes, and application logs.
Business applications
|
v
Continuous event stream
|
v
Stream-processing application
|
+-- Validate
+-- Filter
+-- Transform
+-- Aggregate
+-- Join
+-- Detect patterns
|
v
Output streams and systems
Traditional batch processing collects a finite dataset and processes it at a scheduled time. Stream processing works with continuously arriving events and produces results while the data remains in motion.
Core idea: Stream processing continuously consumes events, applies processing logic, maintains any required state, and produces updated results without waiting for the complete dataset to finish.
Prerequisites
| # | Prerequisite | Why It Is Needed |
|---|---|---|
| 1 | Queues vs logs | Stream-processing applications commonly consume retained event logs. |
| 2 | Kafka, RabbitMQ, and Amazon SQS | Kafka is commonly used as an event-streaming foundation, while queues serve different asynchronous workloads. |
| 3 | Consumer groups | Partitions are distributed among processing instances for parallelism. |
| 4 | Delivery semantics | Failures and recovery can cause records to be replayed. |
| 5 | Ordering | Ordering guarantees are normally scoped to partitions or business keys. |
| 6 | Backpressure | Processing pipelines must remain stable when event arrival exceeds completion capacity. |
| 7 | Batching | Stream systems can use record-at-a-time processing, micro-batches, or bounded groups of records. |
What Is Stream Processing?
Stream processing is a computational model in which an application continuously consumes and processes an ongoing sequence of events.
Event 1 -> Process -> Result
Event 2 -> Process -> Result
Event 3 -> Process -> Result
...
The stream continues.
A stream is often treated as an unbounded dataset because new events can continue arriving while the application remains active.
Typical processing operations include:
- Filtering unwanted events
- Transforming event structure
- Enriching events with reference data
- Grouping events by business key
- Calculating running aggregates
- Joining related streams
- Detecting event patterns
- Writing projections and alerts
Stream vs Batch Processing
| Area | Stream Processing | Batch Processing |
|---|---|---|
| Data model | Continuously arriving events | Finite collection of records |
| Trigger | Event arrival or continuous execution | Schedule, request, or dataset availability |
| Result timing | Results update continuously | Results appear after batch execution |
| Latency direction | Designed for lower processing delay | Usually accepts delayed results |
| State | Can remain active and evolve across events | Can be recomputed from the finite input |
| Failure recovery | Restore state and resume or replay input | Restart the complete job or a failed partition |
| Example | Continuously update learner-progress totals | Generate the previous month's completion report |
Stream-processing Topology
A topology describes how events move through processing operators.
Source
|
v
Validate
|
v
Filter
|
v
Key by learner
|
v
Aggregate progress
|
v
Enrich with course data
|
v
Output topic
|
v
Search, analytics, or notification system
A topology can branch into several outputs or merge several inputs.
Stateless Processing
A stateless operation processes one event using only the information available in that event.
Incoming event
|
v
Inspect event fields
|
v
Produce result
No earlier event information
is required.
Common stateless operations include:
- Filtering
- Field selection
- Data-format conversion
- Validation
- Masking or redaction
- Routing
- Simple one-event calculations
Stateless Transformation
{
"eventType": "LessonCompleted",
"tenantId": 17,
"learnerId": 1042,
"courseId": 42,
"lessonId": 8,
"durationSeconds": 720
}
A stateless processor can convert the duration from seconds to minutes without consulting previous events.
Stateful Processing
A stateful operation remembers information across several events.
LessonCompleted event
|
v
Read current learner state
|
v
Update completed-lesson count
|
v
Store current state
|
v
Emit updated progress
Stateful operations include:
- Counts and totals
- Running averages
- Windowed aggregations
- Stream joins
- Deduplication
- Session tracking
- Pattern detection
- Business state machines
State rule: State must be partitioned, persisted, recoverable, observable, and limited by an explicit lifecycle. Unbounded state can eventually exhaust memory or storage.
Keyed State
Keyed state stores independent state for each business key.
Key:
learner-1042-course-42
State:
completedLessons = 8
totalLessons = 10
progress = 80%
Events with the same key should reach the operator instance responsible for that key.
Different keys can be processed in parallel.
Processor A:
Learners 1, 4, 7
Processor B:
Learners 2, 5, 8
Processor C:
Learners 3, 6, 9
Aggregation
Aggregation combines several events into a continuously updated result.
A running sum can be represented as:
\[ Total_n = Total_{n-1} + Value_n \]
Incoming lesson durations:
10 minutes
15 minutes
8 minutes
Running total:
10
25
33 minutes
The processor must preserve the current total as state between events.
Stream Transformations
| Operation | Purpose |
|---|---|
| Filter | Keep only records that match a condition |
| Map | Convert one event into another form |
| Flat map | Convert one event into zero, one, or several results |
| Group by key | Bring related events into the same processing scope |
| Aggregate | Maintain counts, sums, averages, or other derived state |
| Join | Combine related records from streams or tables |
| Window | Limit unbounded input to a defined time or session scope |
| Branch | Send records to different processing paths |
Stream Joins
A join combines related information from separate streams or from a stream and a current-state table.
Enrollment Stream:
learnerId
courseId
enrolledAt
Course Table:
courseId
courseTitle
category
Joined Output:
learnerId
courseId
courseTitle
category
enrolledAt
A stream join normally requires compatible keys and a bounded matching policy. Without a time or state-retention boundary, unmatched join state can grow indefinitely.
Processing Time and Event Time
| Time Concept | Meaning |
|---|---|
| Event time | When the business event actually occurred |
| Ingestion time | When the event entered the streaming platform |
| Processing time | When a processor handled the event |
Event occurred:
10:00
Mobile device was offline.
Event entered stream:
10:10
Processor handled event:
10:11
Processing-time calculations place this event near 10:11. Event-time calculations place it near 10:00.
Windows
A window divides an unbounded stream into bounded groups for calculation.
Tumbling Window
09:00 to 09:05
09:05 to 09:10
09:10 to 09:15
Each event belongs to
one non-overlapping window.
Sliding Window
Window length:
10 minutes
Window advances:
Every 5 minutes
One event can appear
in several overlapping windows.
Session Window
Learner activity begins
|
v
Events continue
|
v
No activity for configured gap
|
v
Session closes
| Window Type | Typical Use |
|---|---|
| Tumbling | Events per fixed reporting interval |
| Sliding | Moving average or recent activity rate |
| Session | User activity grouped by inactivity gaps |
| Global | Custom application-controlled grouping |
Watermarks
Event-time processing must decide when enough earlier events have probably arrived to calculate a window result.
A watermark represents processing progress through event time.
Current watermark:
10:05
Interpretation:
The processor assumes that
most events earlier than 10:05
have already arrived.
A watermark does not guarantee that no older event will arrive later. Applications still require a late-event policy.
Late Events
A late event arrives after the processor has already advanced beyond the event's expected event-time position.
Window:
10:00 to 10:05
Window result emitted:
100 completions
Late event arrives:
Event time = 10:03
Possible policies include:
- Update the previous result
- Emit a correction event
- Allow lateness within a bounded period
- Send very late events to a side output
- Ignore events that are no longer useful
- Reconcile later through batch processing
Late-data rule: Define late-event behaviour as part of the business contract. Silently dropping late data can produce permanently incomplete results.
Partitioning and Parallelism
Stream-processing systems partition events by a key so that several operator instances can process the stream in parallel.
Event key
|
v
Partitioning function
|
+-- Partition 0 -> Processor A
+-- Partition 1 -> Processor B
+-- Partition 2 -> Processor C
Events requiring per-entity order should use a stable key that places them in the same ordered processing scope.
A low-cardinality or highly popular key can create a hot partition and limit parallelism.
Kafka Streams
Kafka provides built-in stream-processing capabilities for operations such as filters, transformations, joins, aggregations, and event-time processing. Kafka Streams applications read Kafka topics, execute a topology, and write results to Kafka topics or other supported outputs.
Kafka input topic
|
v
Kafka Streams application
|
+-- Filter
+-- Group
+-- Aggregate
+-- Join
|
v
Kafka output topic
Kafka Streams Example
StreamsBuilder builder =
new StreamsBuilder();
KStream<String, ProgressEvent> events =
builder.stream(
"learner-progress-events"
);
KTable<String, Long> completedLessons =
events
.filter(
(key, event) ->
event.isLessonCompleted()
)
.groupByKey()
.count();
completedLessons
.toStream()
.to(
"learner-progress-counts"
);
Production code requires serializers, schemas, error handling, topology testing, security, state-store configuration, delivery-semantics configuration, and observability.
Apache Flink
Apache Flink is a distributed processing engine for stateful computations over bounded and unbounded streams. Flink supports stateful operators, event-time processing, late-data handling, checkpoints, savepoints, and scale-out execution.
Kafka source
|
v
Flink processing topology
|
+-- Keyed state
+-- Event-time windows
+-- Checkpoints
+-- Async enrichment
|
v
Kafka, database,
object store, or API sink
Internal architecture material describes a Flink workflow that consumes Kafka requests, validates events, invokes asynchronous processing, handles timeouts, constructs result messages, and publishes those results through a Kafka sink. 【1-fcfa49】
Kafka Streams vs Apache Flink
| Area | Kafka Streams | Apache Flink |
|---|---|---|
| Form | Stream-processing library used inside an application | Distributed stream-processing framework and engine |
| Primary ecosystem | Kafka topics and Kafka applications | Multiple supported streaming sources and sinks |
| State | Local state stores backed by Kafka topics according to topology design | Managed operator and keyed state with configurable state backends |
| Fault tolerance | Topic replay and state restoration through Kafka integration | Checkpoints, savepoints, and input replay |
| Typical fit | Kafka-centered event-processing applications | Complex stateful pipelines, event-time analysis, and large managed dataflows |
Checkpoints
A checkpoint records a consistent processing position and the associated operator state.
Input positions
+
Operator state
+
Pending processing context
|
v
Consistent checkpoint
After a failure, the processor can restore the checkpointed state and replay records from the corresponding input positions.
Processor failure
|
v
Restore latest valid checkpoint
|
v
Restore operator state
|
v
Replay later input records
|
v
Resume processing
Apache Flink documentation states that Flink uses stream replay and checkpointing for fault-tolerant stateful processing and can restore operator state together with input positions. 【2-f9592b】
Savepoints
A savepoint is an operationally triggered state snapshot that can support controlled application upgrades, migrations, or restarts where the processing framework supports it.
Running stream job
|
v
Create savepoint
|
v
Stop or upgrade job
|
v
Restore from savepoint
Processing Guarantees
| Semantic | General Behaviour |
|---|---|
| At-most-once | Failures can cause records to be skipped. |
| At-least-once | Required records are retried, but processing can be repeated. |
| Exactly-once state processing | State and stream progress are coordinated inside the supported processing boundary. |
| Effectively-once external outcome | Repeated delivery is absorbed through idempotency, deduplication, or conditional writes. |
Exactly-once processing inside a stream engine does not automatically include an external payment provider, email system, or database unless the sink participates in the relevant consistency protocol.
Deduplication
Stream replay and at-least-once delivery can process one logical event more than once.
{
"eventId": "stable-unique-event-id",
"eventType": "LessonCompleted",
"aggregateId": "learner-1042-course-42",
"sourceVersion": 8
}
Use stable event IDs, sequence numbers, compact deduplication state, database constraints, or idempotent sink operations.
Schema Evolution
Stream producers and consumers are often deployed independently, so message schemas must evolve without breaking active processors.
Use:
- Explicit event type and schema version
- Backward-compatible optional fields
- Stable field meanings
- Schema compatibility checks
- Consumer handling for unknown optional fields
- Dead-letter handling for unsupported records
Internal architecture material describes a schema registry used to manage Avro, JSON Schema, and Protocol Buffers contracts for streaming applications. 【3-0ca42f】
Backpressure
Stream-processing pipelines become pressured when input arrives faster than the topology and its sinks can complete processing.
Source rate increases
|
v
Operator cannot keep pace
|
v
Input lag and buffered work grow
|
v
Backpressure slows upstream flow
or limits additional intake
Useful controls include:
- Bounded operator buffers
- Consumer lag monitoring
- Partition-aware scaling
- Async-operation concurrency limits
- Sink connection limits
- Producer throttling
- Retry budgets
- Poison-event isolation
Conceptual Stream-processing Policy
streamProcessing:
source:
topic: learning-domain-events
schemaValidation: required
partitioning:
key: trusted-aggregate-id
orderingScope: per-aggregate
processing:
mode: continuous
stateful: true
idempotent: true
time:
primarySemantics: event-time
watermarkPolicy: approved-policy
allowedLateness: approved-boundary
state:
keyed: true
retention: approved-lifecycle
checkpointing: enabled
failures:
boundedRetries: true
retryBackoff: exponential-with-jitter
deadLetterDestination: approved-dlq
backpressure:
boundedBuffers: true
asyncConcurrency: approved-limit
observability:
sourceLag: enabled
watermarkLag: enabled
lateEvents: enabled
checkpointFailures: enabled
stateSize: enabled
This is a conceptual configuration. Exact properties depend on the stream engine, messaging platform, application topology, and business consistency requirements.
Learning-platform Examples
| Use Case | Stream Operation | State Required |
|---|---|---|
| Learner-progress calculation | Group and aggregate lesson events | Current progress per learner and course |
| Live course popularity | Windowed count of course views | Count per course and time window |
| Assessment monitoring | Filter and aggregate submission events | Current submission metrics |
| Certificate eligibility | Join progress, assessment, and completion streams | Eligibility state per learner-course aggregate |
| Notification trigger | Detect qualifying event pattern | Previously observed lifecycle events |
| Search projection | Transform and materialize course events | Latest indexed version per course |
| Active learner sessions | Session windows | Events grouped until the inactivity gap |
Observability
Useful stream-processing metrics include:
- Records received per second
- Records completed per second
- Consumer lag by partition
- End-to-end processing latency
- Event-time and watermark delay
- Late-event count
- Out-of-order event count
- Operator processing duration
- State size by operator
- Checkpoint duration and failures
- Restart and recovery count
- Backpressure duration
- Retry and dead-letter volume
- Sink failures and throttling
- Hot keys and partition skew
Structured Processing Event
{
"topology": "learner-progress-processing",
"operator": "progress-aggregation",
"partition": "protected-partition-reference",
"eventTime": "event-time",
"processingResult": "success",
"checkpointReference": "protected-checkpoint-reference",
"schemaVersion": 1
}
Avoid including credentials, complete event payloads, or unnecessary personal information in stream-processing logs.
Alert Conditions
Alert when:
- Source lag continues increasing
- Processing completion rate remains below the input rate
- Watermarks stop advancing
- Late-event volume increases unexpectedly
- One partition or key becomes hot
- State size grows without a defined boundary
- Checkpoints fail repeatedly
- Checkpoint duration increases significantly
- A processor repeatedly restarts
- Retries or dead-letter events increase
- A sink becomes throttled or unavailable
- End-to-end event latency exceeds the processing objective
Troubleshooting Workflow
- Identify the source topics and affected topology.
- Compare event-arrival and processing-completion rates.
- Inspect lag for every source partition.
- Check key and partition distribution.
- Identify the slow or failing operator.
- Check state size and state-retention policy.
- Check watermark and late-event behaviour.
- Check checkpoint duration and failure history.
- Check retries, poison events, and DLQ volume.
- Check sink latency, throttling, and availability.
- Check whether replay creates duplicate external effects.
- Scale or reshard only when useful downstream capacity exists.
Common Stream-processing Mistakes
Using Processing Time When Event Time Is Required
Delayed events are assigned to the wrong business-time calculation.
Ignoring Late Events
Previously emitted aggregates remain permanently incomplete.
Allowing Unbounded State
Keys, join entries, or window data accumulate until memory or storage is exhausted.
Choosing a Hot Key
One processing instance becomes overloaded while other instances remain underused.
Assuming Global Ordering
Separate partitions progress independently and do not provide one total sequence.
Ignoring Replay and Duplicate Processing
Recovery repeats external database, API, or notification effects.
Calling Slow APIs without Concurrency Limits
Pending asynchronous calls consume memory and create pipeline backpressure.
Using Incompatible Schemas
Consumer deserialization fails and valid later events can be delayed.
Checkpointing without Monitoring
Recovery assumptions fail because checkpoints are slow, failing, or unavailable.
Keeping State Forever
Inactive keys remain stored despite having no continuing business value.
Scaling without Checking the Sink
Additional processors overload the destination database or API.
Claiming Exactly-Once without Naming the Boundary
The engine can coordinate state internally while an external side effect remains outside the guarantee.
Recommended Test Cases
| Test | Expected Evidence |
|---|---|
| Stateless transformation | Every valid input produces the expected transformed event. |
| Keyed aggregation | Each business key maintains independent correct state. |
| Window boundary | Events are assigned to the expected time window. |
| Late event | The configured correction, side-output, or rejection policy runs. |
| Out-of-order event | Version or sequence validation protects current state. |
| Processor failure | State and processing positions recover from the valid checkpoint. |
| Repeated event | Idempotency prevents a duplicate business effect. |
| Hot key | Partition and operator metrics identify the skew. |
| Sink slowdown | Backpressure and concurrency limits protect the destination. |
| Invalid schema | The record follows the approved failure path. |
| State retention | Expired inactive state is removed according to policy. |
| Application rescaling | Keyed state moves consistently to the new parallel instances. |
Stream-processing Best Practices
Recommended Practices
- Define whether the workload requires batch or continuous processing.
- Separate stateless and stateful operators.
- Use stable high-cardinality business keys.
- Preserve required ordering within each key.
- Use event time when business time matters.
- Define watermark and late-event policies.
- Bound windows, joins, deduplication, and keyed state.
- Define state-retention and cleanup rules.
- Enable and monitor fault-tolerant checkpoints.
- Verify checkpoint restoration regularly.
- Make sinks idempotent or transactional.
- Use stable event IDs and source versions.
- Use compatible and versioned message schemas.
- Apply bounded retries with backoff and jitter.
- Move poison records to a controlled DLQ.
- Limit asynchronous external calls.
- Monitor lag, event-time delay, state size, and backpressure.
- Scale only within partition and downstream capacity.
- Name the exact boundary of processing guarantees.
- Test failure, replay, late data, rescaling, and recovery.
Practice Exercise
Design a real-time learner-progress processing pipeline for your online learning platform.
Requirements
- Create an input topic for learning events.
- Use learner ID and course ID as the aggregate key.
- Validate every event against a versioned schema.
- Filter events unrelated to progress calculation.
- Maintain completed-lesson count as keyed state.
- Calculate the latest progress percentage.
- Create a five-minute course-activity window.
- Define event-time and watermark behaviour.
- Define a bounded late-event policy.
- Publish current progress to an output topic.
- Write the progress projection idempotently to a database.
- Enable checkpointing.
- Simulate processor failure and restore state.
- Simulate an out-of-order event.
- Simulate one hot learner or course key.
- Monitor lag, state size, late events, and checkpoint health.
Stream-design Template
| Decision | Selected Direction | Risk Controlled |
|---|---|---|
| Source | Partitioned durable event topic | Supports scalable consumption and replay |
| Key | Stable learner-course aggregate ID | Preserves related state and ordering |
| Time semantics | Event time with bounded lateness | Handles delayed event delivery |
| State | Keyed progress state with retention | Prevents unbounded state growth |
| Recovery | Checkpoint and input replay | Restores consistent state after failure |
| Sink | Version-aware idempotent upsert | Prevents duplicate and stale updates |
| Failure handling | Bounded retry and DLQ | Prevents poison events from blocking the pipeline |
| Monitoring | Lag, watermarks, state, checkpoints, and sink health | Detects pipeline delay and instability |
Frequently Asked Questions
What is stream processing?
Stream processing continuously consumes and transforms events while the data is being produced.
What is a data stream?
A data stream is an ongoing sequence of events that can continue without a predefined end.
What is stateless processing?
Stateless processing handles each event without remembering information from earlier events.
What is stateful processing?
Stateful processing maintains information across events for aggregations, joins, windows, pattern detection, and current business state.
What is event time?
Event time is the time at which the business event originally occurred.
What is a window?
A window divides an unbounded stream into bounded time or session groups for calculation.
What is a watermark?
A watermark represents progress through event time and helps determine when event-time windows can produce results.
What is a late event?
A late event arrives after the processor has advanced beyond its expected event-time position.
What is keyed state?
Keyed state maintains separate processing state for every business key.
What is checkpointing?
Checkpointing records processing positions and operator state so a failed application can recover consistently.
How is Kafka Streams different from Kafka?
Kafka stores and distributes event streams. Kafka Streams is a library for building applications that process records from Kafka topics.
What is the safest stream-processing design?
Use stable keys, versioned schemas, bounded state, event-time policies, fault-tolerant checkpoints, idempotent sinks, bounded retries, backpressure, DLQs, and complete observability.
Key Takeaway
Stream processing continuously transforms events as they arrive instead of waiting for a complete scheduled batch. Stateless operators handle each event independently, while stateful operators maintain counts, windows, joins, sessions, deduplication data, and business state across events. Partition streams using stable business keys so related events and state remain together while unrelated keys progress in parallel. Use event time when the time of the business event matters, and define watermarks, windows, allowed lateness, and correction behaviour explicitly. Protect state using retention, checkpoints, savepoints, and tested restoration. Make output sinks idempotent because replay can repeat processing. Use compatible schemas, bounded retries, DLQs, backpressure, and concurrency limits. Finally, monitor source lag, event-time delay, late events, state size, checkpoints, hot keys, processing failures, and sink health to ensure that the pipeline remains correct, recoverable, and scalable.