Table of Contents

    stream processing

    MESSAGING & ASYNCHRONOUS 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-processing Flow
    ingest events → validate schema → assign keys → transform or aggregate → manage state → publish results

    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

    1. Identify the source topics and affected topology.
    2. Compare event-arrival and processing-completion rates.
    3. Inspect lag for every source partition.
    4. Check key and partition distribution.
    5. Identify the slow or failing operator.
    6. Check state size and state-retention policy.
    7. Check watermark and late-event behaviour.
    8. Check checkpoint duration and failure history.
    9. Check retries, poison events, and DLQ volume.
    10. Check sink latency, throttling, and availability.
    11. Check whether replay creates duplicate external effects.
    12. Scale or reshard only when useful downstream capacity exists.

    Common Stream-processing Mistakes

    1

    Using Processing Time When Event Time Is Required

    Delayed events are assigned to the wrong business-time calculation.

    2

    Ignoring Late Events

    Previously emitted aggregates remain permanently incomplete.

    3

    Allowing Unbounded State

    Keys, join entries, or window data accumulate until memory or storage is exhausted.

    4

    Choosing a Hot Key

    One processing instance becomes overloaded while other instances remain underused.

    5

    Assuming Global Ordering

    Separate partitions progress independently and do not provide one total sequence.

    6

    Ignoring Replay and Duplicate Processing

    Recovery repeats external database, API, or notification effects.

    7

    Calling Slow APIs without Concurrency Limits

    Pending asynchronous calls consume memory and create pipeline backpressure.

    8

    Using Incompatible Schemas

    Consumer deserialization fails and valid later events can be delayed.

    9

    Checkpointing without Monitoring

    Recovery assumptions fail because checkpoints are slow, failing, or unavailable.

    10

    Keeping State Forever

    Inactive keys remain stored despite having no continuing business value.

    11

    Scaling without Checking the Sink

    Additional processors overload the destination database or API.

    12

    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

    1. Create an input topic for learning events.
    2. Use learner ID and course ID as the aggregate key.
    3. Validate every event against a versioned schema.
    4. Filter events unrelated to progress calculation.
    5. Maintain completed-lesson count as keyed state.
    6. Calculate the latest progress percentage.
    7. Create a five-minute course-activity window.
    8. Define event-time and watermark behaviour.
    9. Define a bounded late-event policy.
    10. Publish current progress to an output topic.
    11. Write the progress projection idempotently to a database.
    12. Enable checkpointing.
    13. Simulate processor failure and restore state.
    14. Simulate an out-of-order event.
    15. Simulate one hot learner or course key.
    16. 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

    1

    What is stream processing?

    Stream processing continuously consumes and transforms events while the data is being produced.

    2

    What is a data stream?

    A data stream is an ongoing sequence of events that can continue without a predefined end.

    3

    What is stateless processing?

    Stateless processing handles each event without remembering information from earlier events.

    4

    What is stateful processing?

    Stateful processing maintains information across events for aggregations, joins, windows, pattern detection, and current business state.

    5

    What is event time?

    Event time is the time at which the business event originally occurred.

    6

    What is a window?

    A window divides an unbounded stream into bounded time or session groups for calculation.

    7

    What is a watermark?

    A watermark represents progress through event time and helps determine when event-time windows can produce results.

    8

    What is a late event?

    A late event arrives after the processor has advanced beyond its expected event-time position.

    9

    What is keyed state?

    Keyed state maintains separate processing state for every business key.

    10

    What is checkpointing?

    Checkpointing records processing positions and operator state so a failed application can recover consistently.

    11

    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.

    12

    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.