Working with Streaming Data

1. Understanding Stream Processing Models

ModelDescriptionExamples
Record-at-a-timeProcess each event individuallyKafka Streams, Flink
Micro-batchSmall batches (sec-scale)Spark Streaming
ContinuousTrue streaming, ms-scaleFlink, Beam, Materialize
Lambda architectureBatch + speed layersHadoop + Storm (legacy)
Kappa architectureSingle stream layerModern default

2. Implementing Exactly-Once Stream Processing

SystemMechanism
Kafka StreamsTransactions + idempotent producer
FlinkAsync snapshots + 2PC sink
Spark Structured StreamingIdempotent sinks + WAL
BeamRunner-dependent (Dataflow has EOS)

3. Implementing Windowing Operations

Window TypeBehaviorExample
TumblingFixed, non-overlapping5-min counts
SlidingFixed size, overlappingLast 5 min, every 1 min
SessionGap of inactivity defines windowUser session = 30-min idle gap
GlobalAll elements; needs triggerCustom triggers

4. Implementing Stream Joins

Join TypeDetail
Stream-streamBoth bounded by window
Stream-tableEnrich with current state
Table-tableMaterialized view join
Temporal joinAs-of-time semantics

5. Implementing Stateful Stream Processing

State BackendDetail
In-memoryFast; lost on crash unless checkpointed
RocksDB (Flink)Off-heap, large state
Changelog topic (Kafka Streams)Recover by replay
CheckpointingPeriodic snapshots

6. Implementing Backpressure Handling

MechanismSystem
Pull-basedConsumer requests N (Reactive Streams)
Watermark to producerFlink's credit-based flow control
Buffer + drop / spillBounded queues
Adaptive throughputReduce ingestion rate

7. Implementing Stream Partitioning

StrategyUse
Hash by keyCo-locate by entity
Round-robinEven distribution; no key locality
CustomGeo, tenant, content-based
Sticky (producer)Batch efficiency

8. Understanding Watermarks and Late Data

ConceptDetail
Watermark W(t)"No event with ts ≤ t will arrive"
Allowed latenessWindow remains open beyond watermark
Late eventsRe-fire window or sent to side output
Punctuated watermarksEmbedded in stream

9. Implementing Event Time vs Processing Time

TypeSourceUse
Event timeEmbedded timestamp from producerCorrectness, replay
Ingestion timeBroker-assignedCompromise
Processing timeWall clock at processorLow-latency, no correctness

10. Implementing Stream-Table Duality

ViewDefinition
Stream → TableAggregate / fold events into current state
Table → StreamChangelog of every update
CompactionKafka log compaction = table view
KTable / GlobalKTableKafka Streams abstraction