Working with Streaming Data
1. Understanding Stream Processing Models
| Model | Description | Examples |
|---|---|---|
| Record-at-a-time | Process each event individually | Kafka Streams, Flink |
| Micro-batch | Small batches (sec-scale) | Spark Streaming |
| Continuous | True streaming, ms-scale | Flink, Beam, Materialize |
| Lambda architecture | Batch + speed layers | Hadoop + Storm (legacy) |
| Kappa architecture | Single stream layer | Modern default |
2. Implementing Exactly-Once Stream Processing
| System | Mechanism |
|---|---|
| Kafka Streams | Transactions + idempotent producer |
| Flink | Async snapshots + 2PC sink |
| Spark Structured Streaming | Idempotent sinks + WAL |
| Beam | Runner-dependent (Dataflow has EOS) |
3. Implementing Windowing Operations
| Window Type | Behavior | Example |
|---|---|---|
| Tumbling | Fixed, non-overlapping | 5-min counts |
| Sliding | Fixed size, overlapping | Last 5 min, every 1 min |
| Session | Gap of inactivity defines window | User session = 30-min idle gap |
| Global | All elements; needs trigger | Custom triggers |
4. Implementing Stream Joins
| Join Type | Detail |
|---|---|
| Stream-stream | Both bounded by window |
| Stream-table | Enrich with current state |
| Table-table | Materialized view join |
| Temporal join | As-of-time semantics |
5. Implementing Stateful Stream Processing
| State Backend | Detail |
|---|---|
| In-memory | Fast; lost on crash unless checkpointed |
| RocksDB (Flink) | Off-heap, large state |
| Changelog topic (Kafka Streams) | Recover by replay |
| Checkpointing | Periodic snapshots |
6. Implementing Backpressure Handling
| Mechanism | System |
|---|---|
| Pull-based | Consumer requests N (Reactive Streams) |
| Watermark to producer | Flink's credit-based flow control |
| Buffer + drop / spill | Bounded queues |
| Adaptive throughput | Reduce ingestion rate |
7. Implementing Stream Partitioning
| Strategy | Use |
|---|---|
| Hash by key | Co-locate by entity |
| Round-robin | Even distribution; no key locality |
| Custom | Geo, tenant, content-based |
| Sticky (producer) | Batch efficiency |
8. Understanding Watermarks and Late Data
| Concept | Detail |
|---|---|
| Watermark W(t) | "No event with ts ≤ t will arrive" |
| Allowed lateness | Window remains open beyond watermark |
| Late events | Re-fire window or sent to side output |
| Punctuated watermarks | Embedded in stream |
9. Implementing Event Time vs Processing Time
| Type | Source | Use |
|---|---|---|
| Event time | Embedded timestamp from producer | Correctness, replay |
| Ingestion time | Broker-assigned | Compromise |
| Processing time | Wall clock at processor | Low-latency, no correctness |
10. Implementing Stream-Table Duality
| View | Definition |
|---|---|
| Stream → Table | Aggregate / fold events into current state |
| Table → Stream | Changelog of every update |
| Compaction | Kafka log compaction = table view |
| KTable / GlobalKTable | Kafka Streams abstraction |