Apache Kafka, Flink, RocksDB State Backend, and Exactly-Once Processing Semantics Under Failure

Distributed Architecture Takeaway

Achieving 1,000,000 events per second with stateful windowing requires eliminating network serializations and disk I/O. Apache Flink with RocksDB state backends and Kafka transactional producers delivers deterministic Exactly-Once Semantics (EOS).

Empirical Architecture Comparison: Batch Processing vs. Micro-Batching vs. True Stream Processing

DimensionTraditional Batch (Apache Spark)Micro-Batching (Spark Streaming)Continuous Streaming (Apache Flink)
Processing ParadigmHours / Nightly batch jobsSub-second micro-batches (200ms - 2s)Event-driven: Processes individual events continuously
Latency ProfileMinutes to Hours500 ms - 2,000 msSub-10 milliseconds (Real-time)
State ManagementEphemeral memory tablesCheckpointing RDD memoryIncremental checkpointing via embedded RocksDB
Failure RecoveryRe-executes entire batch stageRecomputes failed micro-batch RDDsChandy-Lamport distributed snapshot barrier rollback

1. High-Throughput Event Ingestion with Apache Kafka

Sustaining a million incoming JSON payloads per second requires optimizing Kafka producers and partition topologies. Individual synchronous network writes collapse under socket overhead. Producers must batch records using snappy compression and configure memory buffers:
# High-throughput Kafka producer configuration
linger.ms=10
batch.size=131072 # 128KB batch buffer
compression.type=snappy
acks=all
max.in.flight.requests.per.connection=5
enable.idempotence=true
With 32 topic partitions distributed across a 6-node broker cluster, write throughput scales linearly without disk bottlenecking.

2. Stateful Windowing and RocksDB State Backends

When computing rolling 10-minute aggregations across millions of user keys, state cannot fit into JVM heap memory without triggering devastating Garbage Collection pauses. Apache Flink offloads managed state to an out-of-core embedded RocksDB instance running directly on each TaskManager node: $$\text{Memory Overhead} = O(\text{Active Key Working Set}) \quad \text{vs.} \quad O(\text{Total Keys in RocksDB on NVMe})$$ RocksDB stores key-value state in Log-Structured Merge (LSM) trees, writing sequential append-only logs and flushing SSTables to NVMe storage with near-zero latency penalty.

3. The Chandy-Lamport Distributed Snapshot Algorithm

To guarantee Exactly-Once Semantics (EOS) across network partitions and node crashes, Flink implements the Chandy-Lamport distributed snapshot algorithm via checkpoint barriers:
  1. A checkpoint coordinator periodically injects checkpoint barriers $S_n$ into Kafka source streams.
  2. Barriers flow through operator DAGs alongside normal data tuples.
  3. When an operator receives barrier $S_n$ across all input channels, it freezes state and asynchronously writes a snapshot delta to distributed S3 storage.
  4. Upon node crash, all operators rewind their Kafka consumer offsets to the exact state saved in checkpoint $S_n$.

4. Two-Phase Commit (2PC) Sinks for End-to-End Exactly-Once

Ensuring EOS inside Flink is insufficient if the downstream database receives duplicate writes during crashes. Flink utilizes a Two-Phase Commit protocol for sink connectors: transactions are pre-committed on every barrier, and formally committed only when the global checkpoint coordinator confirms that all cluster operators have successfully persisted state.