Batch and stream processing skill

Data processing systems fall into three categories: services (online systems that handle requests), batch processing (offline systems that…

by wondelai·MIT license·★ 2,235 Stars on the repo·GitHub ↗

Use now

Files of Batch and stream processing

wondelai/main1 file
batch-stream.md
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

  1. Batch Processing: MapReduce and Beyond
  2. Dataflow Engines: Beyond MapReduce
  3. Event Sourcing
  4. Change Data Capture (CDC)
  5. Stream-Table Duality
  6. Exactly-Once Semantics
  7. Time Windowing
  8. 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:

  1. Map phase: Read input data, extract key-value pairs. Each mapper processes a portion of the input independently.
  2. Shuffle phase: Framework groups all values by key, distributing them to reducers.
  3. 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/")

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 
3Data 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
71. [Batch Processing: MapReduce and Beyond](#batch-processing-mapreduce-and-beyond)
82. [Dataflow Engines: Beyond MapReduce](#dataflow-engines-beyond-mapreduce)
93. [Event Sourcing](#event-sourcing)
104. [Change Data Capture (CDC)](#change-data-capture-cdc)
115. [Stream-Table Duality](#stream-table-duality)
126. [Exactly-Once Semantics](#exactly-once-semantics)
137. [Time Windowing](#time-windowing)
148. [Architecture Patterns](#architecture-patterns)
15 
16---
17 
18## Batch Processing: MapReduce and Beyond
19 
20### The MapReduce Paradigm
21 
22MapReduce, popularized by Google's 2004 paper, processes large datasets by breaking computation into two phases:
23 
241. **Map phase:** Read input data, extract key-value pairs. Each mapper processes a portion of the input independently.
252. **Shuffle phase:** Framework groups all values by key, distributing them to reducers.
263. **Reduce phase:** For each key, combine all values into a result.
27 
28```
29Input: ["the cat sat on the mat"]
30 
31Map:
32 "the" -> 1
33 "cat" -> 1
34 "sat" -> 1
35 "on" -> 1
36 "the" -> 1
37 "mat" -> 1
38 
39Shuffle: Group by key
40 "the" -> [1, 1]
41 "cat" -> [1]
42 "sat" -> [1]
43 "on" -> [1]
44 "mat" -> [1]
45 
46Reduce: 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 
73Spark 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```python
82# Spark: Word count
83text_file = spark.read.text("hdfs://input/")
84word_counts = (text_file
85 .select(explode(split(col("value"), " ")).alias("word"))
86 .groupBy("word")
87 .count()
88 .orderBy(desc("count")))
89word_counts.write.parquet("hdfs://output/")
90```
91 
92### Apache Flink
93 
94Flink 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 
118Instead 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```
121Traditional (mutable state):
122 Account { id: 1, balance: 150 }
123 
124Event 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 
162CDC 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```
165Application -> 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 
195A common anti-pattern is writing to two systems directly:
196 
197```
198Application -> writes to PostgreSQL
199 -> writes to Elasticsearch
200 
201Problem: If the Elasticsearch write fails after the PostgreSQL write succeeds,
202the 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 
213A 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```
219Stream:
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 
225Table (materialized from stream):
226 | id | name |
227 |----|-------------|
228 | 1 | Alice Chen |
229```
230 
231### Practical Application: Kafka Log Compaction
232 
233Kafka's log compaction feature retains only the latest value for each key, effectively converting a stream into a table snapshot:
234 
235```
236Before compaction:
237 key=1, value="Alice"
238 key=2, value="Bob"
239 key=1, value="Alice Chen"
240 key=2, value=null (tombstone)
241 
242After compaction:
243 key=1, value="Alice Chen"
244 (key=2 is deleted because value is null)
245```
246 
247This 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 
255In 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 
259True 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 
274Unbounded 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 
287Events 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```
294Window: 10:00 - 10:05
295Watermark: 10:06 (all events before 10:06 expected)
296Allowed lateness: 1 hour
297 
298Event at 10:03 arriving at 10:07: accepted (within allowed lateness)
299Event 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 
308Run both batch and stream processing pipelines in parallel:
309 
310```
311Raw Data -> Batch Layer (Spark) -> Batch Views
312 -> Speed Layer (Flink) -> Real-time Views
313 
314Query: 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 
322Use a single stream processing pipeline for everything. Reprocess historical data by replaying the event log:
323 
324```
325Event Log (Kafka) -> Stream Processor (Flink) -> Derived Views
326 
327Reprocessing: 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 
335Use **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 
339Use **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