Designing Streaming and Real-Time Data Processing

1. Designing Stream Processing Architecture

ComponentDetail
SourceKafka, Kinesis, Pulsar, CDC
ProcessorFlink, Spark Streaming, Kafka Streams
State storeRocksDB local + remote backup
SinkDB, search, downstream topic
DeliveryAt-least-once or exactly-once

2. Designing Event Streaming Platform

PropertyDetail
Durable logReplicated, retention by time/size
ThroughputMillions of events/sec horizontally
ReplayFrom any offset
ToolsKafka, Pulsar, Redpanda, Kinesis

3. Designing Real-Time ETL Pipeline

   Sources ─▶ CDC/Connectors ─▶ Kafka ─▶ Stream Processor ─▶ Sinks (DW, Search, Cache)
                                          │
                                          └─▶ Side outputs (alerts, metrics)
      
StageTool
IngestDebezium, Kafka Connect
TransformFlink SQL, ksqlDB, Spark
SinkIceberg, ES, ClickHouse

4. Designing Stream Windowing Strategy

WindowDetail
TumblingFixed, non-overlapping
SlidingOverlapping; smoother stats
SessionActivity-based gap
Global + triggerCustom firing
Time semanticsEvent time (recommended) vs processing time

5. Designing Stream Aggregation Patterns

PatternDetail
Count / sum / avg per windowBasic
Top-KHeavy hitters; Count-Min Sketch
Approx distinctHyperLogLog
Quantilest-digest, KLL

6. Designing Complex Event Processing

CapabilityDetail
Pattern matching"A then B within 5min and not C"
CEP enginesFlink CEP, Esper, Drools Fusion
Use casesFraud, IoT alerts, security

7. Designing Stream State Management

AspectDetail
Keyed statePer partition key
StorageLocal RocksDB; remote checkpoint
TTLAuto-expire stale state
Size mgmtAvoid unbounded growth

8. Designing Stream Join Operations

TypeDetail
Stream-streamBoth windowed; bounded state
Stream-tableLookup join (KTable)
Temporal joinVersioned table at event time
Interval joinWithin time window

9. Designing Late Data Handling

MechanismDetail
WatermarksEstimate of event-time progress
Allowed latenessWindow stays open extra duration
Side outputRoute late events for special handling
Restate / reprocessReplay batch for corrections

10. Designing Stream Backpressure Management

StrategyDetail
Reactive backpressureSlow producer when consumer slow (Flink)
Buffer + spillDisk overflow
Drop / sampleFor metrics streams
Auto-scale workersAdd tasks under lag

11. Designing Stream Checkpointing

AspectDetail
Frequency10s–1min typical
StorageS3 / HDFS / GCS
ModeAligned (exactly-once) vs unaligned (latency)
RestoreResume from latest successful checkpoint

12. Designing Stream Fault Tolerance

ElementDetail
Source replayKafka offsets persistent
State checkpointsRestorable state
Idempotent sinksOr transactional 2PC
Task restartAuto by orchestrator (Flink/K8s)