Designing Streaming and Real-Time Data Processing
1. Designing Stream Processing Architecture
| Component | Detail |
|---|---|
| Source | Kafka, Kinesis, Pulsar, CDC |
| Processor | Flink, Spark Streaming, Kafka Streams |
| State store | RocksDB local + remote backup |
| Sink | DB, search, downstream topic |
| Delivery | At-least-once or exactly-once |
2. Designing Event Streaming Platform
| Property | Detail |
|---|---|
| Durable log | Replicated, retention by time/size |
| Throughput | Millions of events/sec horizontally |
| Replay | From any offset |
| Tools | Kafka, Pulsar, Redpanda, Kinesis |
3. Designing Real-Time ETL Pipeline
Sources ─▶ CDC/Connectors ─▶ Kafka ─▶ Stream Processor ─▶ Sinks (DW, Search, Cache)
│
└─▶ Side outputs (alerts, metrics)
| Stage | Tool |
|---|---|
| Ingest | Debezium, Kafka Connect |
| Transform | Flink SQL, ksqlDB, Spark |
| Sink | Iceberg, ES, ClickHouse |
4. Designing Stream Windowing Strategy
| Window | Detail |
|---|---|
| Tumbling | Fixed, non-overlapping |
| Sliding | Overlapping; smoother stats |
| Session | Activity-based gap |
| Global + trigger | Custom firing |
| Time semantics | Event time (recommended) vs processing time |
5. Designing Stream Aggregation Patterns
| Pattern | Detail |
|---|---|
| Count / sum / avg per window | Basic |
| Top-K | Heavy hitters; Count-Min Sketch |
| Approx distinct | HyperLogLog |
| Quantiles | t-digest, KLL |
6. Designing Complex Event Processing
| Capability | Detail |
|---|---|
| Pattern matching | "A then B within 5min and not C" |
| CEP engines | Flink CEP, Esper, Drools Fusion |
| Use cases | Fraud, IoT alerts, security |
7. Designing Stream State Management
| Aspect | Detail |
|---|---|
| Keyed state | Per partition key |
| Storage | Local RocksDB; remote checkpoint |
| TTL | Auto-expire stale state |
| Size mgmt | Avoid unbounded growth |
8. Designing Stream Join Operations
| Type | Detail |
|---|---|
| Stream-stream | Both windowed; bounded state |
| Stream-table | Lookup join (KTable) |
| Temporal join | Versioned table at event time |
| Interval join | Within time window |
9. Designing Late Data Handling
| Mechanism | Detail |
|---|---|
| Watermarks | Estimate of event-time progress |
| Allowed lateness | Window stays open extra duration |
| Side output | Route late events for special handling |
| Restate / reprocess | Replay batch for corrections |
10. Designing Stream Backpressure Management
| Strategy | Detail |
|---|---|
| Reactive backpressure | Slow producer when consumer slow (Flink) |
| Buffer + spill | Disk overflow |
| Drop / sample | For metrics streams |
| Auto-scale workers | Add tasks under lag |
11. Designing Stream Checkpointing
| Aspect | Detail |
|---|---|
| Frequency | 10s–1min typical |
| Storage | S3 / HDFS / GCS |
| Mode | Aligned (exactly-once) vs unaligned (latency) |
| Restore | Resume from latest successful checkpoint |
12. Designing Stream Fault Tolerance
| Element | Detail |
|---|---|
| Source replay | Kafka offsets persistent |
| State checkpoints | Restorable state |
| Idempotent sinks | Or transactional 2PC |
| Task restart | Auto by orchestrator (Flink/K8s) |