Files of Batch and stream processing
wondelai/
Show the full text343 lines
Batch and Stream Processing
Data processing systems fall into three categories: services (online systems that handle requests), batch processing (offline systems that process large volumes of accumulated data), and stream processing (near-real-time systems that process data as it arrives). Understanding when and how to use batch and stream processing is essential for building data pipelines, analytics systems, and derived data stores.
Table of Contents
- Batch Processing: MapReduce and Beyond
- Dataflow Engines: Beyond MapReduce
- Event Sourcing
- Change Data Capture (CDC)
- Stream-Table Duality
- Exactly-Once Semantics
- Time Windowing
- Architecture Patterns
Batch Processing: MapReduce and Beyond
The MapReduce Paradigm
MapReduce, popularized by Google's 2004 paper, processes large datasets by breaking computation into two phases:
- Map phase: Read input data, extract key-value pairs. Each mapper processes a portion of the input independently.
- Shuffle phase: Framework groups all values by key, distributing them to reducers.
- Reduce phase: For each key, combine all values into a result.
Input: ["the cat sat on the mat"]
Map:
"the" -> 1
"cat" -> 1
"sat" -> 1
"on" -> 1
"the" -> 1
"mat" -> 1
Shuffle: Group by key
"the" -> [1, 1]
"cat" -> [1]
"sat" -> [1]
"on" -> [1]
"mat" -> [1]
Reduce: Sum values
"the" -> 2
"cat" -> 1
"sat" -> 1
"on" -> 1
"mat" -> 1
MapReduce Strengths
- Horizontal scalability: Add more machines to process larger datasets; the framework handles distribution
- Fault tolerance: If a mapper or reducer fails, the framework re-executes that task on another machine. Input data is immutable, so re-execution is safe.
- Simplicity: The programmer only writes map and reduce functions; the framework handles parallelism, distribution, and fault tolerance
MapReduce Weaknesses
- High latency: Each MapReduce job reads from and writes to distributed storage (HDFS), adding significant I/O overhead
- Chaining is awkward: Complex computations require chaining multiple MapReduce jobs, each with its own read/write cycle
- No iteration support: Machine learning algorithms that iterate over data must launch a new MapReduce job for each iteration
- Limited expressiveness: Not all computations fit the map-reduce pattern naturally
Dataflow Engines: Beyond MapReduce
Apache Spark
Spark replaces MapReduce with a more general computation model based on Resilient Distributed Datasets (RDDs) and directed acyclic graphs (DAGs) of operators.
Key improvements over MapReduce:
- In-memory processing: Intermediate results stay in memory instead of being written to disk between stages
- Arbitrary DAGs: Computations can have multiple stages with various operators (map, filter, join, group, sort), not just map and reduce
- Lazy evaluation: Spark builds a computation plan before executing, enabling the optimizer to eliminate unnecessary steps
- Iterative algorithms: RDDs can be cached in memory and reused across iterations
# Spark: Word count
text_file = spark.read.text("hdfs://input/")
word_counts = (text_file
.select(explode(split(col("value"), " ")).alias("word"))
.groupBy("word")
.count()
.orderBy(desc("count")))
word_counts.write.parquet("hdfs://output/")
Apache Flink
Flink treats batch processing as a special case of stream processing (a bounded stream). Its architecture is stream-first.
Key features:
- True streaming: Processes events one at a time (not micro-batches like Spark Streaming)
- Event-time processing: Handles out-of-order events based on when they occurred, not when they arrived
- Exactly-once semantics: Provides exactly-once processing guarantees through checkpointing
- Savepoints: Snapshot the entire pipeline state for upgrades, scaling, or debugging
Comparison
| Feature | MapReduce | Spark | Flink |
|---|---|---|---|
| Processing model | Batch only | Batch + micro-batch streaming | Batch + true streaming |
| Intermediate storage | Disk (HDFS) | Memory (spills to disk) | Memory (checkpoints to disk) |
| Latency | Minutes to hours | Seconds to minutes | Milliseconds to seconds |
| Fault tolerance | Re-execute failed tasks | Recompute lost RDD partitions | Checkpoint-based recovery |
| Best for | Very large batch jobs | Interactive analytics, ML | Real-time streaming, event-time processing |
Event Sourcing
Concept
Instead of storing the current state of an entity, store every state change as an immutable event. The current state is derived by replaying all events.
Traditional (mutable state):
Account { id: 1, balance: 150 }
Event sourcing (immutable log):
AccountCreated { id: 1, balance: 0 }
MoneyDeposited { id: 1, amount: 200 }
MoneyWithdrawn { id: 1, amount: 50 }
Current state = replay events: 0 + 200 - 50 = 150
Benefits
- Complete audit trail: Every change is recorded with timestamp, actor, and context
- Temporal queries: "What was the balance on January 15?" Replay events up to that date
- Event replay: Rebuild read models, fix bugs by replaying with corrected logic, build new views from historical events
- Debugging: Reproduce any state by replaying the exact sequence of events
- Decoupling: Event producers and consumers can evolve independently
Challenges
- Event schema evolution: Once events are stored, changing their schema is hard; use versioned event schemas
- Eventual consistency: Read models derived from events may lag behind the event log
- Storage growth: The event log grows forever; compaction or snapshotting is needed for old events
- Complexity: Building and maintaining projections (read models) adds architectural complexity
Event Store Implementations
| Technology | Type | Key Feature |
|---|---|---|
| EventStoreDB | Purpose-built event store | Projections, subscriptions, optimistic concurrency |
| Apache Kafka | Distributed log | High throughput, log compaction, exactly-once semantics |
| PostgreSQL | Relational DB as event store | ACID transactions on event writes; LISTEN/NOTIFY for subscribers |
| DynamoDB Streams | Change stream | Automatic change capture from DynamoDB tables |
Change Data Capture (CDC)
Concept
CDC observes all writes to a database and extracts them as a stream of change events that can be consumed by other systems. This keeps derived data stores (search indexes, caches, analytics databases) in sync with the source of truth.
Application -> PostgreSQL (source of truth)
|
v (CDC)
Kafka topic
/ | \
v v v
Elasticsearch Redis Data Warehouse
(search) (cache) (analytics)
CDC Implementation Approaches
| Approach | How It Works | Trade-off |
|---|---|---|
| Log-based (WAL parsing) | Read the database's write-ahead log and extract changes | Most reliable; low overhead; captures all changes including those from direct SQL |
| Trigger-based | Database triggers write changes to an outbox table | Works with any database; adds write overhead; may miss changes from bulk operations |
| Polling-based | Periodically query for changed rows (using updated_at timestamp) | Simplest; misses deletes; can miss rapid changes between polls; adds query load |
| Application-level | Application explicitly writes events when modifying data | Full control; risk of forgetting to emit events; dual-write problem |
CDC Tools
| Tool | Source Databases | Sink | Key Feature |
|---|---|---|---|
| Debezium | PostgreSQL, MySQL, MongoDB, SQL Server, Oracle | Kafka | WAL-based; exactly-once; schema registry integration |
| Maxwell | MySQL only | Kafka, RabbitMQ, Redis | Lightweight; MySQL binlog parsing |
| AWS DMS | Most databases | Kafka, S3, Redshift, DynamoDB | Managed service; heterogeneous migration |
| Fivetran/Airbyte | Many sources | Data warehouses | Managed ELT platforms with CDC connectors |
The Dual-Write Problem
A common anti-pattern is writing to two systems directly:
Application -> writes to PostgreSQL
-> writes to Elasticsearch
Problem: If the Elasticsearch write fails after the PostgreSQL write succeeds,
the systems are now inconsistent. Retrying may cause duplicates.
Solution: Write to one system (PostgreSQL) and use CDC to propagate to the other. The CDC pipeline handles retries, ordering, and exactly-once delivery.
Stream-Table Duality
The Core Insight
A stream and a table are two sides of the same coin:
- A stream is the changelog of a table. If you record every INSERT, UPDATE, and DELETE to a table, that sequence of changes is a stream.
- A table is the materialized state of a stream. If you replay a stream of changes from the beginning, applying each change in order, you get the current table state.
Stream:
INSERT user {id: 1, name: "Alice"}
INSERT user {id: 2, name: "Bob"}
UPDATE user {id: 1, name: "Alice Chen"}
DELETE user {id: 2}
Table (materialized from stream):
| id | name |
|----|-------------|
| 1 | Alice Chen |
Practical Application: Kafka Log Compaction
Kafka's log compaction feature retains only the latest value for each key, effectively converting a stream into a table snapshot:
Before compaction:
key=1, value="Alice"
key=2, value="Bob"
key=1, value="Alice Chen"
key=2, value=null (tombstone)
After compaction:
key=1, value="Alice Chen"
(key=2 is deleted because value is null)
This allows a new consumer to read the compacted log and reconstruct the full current state without replaying the entire history.
Exactly-Once Semantics
The Challenge
In distributed systems, messages can be lost, duplicated, or reordered. "Exactly-once" means that the effect of processing each message is reflected exactly once in the output, even in the presence of failures.
Achieving Exactly-Once
True exactly-once requires coordination between the messaging system and the processing logic:
| Approach | How It Works | Example |
|---|---|---|
| Idempotent operations | Design operations so that applying them multiple times has the same effect as applying once | SET balance = 100 is idempotent; INCREMENT balance BY 10 is not |
| Transactional output | Write output and update consumer offset in a single atomic transaction | Kafka Streams: transactional producer commits output records and consumer offsets together |
| Deduplication | Assign a unique ID to each message; recipient ignores messages it has already processed | Store processed message IDs in a set; check before processing |
| Checkpointing | Periodically save processing state; on failure, resume from last checkpoint | Flink savepoints: snapshot operator state and input positions |
Time Windowing
Why Windowing Matters
Unbounded streams have no natural "end," so you can't wait for all data before computing an aggregate. Windowing divides the stream into finite chunks for aggregation.
Window Types
| Window Type | How It Works | Use Case |
|---|---|---|
| Tumbling | Fixed-size, non-overlapping windows (e.g., every 5 minutes) | Hourly metrics, daily summaries |
| Hopping | Fixed-size windows that overlap (e.g., 10-minute windows every 5 minutes) | Smoothed averages, sliding computations |
| Session | Variable-size windows based on activity gaps (e.g., a session ends after 30 minutes of inactivity) | User session analytics, click streams |
| Global | A single window for the entire stream | Running totals, all-time aggregates |
Handling Late Events
Events may arrive after their window has closed (due to network delays, buffering, or clock skew):
- Watermarks: A timestamp that says "I believe all events before this time have arrived." Events arriving after the watermark are considered late.
- Allowed lateness: Accept late events up to a threshold (e.g., 1 hour after window closes), updating the window's result
- Side outputs: Route late events to a separate output for manual or delayed processing
Window: 10:00 - 10:05
Watermark: 10:06 (all events before 10:06 expected)
Allowed lateness: 1 hour
Event at 10:03 arriving at 10:07: accepted (within allowed lateness)
Event at 10:01 arriving at 11:30: discarded or sent to side output
Architecture Patterns
Lambda Architecture
Run both batch and stream processing pipelines in parallel:
Raw Data -> Batch Layer (Spark) -> Batch Views
-> Speed Layer (Flink) -> Real-time Views
Query: Merge batch views + real-time views for complete result
Pros: Batch provides correctness; stream provides speed. Cons: Maintaining two pipelines with the same logic is expensive and error-prone; results may differ between batch and stream.
Kappa Architecture
Use a single stream processing pipeline for everything. Reprocess historical data by replaying the event log:
Event Log (Kafka) -> Stream Processor (Flink) -> Derived Views
Reprocessing: Start a new consumer from the beginning of the log
Pros: Single codebase; simpler architecture; easier to reason about. Cons: Reprocessing large histories can be slow; requires a durable, replayable log (Kafka with long retention).
Choosing Between Lambda and Kappa
Use Lambda when:
- Exact correctness is required and stream processing approximations are unacceptable
- You have existing batch infrastructure and are adding streaming incrementally
Use Kappa when:
- Your stream processing framework provides exactly-once guarantees
- You can afford to reprocess from the log when logic changes
- Simplicity and maintainability are priorities
| 1 | # Batch and Stream Processing |
| 2 | |
| 3 | Data processing systems fall into three categories: services (online systems that handle requests), batch processing (offline systems that process large volumes of accumulated data), and stream processing (near-real-time systems that process data as it arrives). Understanding when and how to use batch and stream processing is essential for building data pipelines, analytics systems, and derived data stores. |
| 4 | |
| 5 | |
| 6 | ## Table of Contents |
| 7 | [Batch Processing: MapReduce and Beyond] |
| 8 | [Dataflow Engines: Beyond MapReduce] |
| 9 | [Event Sourcing] |
| 10 | [Change Data Capture (CDC)] |
| 11 | [Stream-Table Duality] |
| 12 | [Exactly-Once Semantics] |
| 13 | [Time Windowing] |
| 14 | [Architecture Patterns] |
| 15 | |
| 16 | |
| 17 | |
| 18 | ## Batch Processing: MapReduce and Beyond |
| 19 | |
| 20 | ### The MapReduce Paradigm |
| 21 | |
| 22 | MapReduce, popularized by Google's 2004 paper, processes large datasets by breaking computation into two phases: |
| 23 | |
| 24 | **Map phase:** Read input data, extract key-value pairs. Each mapper processes a portion of the input independently. |
| 25 | **Shuffle phase:** Framework groups all values by key, distributing them to reducers. |
| 26 | **Reduce phase:** For each key, combine all values into a result. |
| 27 | |
| 28 | |
| 29 | Input: ["the cat sat on the mat"] |
| 30 | |
| 31 | Map: |
| 32 | "the" -> 1 |
| 33 | "cat" -> 1 |
| 34 | "sat" -> 1 |
| 35 | "on" -> 1 |
| 36 | "the" -> 1 |
| 37 | "mat" -> 1 |
| 38 | |
| 39 | Shuffle: Group by key |
| 40 | "the" -> [1, 1] |
| 41 | "cat" -> [1] |
| 42 | "sat" -> [1] |
| 43 | "on" -> [1] |
| 44 | "mat" -> [1] |
| 45 | |
| 46 | Reduce: Sum values |
| 47 | "the" -> 2 |
| 48 | "cat" -> 1 |
| 49 | "sat" -> 1 |
| 50 | "on" -> 1 |
| 51 | "mat" -> 1 |
| 52 | |
| 53 | |
| 54 | ### MapReduce Strengths |
| 55 | |
| 56 | **Horizontal scalability:** Add more machines to process larger datasets; the framework handles distribution |
| 57 | **Fault tolerance:** If a mapper or reducer fails, the framework re-executes that task on another machine. Input data is immutable, so re-execution is safe. |
| 58 | **Simplicity:** The programmer only writes map and reduce functions; the framework handles parallelism, distribution, and fault tolerance |
| 59 | |
| 60 | ### MapReduce Weaknesses |
| 61 | |
| 62 | **High latency:** Each MapReduce job reads from and writes to distributed storage (HDFS), adding significant I/O overhead |
| 63 | **Chaining is awkward:** Complex computations require chaining multiple MapReduce jobs, each with its own read/write cycle |
| 64 | **No iteration support:** Machine learning algorithms that iterate over data must launch a new MapReduce job for each iteration |
| 65 | **Limited expressiveness:** Not all computations fit the map-reduce pattern naturally |
| 66 | |
| 67 | |
| 68 | |
| 69 | ## Dataflow Engines: Beyond MapReduce |
| 70 | |
| 71 | ### Apache Spark |
| 72 | |
| 73 | Spark replaces MapReduce with a more general computation model based on Resilient Distributed Datasets (RDDs) and directed acyclic graphs (DAGs) of operators. |
| 74 | |
| 75 | **Key improvements over MapReduce:** |
| 76 | **In-memory processing:** Intermediate results stay in memory instead of being written to disk between stages |
| 77 | **Arbitrary DAGs:** Computations can have multiple stages with various operators (map, filter, join, group, sort), not just map and reduce |
| 78 | **Lazy evaluation:** Spark builds a computation plan before executing, enabling the optimizer to eliminate unnecessary steps |
| 79 | **Iterative algorithms:** RDDs can be cached in memory and reused across iterations |
| 80 | |
| 81 | |
| 82 | # Spark: Word count |
| 83 | text_file = spark.read.text("hdfs://input/") |
| 84 | word_counts = (text_file |
| 85 | .select(explode(split(col("value"), " ")).alias("word")) |
| 86 | .groupBy("word") |
| 87 | .count() |
| 88 | .orderBy(desc("count"))) |
| 89 | word_counts.write.parquet("hdfs://output/") |
| 90 | |
| 91 | |
| 92 | ### Apache Flink |
| 93 | |
| 94 | Flink treats batch processing as a special case of stream processing (a bounded stream). Its architecture is stream-first. |
| 95 | |
| 96 | **Key features:** |
| 97 | **True streaming:** Processes events one at a time (not micro-batches like Spark Streaming) |
| 98 | **Event-time processing:** Handles out-of-order events based on when they occurred, not when they arrived |
| 99 | **Exactly-once semantics:** Provides exactly-once processing guarantees through checkpointing |
| 100 | **Savepoints:** Snapshot the entire pipeline state for upgrades, scaling, or debugging |
| 101 | |
| 102 | ### Comparison |
| 103 | |
| 104 | | Feature | MapReduce | Spark | Flink | |
| 105 | |---------|-----------|-------|-------| |
| 106 | | **Processing model** | Batch only | Batch + micro-batch streaming | Batch + true streaming | |
| 107 | | **Intermediate storage** | Disk (HDFS) | Memory (spills to disk) | Memory (checkpoints to disk) | |
| 108 | | **Latency** | Minutes to hours | Seconds to minutes | Milliseconds to seconds | |
| 109 | | **Fault tolerance** | Re-execute failed tasks | Recompute lost RDD partitions | Checkpoint-based recovery | |
| 110 | | **Best for** | Very large batch jobs | Interactive analytics, ML | Real-time streaming, event-time processing | |
| 111 | |
| 112 | |
| 113 | |
| 114 | ## Event Sourcing |
| 115 | |
| 116 | ### Concept |
| 117 | |
| 118 | Instead of storing the current state of an entity, store every state change as an immutable event. The current state is derived by replaying all events. |
| 119 | |
| 120 | |
| 121 | Traditional (mutable state): |
| 122 | Account { id: 1, balance: 150 } |
| 123 | |
| 124 | Event sourcing (immutable log): |
| 125 | AccountCreated { id: 1, balance: 0 } |
| 126 | MoneyDeposited { id: 1, amount: 200 } |
| 127 | MoneyWithdrawn { id: 1, amount: 50 } |
| 128 | |
| 129 | Current state = replay events: 0 + 200 - 50 = 150 |
| 130 | |
| 131 | |
| 132 | ### Benefits |
| 133 | |
| 134 | **Complete audit trail:** Every change is recorded with timestamp, actor, and context |
| 135 | **Temporal queries:** "What was the balance on January 15?" Replay events up to that date |
| 136 | **Event replay:** Rebuild read models, fix bugs by replaying with corrected logic, build new views from historical events |
| 137 | **Debugging:** Reproduce any state by replaying the exact sequence of events |
| 138 | **Decoupling:** Event producers and consumers can evolve independently |
| 139 | |
| 140 | ### Challenges |
| 141 | |
| 142 | **Event schema evolution:** Once events are stored, changing their schema is hard; use versioned event schemas |
| 143 | **Eventual consistency:** Read models derived from events may lag behind the event log |
| 144 | **Storage growth:** The event log grows forever; compaction or snapshotting is needed for old events |
| 145 | **Complexity:** Building and maintaining projections (read models) adds architectural complexity |
| 146 | |
| 147 | ### Event Store Implementations |
| 148 | |
| 149 | | Technology | Type | Key Feature | |
| 150 | |-----------|------|-------------| |
| 151 | | **EventStoreDB** | Purpose-built event store | Projections, subscriptions, optimistic concurrency | |
| 152 | | **Apache Kafka** | Distributed log | High throughput, log compaction, exactly-once semantics | |
| 153 | | **PostgreSQL** | Relational DB as event store | ACID transactions on event writes; LISTEN/NOTIFY for subscribers | |
| 154 | | **DynamoDB Streams** | Change stream | Automatic change capture from DynamoDB tables | |
| 155 | |
| 156 | |
| 157 | |
| 158 | ## Change Data Capture (CDC) |
| 159 | |
| 160 | ### Concept |
| 161 | |
| 162 | CDC observes all writes to a database and extracts them as a stream of change events that can be consumed by other systems. This keeps derived data stores (search indexes, caches, analytics databases) in sync with the source of truth. |
| 163 | |
| 164 | |
| 165 | Application -> PostgreSQL (source of truth) |
| 166 | | |
| 167 | v (CDC) |
| 168 | Kafka topic |
| 169 | / | \ |
| 170 | v v v |
| 171 | Elasticsearch Redis Data Warehouse |
| 172 | (search) (cache) (analytics) |
| 173 | |
| 174 | |
| 175 | ### CDC Implementation Approaches |
| 176 | |
| 177 | | Approach | How It Works | Trade-off | |
| 178 | |----------|-------------|-----------| |
| 179 | | **Log-based (WAL parsing)** | Read the database's write-ahead log and extract changes | Most reliable; low overhead; captures all changes including those from direct SQL | |
| 180 | | **Trigger-based** | Database triggers write changes to an outbox table | Works with any database; adds write overhead; may miss changes from bulk operations | |
| 181 | | **Polling-based** | Periodically query for changed rows (using updated_at timestamp) | Simplest; misses deletes; can miss rapid changes between polls; adds query load | |
| 182 | | **Application-level** | Application explicitly writes events when modifying data | Full control; risk of forgetting to emit events; dual-write problem | |
| 183 | |
| 184 | ### CDC Tools |
| 185 | |
| 186 | | Tool | Source Databases | Sink | Key Feature | |
| 187 | |------|-----------------|------|-------------| |
| 188 | | **Debezium** | PostgreSQL, MySQL, MongoDB, SQL Server, Oracle | Kafka | WAL-based; exactly-once; schema registry integration | |
| 189 | | **Maxwell** | MySQL only | Kafka, RabbitMQ, Redis | Lightweight; MySQL binlog parsing | |
| 190 | | **AWS DMS** | Most databases | Kafka, S3, Redshift, DynamoDB | Managed service; heterogeneous migration | |
| 191 | | **Fivetran/Airbyte** | Many sources | Data warehouses | Managed ELT platforms with CDC connectors | |
| 192 | |
| 193 | ### The Dual-Write Problem |
| 194 | |
| 195 | A common anti-pattern is writing to two systems directly: |
| 196 | |
| 197 | |
| 198 | Application -> writes to PostgreSQL |
| 199 | -> writes to Elasticsearch |
| 200 | |
| 201 | Problem: If the Elasticsearch write fails after the PostgreSQL write succeeds, |
| 202 | the systems are now inconsistent. Retrying may cause duplicates. |
| 203 | |
| 204 | |
| 205 | **Solution:** Write to one system (PostgreSQL) and use CDC to propagate to the other. The CDC pipeline handles retries, ordering, and exactly-once delivery. |
| 206 | |
| 207 | |
| 208 | |
| 209 | ## Stream-Table Duality |
| 210 | |
| 211 | ### The Core Insight |
| 212 | |
| 213 | A stream and a table are two sides of the same coin: |
| 214 | |
| 215 | **A stream is the changelog of a table.** If you record every INSERT, UPDATE, and DELETE to a table, that sequence of changes is a stream. |
| 216 | **A table is the materialized state of a stream.** If you replay a stream of changes from the beginning, applying each change in order, you get the current table state. |
| 217 | |
| 218 | |
| 219 | Stream: |
| 220 | INSERT user {id: 1, name: "Alice"} |
| 221 | INSERT user {id: 2, name: "Bob"} |
| 222 | UPDATE user {id: 1, name: "Alice Chen"} |
| 223 | DELETE user {id: 2} |
| 224 | |
| 225 | Table (materialized from stream): |
| 226 | | id | name | |
| 227 | |----|-------------| |
| 228 | | 1 | Alice Chen | |
| 229 | |
| 230 | |
| 231 | ### Practical Application: Kafka Log Compaction |
| 232 | |
| 233 | Kafka's log compaction feature retains only the latest value for each key, effectively converting a stream into a table snapshot: |
| 234 | |
| 235 | |
| 236 | Before compaction: |
| 237 | key=1, value="Alice" |
| 238 | key=2, value="Bob" |
| 239 | key=1, value="Alice Chen" |
| 240 | key=2, value=null (tombstone) |
| 241 | |
| 242 | After compaction: |
| 243 | key=1, value="Alice Chen" |
| 244 | (key=2 is deleted because value is null) |
| 245 | |
| 246 | |
| 247 | This allows a new consumer to read the compacted log and reconstruct the full current state without replaying the entire history. |
| 248 | |
| 249 | |
| 250 | |
| 251 | ## Exactly-Once Semantics |
| 252 | |
| 253 | ### The Challenge |
| 254 | |
| 255 | In distributed systems, messages can be lost, duplicated, or reordered. "Exactly-once" means that the effect of processing each message is reflected exactly once in the output, even in the presence of failures. |
| 256 | |
| 257 | ### Achieving Exactly-Once |
| 258 | |
| 259 | True exactly-once requires coordination between the messaging system and the processing logic: |
| 260 | |
| 261 | | Approach | How It Works | Example | |
| 262 | |----------|-------------|---------| |
| 263 | | **Idempotent operations** | Design operations so that applying them multiple times has the same effect as applying once | `SET balance = 100` is idempotent; `INCREMENT balance BY 10` is not | |
| 264 | | **Transactional output** | Write output and update consumer offset in a single atomic transaction | Kafka Streams: transactional producer commits output records and consumer offsets together | |
| 265 | | **Deduplication** | Assign a unique ID to each message; recipient ignores messages it has already processed | Store processed message IDs in a set; check before processing | |
| 266 | | **Checkpointing** | Periodically save processing state; on failure, resume from last checkpoint | Flink savepoints: snapshot operator state and input positions | |
| 267 | |
| 268 | |
| 269 | |
| 270 | ## Time Windowing |
| 271 | |
| 272 | ### Why Windowing Matters |
| 273 | |
| 274 | Unbounded streams have no natural "end," so you can't wait for all data before computing an aggregate. Windowing divides the stream into finite chunks for aggregation. |
| 275 | |
| 276 | ### Window Types |
| 277 | |
| 278 | | Window Type | How It Works | Use Case | |
| 279 | |------------|-------------|----------| |
| 280 | | **Tumbling** | Fixed-size, non-overlapping windows (e.g., every 5 minutes) | Hourly metrics, daily summaries | |
| 281 | | **Hopping** | Fixed-size windows that overlap (e.g., 10-minute windows every 5 minutes) | Smoothed averages, sliding computations | |
| 282 | | **Session** | Variable-size windows based on activity gaps (e.g., a session ends after 30 minutes of inactivity) | User session analytics, click streams | |
| 283 | | **Global** | A single window for the entire stream | Running totals, all-time aggregates | |
| 284 | |
| 285 | ### Handling Late Events |
| 286 | |
| 287 | Events may arrive after their window has closed (due to network delays, buffering, or clock skew): |
| 288 | |
| 289 | **Watermarks:** A timestamp that says "I believe all events before this time have arrived." Events arriving after the watermark are considered late. |
| 290 | **Allowed lateness:** Accept late events up to a threshold (e.g., 1 hour after window closes), updating the window's result |
| 291 | **Side outputs:** Route late events to a separate output for manual or delayed processing |
| 292 | |
| 293 | |
| 294 | Window: 10:00 - 10:05 |
| 295 | Watermark: 10:06 (all events before 10:06 expected) |
| 296 | Allowed lateness: 1 hour |
| 297 | |
| 298 | Event at 10:03 arriving at 10:07: accepted (within allowed lateness) |
| 299 | Event at 10:01 arriving at 11:30: discarded or sent to side output |
| 300 | |
| 301 | |
| 302 | |
| 303 | |
| 304 | ## Architecture Patterns |
| 305 | |
| 306 | ### Lambda Architecture |
| 307 | |
| 308 | Run both batch and stream processing pipelines in parallel: |
| 309 | |
| 310 | |
| 311 | Raw Data -> Batch Layer (Spark) -> Batch Views |
| 312 | -> Speed Layer (Flink) -> Real-time Views |
| 313 | |
| 314 | Query: Merge batch views + real-time views for complete result |
| 315 | |
| 316 | |
| 317 | **Pros:** Batch provides correctness; stream provides speed. |
| 318 | **Cons:** Maintaining two pipelines with the same logic is expensive and error-prone; results may differ between batch and stream. |
| 319 | |
| 320 | ### Kappa Architecture |
| 321 | |
| 322 | Use a single stream processing pipeline for everything. Reprocess historical data by replaying the event log: |
| 323 | |
| 324 | |
| 325 | Event Log (Kafka) -> Stream Processor (Flink) -> Derived Views |
| 326 | |
| 327 | Reprocessing: Start a new consumer from the beginning of the log |
| 328 | |
| 329 | |
| 330 | **Pros:** Single codebase; simpler architecture; easier to reason about. |
| 331 | **Cons:** Reprocessing large histories can be slow; requires a durable, replayable log (Kafka with long retention). |
| 332 | |
| 333 | ### Choosing Between Lambda and Kappa |
| 334 | |
| 335 | Use **Lambda** when: |
| 336 | Exact correctness is required and stream processing approximations are unacceptable |
| 337 | You have existing batch infrastructure and are adding streaming incrementally |
| 338 | |
| 339 | Use **Kappa** when: |
| 340 | Your stream processing framework provides exactly-once guarantees |
| 341 | You can afford to reprocess from the log when logic changes |
| 342 | Simplicity and maintainability are priorities |
| 343 |
Discussion
Browse more free Claude skills.