Interview Setup
Interview Prompt
Design a distributed stream processing platform (Apache Flink) that ingests 50M events/sec, runs stateful windowed aggregations, and delivers results to Kafka, Elasticsearch, and ClickHouse with sub-second latency.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Do we need true exactly-once end-to-end, or is at-least-once with idempotent sinks acceptable? | Exactly-once requires transactional Kafka sinks and idempotent writes, which introduces additional latency. While analytics dashboards often tolerate at-least-once delivery, billing counters strictly require exactly-once processing. |
| How do we handle late-arriving events: drop them, route to side-outputs, or update within allowed lateness? | At 50M events/sec, retaining unbounded late data balloons RocksDB state. Allowed lateness of 1 to 5 minutes is a concrete cost trade-off. |
| What is the checkpoint interval vs recovery time budget? | With 50 TB of keyed state, assuming roughly 1% of aggregate state is newly materialized or modified during that interval, each checkpoint cycle writes approximately 500 GB incrementally to S3. Setting the interval too short prevents checkpoints from completing, whereas setting it too long increases the amount of source data that may need to be replayed from the last completed checkpoint. |
| How many stateful jobs and what keyed state per job? | 500 jobs x 100 GB = 50 TB total aggregate state. While 2,000 TaskManagers (8 cores, 32 GB each) serves as an illustrative capacity assumption for this scenario, actual sizing depends on CPU, managed memory, local state and disk capacity, network bandwidth, state distribution, checkpoint throughput, and operational headroom. |
Scope
In scope
- Windowing (tumbling, sliding, session)
- Watermarks & late data
- Exactly-once in Flink/Spark
- Stateful operators
- Checkpointing
- Capacity estimation with shown math
Out of scope (state explicitly)
- Batch ETL warehouse design (Spark offline)
- Exactly-once Kafka broker internals
- ML feature store unless ranking is in scope
Functional Requirements
Start by asking your interviewer about event sources, windowed aggregations, and delivery guarantees (at-least-once vs exactly-once). Whether you expose Flink SQL or only the DataStream API is a common scope pivot.
- Process unbounded (infinite) event streams in real-time with sub-second per-operator processing latency
- Support windowed aggregations: tumbling, sliding, session, and global windows
- Handle event-time semantics with watermarks for out-of-order data
- Support flexible delivery semantics: offer at-least-once as the high-throughput baseline for telemetry and analytics, while supporting end-to-end exactly-once (via two-phase commit or idempotent sinks) for billing, inventory, and ledger workloads where duplicate side effects are unacceptable
- Support stateful operators: keyed state (per-key counters, aggregates, ML features) persisted durably
- Stream-to-stream joins: join two event streams on a key within a time window
- Stream-to-table enrichment joins: enrich streaming events with slowly-changing dimension data
- SQL interface for analysts (Flink SQL) alongside programmatic DataStream API for engineers
- Connectors to sources (Kafka, Kinesis, files) and sinks (Kafka, Elasticsearch, PostgreSQL, S3, ClickHouse)
- Savepoints: manually triggered consistent snapshots for version upgrades and job migration
- Backpressure handling: slow operators should not cause data loss
Non-Functional Requirements
Your interviewer will focus heavily on exactly-once semantics and checkpoint recovery for this problem. Watermarks and event-time processing represent critical areas where senior candidates demonstrate architectural depth, so be sure to highlight them early in the discussion.
- Low Latency: Sub-100ms (p99) operator and non-windowed processing latency for streaming filters and stateless transformations, while windowed event-time end-to-end delivery additionally depends on watermark progression and allowed lateness (< 5s p99 scenario operational SLO)
- High Throughput: 10M+ events/sec per job, while the cluster handles 100M+ events/sec
- Delivery Semantics: Configurable per pipeline—at-least-once baseline for high-volume analytics, and end-to-end exactly-once (or effectively-once via idempotent upserts) for financial and audit ledgers to guarantee zero duplicate side effects
- Scalability: Horizontal scaling by adding TaskManagers to increase parallelism
- Fault Tolerance & RTO: Automatic recovery from TaskManager failures with a target RTO < 5 minutes, and typical illustrative recovery taking approximately 30–120 seconds depending on state size, container startup, and replay
- Stateful at Scale: Terabyte-scale keyed state per job backed by RocksDB
- Elastic Scaling: Dynamic scale up and scale down without losing state via reactive mode
- Backpressure Handling: Credit-based flow control preventing buffer overflow and dropped data under slow consumers
Capacity Estimations
Run this math before you set parallelism. Events per second and keyed state size tell you how many TaskManagers you need and whether RocksDB is mandatory.
| Metric | Calculation | Value |
|---|---|---|
| Events ingested / sec | 50M | |
| Avg event size | 500 bytes | |
| Ingestion throughput | 50M x 500B | 25 GB/sec |
| Stateful jobs | 500 | |
| Avg keyed state per job | 100 GB | |
| Total state | 500 x 100 GB | 50 TB |
| Checkpoints (every 1 min) | ~1% state churn on 50 TB | ~500 GB incremental per checkpoint cycle |
| Checkpoint storage | S3 | ~10 TB (retained checkpoints) |
| TaskManagers | 2,000 (8 cores, 32 GB each) | |
| JobManagers | 3 (HA quorum) | |
| Network (inter-operator shuffle) | ~10 GB/sec cluster-wide |
Architecture Diagram
In the interview: clarify delivery guarantees upfront because at-least-once delivery is simpler, whereas end-to-end exactly-once fundamentally alters your checkpointing and sink design.
Walk your interviewer through the diagram by data flow. Events enter from Kafka, get partitioned by key, and flow through a DAG of stateful operators. We checkpoint operator state to durable storage so TaskManager crashes recover without data loss. Windowed aggregations fire on event time using watermarks for out-of-order events. Results sink to databases or downstream topics with idempotent or two-phase-commit writes for end-to-end exactly-once.
Component Deep Dives
Next we walk through each box on the diagram. Start with cluster coordination, then processing semantics, then fault tolerance.
JobManager: The Brain
The JobManager coordinates the entire cluster as a centralized coordinator for state and execution without bottlenecking the data path (conceptual internal coordinator architecture):
The JobManager serves as the control plane for the Flink cluster, comprising conceptual internal coordinator components:
Core sub-components:
1. Dispatcher: Accepts job submissions containing compiled job graphs and JARs, initiating JobMaster instances.
2. JobMaster (one instance per active job): Manages the execution lifecycle of a single job.
- Compiles the logical stream graph into a physical execution DAG.
- Assigns operator parallelism across task slots.
- Requests and releases slot allocations from the ResourceManager.
- Deploys task execution pipelines to TaskManagers.
- Coordinates distributed checkpointing.
3. ResourceManager: Manages cluster-wide task slot availability.
- Kubernetes mode: dynamically provisions new TaskManager worker pods.
- YARN mode: requests new container allocations from the YARN NodeManager.
- Standalone mode: coordinates across pre-provisioned static TaskManager instances.
4. Checkpoint Coordinator (internal coordinator component): Periodically injects checkpoint barriers and collects distributed acknowledgments.
High Availability (HA) deployment:
Multiple JobManagers are deployed in an active-standby quorum using ZooKeeper or Kubernetes leader election.
If the active JobManager fails, a standby instance is elected and reconstructs the job graph from persistent metadata in ZooKeeper or ConfigMaps, resuming execution from the latest completed checkpoint.TaskManager: The Muscle
TaskManagers execute the data pipeline directly where each worker process hosts operator subtasks alongside local RocksDB state.
Each TaskManager runs as a dedicated JVM worker process configured with: - Task slots: resource-allocation and scheduling units, typically sized relative to physical CPU core count. - Subtask execution: operators and subtasks execute within JVM threads, sharing slots through operator chaining or slot-sharing groups. - Managed memory: off-heap memory region dedicated to RocksDB block caches, hash tables, and sorting. - Network buffers: credit-based flow control buffers for inter-worker data exchange. - State backend: EmbeddedRocksDB utilizing local NVMe SSDs for terabyte-scale state, or heap HashMap for small state. Task slot resource isolation: - Task slots share the same JVM heap but maintain isolated managed memory boundaries. - Task slots provide resource-allocation units for operator subtasks. Chained operators can execute together within a task thread, while multiple subtasks/operators from the same job may share a slot depending on scheduling and chaining configuration. - Operators from DIFFERENT jobs are deployed to separate slots to guarantee resource isolation.
Operator Chaining (Key Optimization)
Operator chaining eliminates network serialization between lightweight transformations, significantly reducing end-to-end processing latency.
Without operator chaining: Source -> [network serialization] -> Map -> [network serialization] -> Filter -> [network serialization] -> Sink Every inter-operator boundary incurs thread context switching, buffer copy, and network serialization costs. With operator chaining: [Source -> Map -> Filter] -> [network serialization] -> [Sink] Chained operators execute within the SAME thread and pass Java object references directly in memory. Operator chaining can materially improve throughput for narrow transformation pipelines by avoiding serialization and network boundaries. Operator chaining boundaries: - keyBy(): requires hash-based network shuffling across partitions - rebalance(): requires round-robin redistribution across downstream tasks - Disparate parallelism: chaining breaks whenever adjacent operators have different parallelism configurations
Windowing (The Heart of Stream Processing)
Windowing defines how continuous event streams group into discrete buckets, balancing exactness against memory consumption.
TUMBLING WINDOW (fixed-size, non-overlapping):
[0-5min] [5-10min] [10-15min] ...
Each event belongs to exactly one discrete window.
Use case: Total completed purchase orders per 5-minute interval.
events.keyBy(e -> e.userId)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.sum("amount");
SLIDING WINDOW (fixed-size, overlapping):
[0-10min] [5-15min] [10-20min] ... (size = 10min, slide = 5min)
Each event belongs to multiple concurrent windows (size / slide = 2 windows).
Use case: Moving average of CPU utilization over the last 10 minutes, refreshed every 5 minutes.
SESSION WINDOW (gap-based, dynamic boundaries):
Events are clustered by continuous activity, and the window closes after an inactivity gap.
User A: [click, click, click] <5min gap> [click, click] <5min gap> [click]
Produces 3 distinct session windows.
Use case: User session duration and journey analysis for web applications.
events.keyBy(e -> e.userId)
.window(EventTimeSessionWindows.withGap(Time.minutes(5)))
.aggregate(new SessionDurationAggregator());
GLOBAL WINDOW (unbounded time):
Retains all events for a key in a single window, firing upon custom triggers.
Use case: Emitting an alert after every 100 consecutive error events per service.Event Time vs Processing Time vs Ingestion Time
Choosing between event time, ingestion time, and processing time determines determinism during replays and shapes the required watermark strategy.
PROCESSING TIME: Evaluated using the local wall clock of the machine processing the record. Advantages: Simple implementation that requires no watermarks or timestamp extractors. Trade-offs: Non-deterministic execution where replaying historical streams produces divergent results. Out-of-order records cause events to land in incorrect time windows. EVENT TIME: Evaluated using the original timestamp embedded in the event by the producing client or sensor. Advantages: Fully deterministic execution where replaying historical streams yields identical aggregations. Accurately handles out-of-order arrivals and network delays. Trade-offs: Requires watermark coordination and explicit policies for late data. INGESTION TIME: Evaluated using the timestamp assigned when the record first enters the streaming system (such as the ingestion broker or stream source operator). Advantages: Avoids requiring producer-side event timestamps and can be useful when event-time timestamps are completely unavailable or unreliable. Trade-offs: Does not represent the original event occurrence time and does not provide globally monotonic timestamps across distributed partitions, nor can it account for producer-side network delays or mobile offline buffering. RECOMMENDATION: Always standardize on EVENT TIME for production architectures. Reserve Processing Time strictly for local testing or for pipelines where event timestamps are completely unavailable.
Watermarks & Late Data (The Hardest Concept)
Watermarks track progress in event time, allowing streaming operators to close windows deterministically despite network delays and out-of-order arrivals.
The out-of-order problem:
Events arrive out of order due to network jitter, mobile client buffering, and distributed partitions.
An event timestamped at 10:05 might arrive after an event timestamped at 10:08.
A stream processor must determine when it is safe to close and materialize the 10:00-10:05 window.
Watermark concept:
A watermark is a monotonically advancing event-time progress signal indicating that the system considers timestamps up to W sufficiently complete according to configured watermark-generation and arrival assumptions.
When a watermark of W=10:05 arrives, downstream operators treat records with event timestamp t <= 10:05 as complete enough to trigger window evaluation.
The 10:00-10:05 event-time window can now execute its aggregation and emit results.
Because watermarks are heuristic progress signals rather than absolute physical guarantees of total stream completeness, late events exceeding the configured tolerance can still arrive and require explicit handling policies (drop, allowed lateness, side output).
Watermark generation strategies:
Bounded out-of-orderness (periodic):
Emits watermark = max_observed_timestamp - max_out_of_orderness.
Example: WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5))
If the latest observed event has timestamp 10:10, the watermark is set to 10:05.
This provides a 5-second arrival tolerance for delayed records.
Punctuated:
Emits watermarks conditioned on specific control events embedded in the stream, such as periodic heartbeats.
Handling late-arriving data (events with timestamp < current watermark):
Option 1: DROP (default behavior): Late records are permanently ignored and dropped from calculations.
Option 2: ALLOWED LATENESS: Keeps window state active for an additional grace period after the watermark passes.
Example:
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.allowedLateness(Time.minutes(1))
The window fires its initial result when the watermark passes, but accepts late records for 1 more minute, emitting an updated aggregation on each late arrival.
Option 3: SIDE OUTPUT: Routes late events to an isolated side-output stream.
Example: .sideOutputLateData(lateOutputTag)
Allows downstream jobs to log delayed data or route records to dead-letter queues for auditing.
The operational trade-off:
Small out-of-orderness tolerance: Delivers low processing latency, but increases the probability of dropping late records.
Large out-of-orderness tolerance: Accommodates high network delay and jitter, but delays window emissions cluster-wide.
Standard configurations range between 5 and 30 seconds for real-time dashboards, and multiple minutes for financial settlement jobs.State Management (What Makes Flink Special)
State management allows operators to accumulate context across events, backing high-scale aggregations with RocksDB and off-heap memory.
Keyed state (partitioned per key, per operator):
ValueState<T>: stores a single value per key (such as last known location)
ListState<T>: stores a list of elements per key (such as recent user events)
MapState<K,V>: stores key-value mappings per key (such as user feature vectors)
ReducingState<T>: stores pre-aggregated values combined via a reduce function
AggregatingState: stores custom intermediate accumulator values
Example: Count events per user over the last 1 hour:
ValueState<Long> count;
processElement(Event e) {
count.update(count.value() + 1);
}
State backends:
HashMapStateBackend (in-memory):
Advantages: Ultra-fast read and write access via Java objects on JVM heap.
Limitations: Bounded by available JVM heap and constrained by GC pauses, serialization overhead, and process memory limits, requiring full serialization on checkpoints.
Best for: Small keyed state, low-latency stateless routing, and minimal aggregation windows.
EmbeddedRocksDBStateBackend (out-of-core on local NVMe SSD):
Advantages: Supports multi-terabyte state by spilling to local SSD, performs incremental checkpoints by writing only newly flushed SST files to S3.
Limitations: Lookups can be substantially slower than heap access for some workloads because of serialization, block cache misses, and disk I/O, so performance should be benchmarked for the target workload.
Best for: Large-scale stateful pipelines, high-cardinality keys, and long retention windows.
Production recommendation for this workload: use EmbeddedRocksDBStateBackend:
- Driven by this scenario's 50 TB aggregate keyed state, high-cardinality keys, and local NVMe SSD disk capacity requirements.
- While lightweight pipelines with small in-memory state benefit from HashMapStateBackend's lower lookup latency, large/long-lived keyed state makes heap backends prone to sudden out-of-memory crashes.
- Incremental checkpoints reduce network and S3 I/O by uploading only newly generated SST files, writing approximately 5 GB rather than full 100 GB snapshots.
- Recommended tuning: set state.backend.rocksdb.block.cache-size to 256MB per task slot.Checkpointing: Exactly-Once State Consistency (Deep Dive)
Checkpointing adapts the Chandy-Lamport distributed snapshot algorithm through stream barrier alignment to guarantee consistent internal state recovery without stopping processing.
Flink's core snapshot mechanism: adapted Chandy-Lamport distributed snapshot algorithm.
Execution flow:
1. Checkpoint Coordinator in the JobManager triggers checkpoint N.
2. Coordinator injects BARRIER-N markers into all source operator input channels.
3. Sources emit the barrier inline with event streams:
[event] [event] [BARRIER-N] [event] [event] [BARRIER-N+1] ...
4. When an operator receives BARRIER-N from all input channels:
a. It initiates an asynchronous local state snapshot write to S3 (state snapshots are non-blocking and asynchronous).
b. It forwards the barrier downstream to next-hop operators.
c. It continues processing incoming events without pausing the execution pipeline.
5. Checkpoint completion occurs only after all required subtask state acknowledgments succeed and metadata is committed by the coordinator.
6. Checkpoint N becomes the latest globally consistent recovery snapshot.
Barrier alignment for exactly-once processing (Aligned Checkpoints):
Consider an operator with two input streams (Channel A and Channel B):
Channel A: [event] [BARRIER-N] [event] ...
Channel B: [event] [event] [event] [BARRIER-N] ...
BARRIER-N arrives on Channel A first: the operator temporarily blocks or pauses consumption from Channel A and buffers incoming events to prevent state divergence.
The operator continues reading from Channel B until BARRIER-N arrives on B.
Once both barriers are received, the operator takes an asynchronous state snapshot, emits BARRIER-N downstream, and resumes processing buffered records from A.
This barrier alignment ensures that state snapshots correspond exactly to the barrier boundary, preventing events from being double-counted across checkpoints.
Unaligned checkpoints (Flink 1.11+):
Operators do not block consumption or wait for barrier alignment during backpressure.
Instead, in-flight channel buffers are captured directly within the checkpoint snapshot alongside operator state, greatly reducing alignment delays.
Advantages: Eliminates processing pauses and keeps latency predictable during backpressure.
Trade-offs: Increases checkpoint artifact size because in-flight network buffers are written to S3.
Checkpoint vs Savepoint:
Checkpoint: Automatically triggered, periodic, designed for failure recovery, and optimized for minimal overhead.
Savepoint: Manually triggered, designed for version upgrades, cluster migrations, and state schema evolution.
Stateful upgrade semantics:
Savepoints preserve application state for controlled migration and upgrades with minimal planned downtime.
A conventional cancel-and-restore workflow introduces a brief maintenance gap while the existing job cancels and the new job reloads state.
Achieving true zero-downtime upgrades requires an additional operational strategy, such as running parallel blue-green jobs with coordinated downstream cutover or dual-sink writes with deduplication.
Use savepoints when upgrading application code, rescaling parallelism, or migrating clusters.
Use checkpoints for automated crash recovery during runtime worker failures.API Design
These operator-facing interfaces include pipeline definitions, declarative queries, and administrative REST endpoints that showcase how production teams manage running streaming jobs.
Stream Processing Pipeline Contract (TypeScript)
Illustrative platform-level contract modeling pipeline topologies, watermark policies, and windowed stream transformations (conceptual contract for architectural discussion, distinct from Flink's Java DataStream SDK):
// Illustrative platform-level streaming contract (for conceptual architecture modeling, distinct from Flink's Java DataStream SDK)
// Stream execution topology configuration
export interface StreamExecutionConfig {
jobName: string;
parallelism: number;
checkpointIntervalMs: number;
checkpointMode: "EXACTLY_ONCE" | "AT_LEAST_ONCE";
minPauseBetweenCheckpointsMs: number;
stateBackend: "EmbeddedRocksDB" | "HashMap";
}
// Event-time watermark extraction strategy
export interface WatermarkStrategy<T> {
maxOutOfOrdernessMs: number;
extractTimestamp: (event: T) => number;
idleTimeoutMs?: number;
}
// Windowing definition across streaming events
export interface WindowDefinition {
type: "tumbling" | "sliding" | "session" | "global";
sizeMs: number;
slideMs?: number;
sessionGapMs?: number;
allowedLatenessMs?: number;
}
// Streaming pipeline operator contract
export interface StreamOperator<TInput, TOutput> {
filter(predicate: (event: TInput) => boolean): StreamOperator<TInput, TInput>;
keyBy<TKey>(keySelector: (event: TInput) => TKey): KeyedStreamOperator<TKey, TInput>;
map<TNext>(mapper: (event: TInput) => TNext): StreamOperator<TInput, TNext>;
sink(destination: string, semantic: "EXACTLY_ONCE" | "AT_LEAST_ONCE"): void;
}
export interface KeyedStreamOperator<TKey, TValue> {
window(definition: WindowDefinition): WindowedStreamOperator<TKey, TValue>;
}
export interface WindowedStreamOperator<TKey, TValue> {
aggregate<TResult>(aggregator: (accumulator: TResult, event: TValue) => TResult): StreamOperator<TValue, TResult>;
sideOutputLateData(tag: string): WindowedStreamOperator<TKey, TValue>;
}DataStream API (Java)
Fraud detection pipeline aggregating payment failures across sliding event-time windows:
// Fraud detection: flag users with >5 failed payments in 10 minutes
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setStateBackend(new EmbeddedRocksDBStateBackend());
env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000);
KafkaSource<PaymentEvent> source = KafkaSource.<PaymentEvent>builder()
.setBootstrapServers("kafka-broker:9092")
.setTopics("payment-events")
.setGroupId("fraud-detection-group")
.setStartingOffsets(OffsetsInitializer.earliest())
.setValueOnlyDeserializer(new PaymentEventDeserializationSchema())
.build();
WatermarkStrategy<PaymentEvent> watermarkStrategy = WatermarkStrategy
.<PaymentEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, ts) -> event.getTimestamp());
DataStream<PaymentEvent> events = env.fromSource(
source,
watermarkStrategy,
"KafkaSource"
);
DataStream<FraudAlert> alerts = events
.filter(e -> e.getStatus().equals("FAILED"))
.keyBy(PaymentEvent::getUserId)
.window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)))
.aggregate(new CountAggregator())
.filter(count -> count.getValue() > 5)
.map(count -> new FraudAlert(count.getUserId(), count.getValue()));
KafkaSink<FraudAlert> sink = KafkaSink.<FraudAlert>builder()
.setBootstrapServers("kafka-broker:9092")
.setRecordSerializer(
KafkaRecordSerializationSchema.builder()
.setTopic("fraud-alerts")
.setValueSerializationSchema(new AlertSerializationSchema())
.build()
)
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("fraud-alerts-sink")
.build();
alerts.sinkTo(sink);
env.execute("Fraud Detection Job");Flink SQL
Declarative streaming aggregation using sliding HOP windows and event-time watermarks:
-- Fraud detection using sliding HOP window in SQL
CREATE TABLE payment_events (
user_id STRING,
amount DECIMAL(10,2),
status STRING,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'payment-events',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json'
);
SELECT user_id, COUNT(*) AS fail_count
FROM payment_events
WHERE status = 'FAILED'
GROUP BY user_id, HOP(event_time, INTERVAL '1' MINUTE, INTERVAL '10' MINUTE)
HAVING COUNT(*) > 5;Job Management REST API
HTTP endpoints exposed by the JobManager dispatcher for lifecycle and state control:
POST /jars/upload HTTP/1.1
# Upload compiled job JAR to the cluster dispatcher
POST /jars/{jar_id}/run?savepointPath=s3://checkpoints/savepoint-01¶llelism=16 HTTP/1.1
# Submit and execute job graph with optional savepoint restore
GET /jobs HTTP/1.1
# List all active, completed, and failed jobs
GET /jobs/{job_id} HTTP/1.1
# Retrieve runtime execution graph status, vertex metrics, and backpressure
POST /jobs/{job_id}/savepoints HTTP/1.1
# Trigger an asynchronous savepoint to durable object storage
PATCH /jobs/{job_id} HTTP/1.1
# Gracefully cancel or drain a running stream processing job
GET /jobs/{job_id}/checkpoints HTTP/1.1
# Retrieve checkpoint history, alignment durations, and snapshot sizes
GET /taskmanagers HTTP/1.1
# List registered TaskManager worker instances and available task slotsData Model
Stream processing architectures do not rely on traditional relational tables. Instead, explain how Kafka offsets, RocksDB key-group storage, and checkpoint metadata coordinate to guarantee persistence.
State Backend: RocksDB Internal Structure
Local directory hierarchy and conceptual binary key layout used by RocksDB to store keyed state on NVMe storage (illustrative internal architecture):
Per TaskManager, per keyed operator directory layout:
/tmp/flink-state/job-abc123/op-window-agg/
├── 000042.sst (Sorted String Table: immutable, sorted key-value pairs)
├── 000043.sst
├── 000044.sst
├── MANIFEST (tracks which SST files are current)
├── WAL (write-ahead log for crash recovery)
└── OPTIONS (RocksDB configuration)
Conceptual binary key layout: [key_group (2 bytes)] [key_namespace] [user_key]
Conceptual binary value layout: [serialized state value]
(Note: Illustrated as a conceptual layout, where exact internal byte prefixing, namespace delimiters, and serializers depend on Flink runtime version and TypeInformation configuration.)
Key groups: Keys are deterministically mapped to key groups (conceptually modeled as hash(key) % max_parallelism, implemented via Flink's Murmur3/key-group assignment logic).
Each parallel subtask owns a set/range of key groups:
Subtask 0: key_groups 0-31
Subtask 1: key_groups 32-63
Rescaling: redistributes key groups among subtasks without re-hashing recordsCheckpoint Metadata
Illustrative manifest schema serialized to durable S3 storage upon successful checkpoint completion (conceptual representation):
{
"checkpoint_id": 42,
"job_id": "abc123",
"timestamp": "2026-03-14T10:01:00Z",
"duration_ms": 3500,
"state_size_bytes": 5368709120,
"is_incremental": true,
"operator_states": [
{
"operator_id": "source-kafka",
"subtask_states": [
{ "subtask": 0, "offset": { "partition-0": 5000150, "partition-1": 4200300 } },
{ "subtask": 1, "offset": { "partition-2": 3100200, "partition-3": 2800100 } }
]
},
{
"operator_id": "window-aggregate",
"subtask_states": [
{ "subtask": 0, "state_handle": "s3://checkpoints/job-abc/chk-42/tm-1/sst-files/" },
{ "subtask": 1, "state_handle": "s3://checkpoints/job-abc/chk-42/tm-2/sst-files/" }
]
}
]
}Example Job: Source & Sink Topics
Topic topology and partition alignments connecting streaming ingestion to downstream datastores:
# Source topic configuration
source_topic:
name: payment-events
partitions: 128 # Partitioned by user_id for per-user ordering in session windows
event_fields:
- event_id
- user_id
- amount
- merchant_id
- timestamp
# Sink destinations matching architecture diagram
sink_destinations:
fraud-alerts:
target: Kafka (partitioned by user_id for real-time alert consumers)
fraud-scores:
target: Elasticsearch (analyst search dashboard)
daily-aggregates:
target: ClickHouse (OLAP rollups)
flagged-users:
target: PostgreSQL (serving database for blocklist API)
raw-archive:
target: S3 / Parquet (data lake bounded batch sink)
# Source partition count and operator parallelism alignment:
# Kafka partitioning by user_id preserves per-user ordering at the Kafka-source level.
# Matching 128 Kafka partitions with parallelism 128 simplifies capacity planning and prevents source-side worker imbalance.
# However, Flink keyBy(user_id) independently maps keys to Flink key groups and downstream subtasks.
# Therefore, a keyed network repartition/shuffle can still occur even when Kafka partitions and Flink parallelism match numerically.Kafka Source Offset Tracking
Atomic coordination mechanism between Kafka consumer offsets and checkpointed state:
Flink does not use Kafka internal __consumer_offsets topic for exactly-once processing. Instead: 1. Kafka partition offsets are stored directly in Flink checkpoint state. 2. During checkpointing, the operator snapshots partition-to-offset mappings alongside keyed state. 3. During recovery, Flink seeks Kafka consumers directly to the checkpointed offsets. 4. Offsets and operator state remain atomically consistent because both commit within the same snapshot epoch. 5. This atomic snapshot mechanism guarantees exactly-once stateful processing and recovery within Flink. 6. It does not automatically guarantee external exactly-once side effects, as true end-to-end guarantees require downstream transactional sinks (such as Kafka two-phase commit) or idempotent write patterns (such as PostgreSQL upserts) to ensure external consumers observe exactly-once or effectively-once results.
Fault Tolerance
Checkpoint failures, backpressure, and state recovery after worker crash are core topics.
Fault-Tolerance Failure Matrix
Recovery strategies and operational safeguards across infrastructure failures and data anomalies:
| Concern | Solution |
|---|---|
| TaskManager crash | The JobManager detects missing heartbeats, restores the job graph from the latest completed checkpoint, and redeploys subtasks to surviving or newly provisioned TaskManagers. |
| JobManager crash | A standby JobManager is elected via ZooKeeper or Kubernetes leader election, recovering job metadata and the latest checkpoint state directly from durable storage. |
| State loss | State is continuously persisted to S3 through incremental checkpoints, enabling RocksDB state to be reconstructed from SST files during recovery. |
| Slow operator (backpressure) | Credit-based flow control ensures upstream operators suspend message transmission when downstream buffers fill, preventing dropped packets and memory exhaustion. |
| Out-of-order events | Event-time watermarks combined with allowed lateness ensure that windows fire accurately while accommodating delayed records. |
| Kafka offset drift | Consumer partition offsets are persisted inside Flink checkpoint state rather than Kafka commit logs, guaranteeing exactly-once alignment between consumer offsets and operator state during recovery, while external side-effect guarantees depend on the downstream sink implementation. |
| Checkpoint failure | The coordinator retries failed snapshot attempts. If storage remains unavailable, the job continues processing but cannot recover safely from subsequent worker crashes until a snapshot succeeds. |
| Skewed keys | Invoking rebalance() redistributes unkeyed stream records evenly but does not solve a hot logical key because keyBy() immediately repartitions records by key hash back to the same subtask. Genuine hot keys are mitigated by two-stage aggregation (local pre-aggregation or key salting with sub-key shards followed by a second-stage merge), or by redesigning the partitioning key when business semantics permit. |
Recovery Flow (Detailed)
Step-by-step sequence executed by the JobManager when a TaskManager node fails:
1. TaskManager-3 crashes due to an out-of-memory error or hardware failure.
2. JobManager detects the failed worker after heartbeat/health-check timeout (configuration dependent).
3. JobManager cancels all running subtasks for the affected job graph.
4. JobManager requests new task slots from the ResourceManager (Kubernetes provisions a new pod or YARN allocates a container).
5. JobManager reads the latest completed checkpoint metadata from S3.
6. JobManager deploys all operators to available task slots across surviving and new TaskManagers.
7. Each operator subtask restores its state:
- Sources seek Kafka consumer partitions to checkpointed offsets.
- Stateful operators download RocksDB SST files from S3 and rebuild local indexes.
- Sinks abort uncommitted transactions and prepare for replayed records.
8. Event processing resumes from the exact checkpoint boundary.
Recovery time (Target RTO < 5 minutes):
State download: 100 GB state ÷ 1 GB/s network bandwidth = ~100 seconds.
Typical illustrative recovery: approximately 30-120 seconds, explicitly dependent on TaskManager provisioning, state restoration, checkpoint size, local recovery availability, event replay from Kafka, and network/infrastructure performance.
Note that 30-120 seconds is an illustrative operational estimate rather than a universal Flink guarantee.
Optimization with local recovery:
TaskManagers store secondary state snapshots on local NVMe SSDs alongside durable S3 checkpoints.
When the same TaskManager slot is reassigned, the worker avoids full remote S3 downloads, reducing the state restoration phase significantly.Exactly-Once Sinks (The Last Mile)
Protocols required across downstream datastores to ensure end-to-end exactly-once semantics:
The sink delivery challenge:
Flink processes events exactly once internally, but downstream sinks can receive duplicate writes if a worker crashes after writing but before the checkpoint completes.
Sink implementations by storage type:
Apache Kafka Sink (Two-Phase Commit Protocol):
1. During checkpointing, Flink pre-commits a Kafka transaction so writes are visible only to read_committed consumers.
2. Once the JobManager confirms all operators acknowledged the checkpoint, Flink commits the Kafka transaction, making records visible.
3. If a failure occurs before completion, Flink aborts the transaction on restart, cleanly rolling back uncommitted records.
4. This guarantees true end-to-end exactly-once delivery from source to sink.
Relational JDBC and PostgreSQL Sinks (Idempotent Upserts):
- Implements upsert operations (INSERT ... ON CONFLICT DO UPDATE) driven by deterministic primary keys.
- Upon recovery, replayed records execute identical updates without producing duplicate rows.
Amazon S3 and Object Storage Sinks (Checkpoint-Coordinated Commit Protocols):
1. In-flight records are written as in-progress or pending part objects (such as multipart upload parts or temporary staging paths).
2. When the checkpoint barrier arrives, the sink closes in-progress parts into pending committables.
3. Upon successful checkpoint notification from the JobManager, a committer operator finalizes object visibility (such as completing S3 multipart uploads or publishing manifest-committed files), without relying on POSIX-style atomic renames which S3 does not natively support.
4. On uncommitted worker failure, pending parts are discarded or aborted during recovery, ensuring downstream readers only observe fully committed checkpoint epochs. Exact commit semantics depend on whether the FileSink, Iceberg, or Delta connector is utilized.
Redis and Key-Value Sinks:
- Employs idempotent SET and HSET operations with deterministic key structures.
Elasticsearch Sink:
- Derives document IDs as a deterministic hash of the user key and window timestamp boundary.Additional Considerations
Complex event processing patterns, side outputs, and dynamic scaling are senior follow-ups.
Flink vs Spark Structured Streaming vs Kafka Streams
Detailed comparison of processing paradigms, state capacity, and runtime cluster overhead across popular engines:
| Aspect | Flink | Spark Structured Streaming | Kafka Streams |
|---|---|---|---|
| Processing model | True streaming (event-at-a-time) | Micro-batch (100ms+ intervals) | True streaming (event-at-a-time) |
| Latency | True streaming execution suitable for millisecond-scale workloads (chaining and commit dependent) | Micro-batch latency dependent on trigger/batch interval (sub-second to seconds) | Embedded event-at-a-time processing capable of low-latency delivery (poll and batch dependent) |
| State management | RocksDB (TB-scale, incremental checkpoints) | State store (limited) | RocksDB (embedded) |
| Exactly-once | State consistency via checkpoints, while end-to-end guarantees depend on sink semantics | Depends on source, checkpointing, query, and sink semantics | Via Kafka transactions when EOS is configured |
| Windowing | Rich (tumbling, sliding, session, custom) | Tumbling, sliding (limited session) | Tumbling, sliding, session |
| Event time | First-class (watermarks, late data) | Supported (watermarks) | Timestamp/stream-time based (no watermarks) |
| Deployment | Standalone, YARN, K8s | Spark cluster (YARN, K8s) | Embedded in app (no cluster) |
| SQL support | Flink SQL (full SQL) | Spark SQL (mature) | KSQL / ksqlDB |
| Best for | Complex event processing, large state | Batch + stream unification | Simple stream processing inside services |
| Cluster overhead | Full cluster (JMs + TMs) | Full Spark cluster | None (library) |
Backpressure: Credit-Based Flow Control
Credit-based buffer allocation preventing memory exhaustion during throughput spikes:
The backpressure challenge:
When a downstream operator processes records slower than upstream operators produce them, intermediate buffers fill rapidly.
Comparison of backpressure mechanisms:
- Push-based systems without flow control: buffers accumulate indefinitely until the JVM encounters an OutOfMemoryError.
- TCP-based flow control: relies on TCP window zero-advertisement, which is coarse-grained and blocks entire worker connections.
- Flink credit-based flow control:
Each downstream operator advertises available buffer credits to its upstream producers.
Upstream operators transmit records only when positive credits exist.
When credits reach zero, upstream transmission halts immediately, propagating backpressure upstream without dropped packets.
Diagnosing bottlenecks in the Flink Web UI:
Source: 100% busy (backpressured while waiting for downstream credits).
Map: 50% busy (healthy throughput).
Window Aggregation: 100% busy (the primary processing bottleneck).
Sink: 20% busy (waiting on aggregation output).
Remediation strategies:
Increase the parallelism of the bottlenecked operator, optimize window aggregation state size, or allocate additional CPU and memory resources to the affected TaskManagers.Rescaling (Changing Parallelism Without Data Loss)
Redistributing key groups across an expanded or contracted TaskManager fleet:
Scenario: Job running with parallelism=8, scaling up to parallelism=16:
1. Trigger savepoint: POST /jobs/{id}/savepoints
2. Cancel job gracefully
3. Modify job submission configuration: set parallelism=16
4. Submit job from savepoint: POST /jars/{id}/run?savepointPath=s3://checkpoints/savepoint-01
How state is redistributed:
Flink uses KEY GROUPS (configured to 128 by default).
Parallelism=8: each subtask owns 16 key groups
Subtask 0: key groups 0-15
Subtask 1: key groups 16-31
...
Parallelism=16: each subtask owns 8 key groups
Subtask 0: key groups 0-7 (subset of previous subtask 0 state)
Subtask 1: key groups 8-15 (remaining subset of previous subtask 0 state)
...
Key groups serve as the atomic unit of state redistribution across workers.
max_parallelism determines the total number of key groups and cannot be modified after job creation.
Always configure max_parallelism sufficiently high at startup (such as 128 or 256) to ensure future scale-out capability.
Operational note on downtime:
A conventional cancel-and-restore workflow from a savepoint incurs minimal planned downtime while task slots re-provision and state re-hydrates.
Achieving true zero-downtime scaling or upgrades requires running parallel blue-green jobs and cutting over downstream consumer traffic once state catch-up completes.Production Monitoring
Critical operational indicators and alert thresholds for cluster health:
critical_metrics:
checkpointDuration:
description: Alert if trending upward, indicating excessive state size or S3 throttling
lastCheckpointSize:
description: Alert if growing unbounded, signaling state retention leaks
numRecordsInPerSecond:
description: Monitors raw ingress event velocity
numRecordsOutPerSecond:
description: Monitors sink egress throughput
currentInputWatermark:
description: Stalled values indicate upstream source delays or paused partitions
busyTimeMsPerSecond:
description: Values exceeding 900ms per second indicate subtask CPU saturation
backPressuredTimeMsPerSecond:
description: Values greater than zero identify downstream processing bottlenecks
warning_metrics:
numLateRecordsDropped:
description: Tracks events discarded after exceeding allowed lateness
numberOfFailedCheckpoints:
description: Tracks consecutive snapshot failures against durable storage
fullRestarts:
description: Monitors worker crash and recovery frequency
managedMemoryUsage:
description: Tracks RocksDB off-heap cache and buffer pool consumption
alert_rules:
checkpoint_duration_p99:
condition: "> 2 * checkpoint_interval"
severity: critical
consumer_lag:
condition: "increasing monotonically over 5 minutes"
severity: high
watermark_lag:
condition: "> 5 minutes"
severity: highCommon Anti-Patterns
Frequent architectural mistakes encountered in production stream processing deployments:
1. Using processing time when event time is available Causes non-deterministic results and incorrect window assignments when events arrive out of order. 2. Not setting max_parallelism at job creation Defaults to 128 key groups. If the workload later requires a parallelism of 200, rescaling becomes impossible without reprocessing history. 3. Storing large composite objects in ValueState instead of MapState Every read or write deserializes the entire object. Use MapState to mutate individual nested keys efficiently. 4. Not enabling incremental checkpoints with RocksDB A 100 GB state uploads 100 GB to S3 every minute, causing storage costs and network bandwidth to surge. Incremental snapshots upload only 1-5 GB of newly created SST files. 5. Applying keyBy on a high-cardinality field without accounting for key skew If a single hot key receives 90% of events, a single subtask performs 90% of the work and starves the cluster. Invoking rebalance() before keyBy() does not resolve this because keyBy() immediately re-hashes records back to the same subtask. Mitigate genuine hot keys using local pre-aggregation, sub-key sharding, or key salting (appending random prefixes or shard identifiers), explicitly running a second aggregation stage to combine the salted partial results into the final logical key. 6. Ignoring backpressure signals in operational dashboards The job appears healthy because throughput matches ingress, but the sink is 100% saturated and buffering data, driving processing latency from milliseconds to minutes. Always monitor the Flink Web UI backpressure indicators.
Related Problems
Stream processing patterns explored here apply directly to event analytics pipelines in User Analytics Pipeline and sliding-window fraud evaluation in Fraud Detection System. Foundational messaging, state distribution, and storage primitives connect to Kafka Architecture and Guarantees, Stream Processing Basics, Message Queues Fundamentals, and CAP Theorem and Consistency Models.
Interview Walkthrough
- 25-minute cut
Skip RocksDB internals unless staff.
- Event-time vs processing-time + watermarks (5 min)
- at-least-once vs exactly-once + idempotent sinks (7 min)
- tumbling/sliding/session windows with one example each (8 min)
- checkpointing sketch (5 min)
- Open with the distinction between event time and processing time, because explaining watermarks and late-arriving data provides the primary signal of senior-level depth.
- Diagram the end-to-end data flow from source to Kafka, through the stream processor, and into downstream sinks, clearly labeling where at-least-once vs exactly-once semantics are required.
- Explain checkpointing and idempotent sinks as the practical path to exactly-once without two-phase commit everywhere.
- Compare tumbling, sliding, and session windows with a concrete use case for each (billing, metrics, user sessions).
- Mention backpressure and credit-based flow control when downstream sinks (DB, API) cannot keep pace with ingress.
- Cover state store sizing with RocksDB and changelog topics for fault-tolerant keyed state.
- Avoid the common pitfall of treating stream processing as micro-batching with smaller files, and explicitly address how to handle unbounded state accumulation and out-of-order event streams.
Engineering Trade-offs
Your interviewer will push on framework choice and delivery semantics. Walk through Flink vs Spark Streaming and state what you'd pick for sub-100ms event processing.
Flink vs Kafka Streams vs Spark Structured Streaming
Evaluating stream processing engines based on deployment topologies, state backing models, and latency guarantees:
| Dimension | Apache Flink | Kafka Streams | Spark Structured Streaming |
|---|---|---|---|
| Deployment | Separate cluster (JobManager/TaskManager) | Embedded in app (library, no cluster) | Separate cluster (Driver + Executors) |
| State management | Native RocksDB state backend, large state (TB scale) | RocksDB state stores per partition | In-memory + checkpoints (less suited for huge state) |
| Event-time processing | First-class (watermarks, out-of-order handling) | Timestamp and stream-time based, supporting event-time windows without using Flink-style watermarks | Event-time via watermarks (added in Spark 2.1) |
| Exactly-once | State consistency via checkpoints, where end-to-end guarantees depend on sink semantics | Via Kafka transactions when EOS is configured | Depends on source, checkpointing, query, and sink semantics |
| Latency | True streaming execution suitable for millisecond-scale workloads (chaining and commit dependent) | Embedded event-at-a-time processing capable of low-latency delivery (poll and batch dependent) | Micro-batch latency dependent on trigger and batch intervals ranging from sub-second to seconds, whereas continuous processing has distinct operational constraints |
| Operational complexity | High (separate cluster to manage) | Low (just a library, no infra) | Medium (Spark cluster) but familiar to batch teams |
Best For:
- Apache Flink: Best suited for complex event processing, large stateful stream joins, and millisecond latency SLAs.
- Kafka Streams: Ideal for lightweight per-event enrichment, filtering, and microservices embedded directly within the Kafka ecosystem.
- Spark Structured Streaming: Optimal for engineering teams already maintaining Spark clusters who require unified batch and streaming pipelines.
Decision Framework:
- Consider Spark Structured Streaming when batch and stream processing unification and existing Spark operational expertise outweigh ultra-low-latency streaming requirements.
- Prefer Apache Flink when complex stateful processing, rich event-time watermark semantics, large-scale state management, or sub-second streaming latency requirements make its model advantageous.
- Prefer Kafka Streams when lightweight stream transformations embedded directly within Kafka-centric microservices and operational simplicity without a dedicated cluster are the priorities.
Tumbling vs Sliding vs Session Windows
Architectural breakdown of time-windowing models, memory footprints, and late-arrival handling rules:
Window type determines how events are grouped for aggregation:
Tumbling Window (non-overlapping, fixed size):
|──── 1min ────|──── 1min ────|──── 1min ────|
Events in window [0-60s] are aggregated and emitted at T=60s.
Use case: Revenue per minute, error counts per 5-minute interval.
Properties:
- Each event belongs to exactly one window interval.
- Generally requires less concurrent window state than overlapping sliding windows because state is cleared upon window evaluation (plus allowed lateness).
- Window assignment semantics avoid double-counting across window intervals when configured correctly without unhandled late-data retractions.
Sliding Window (overlapping, fixed size):
|──── 5min ────|
|──── 5min ────|
|──── 5min ────|
Window size = 5 minutes, slide = 1 minute: emits the last 5 minutes of data every minute.
Use case: Moving average of CPU load over the last 5 minutes.
Properties:
- Each event belongs to (window_size / slide) windows (5 concurrent windows in this example).
- Concurrent window state can scale roughly with (window_size / slide) (5x in straightforward implementations), though actual memory depends on pre-aggregation, state representation, triggers, and cleanup policies.
- Yields smoother rolling trends with minimal boundary edge effects.
Session Window (activity-based, dynamic size):
User activity: [click]--5s--[click]--2s--[click] (20s gap) [click]--1s--[click]
Session 1: 3 events (gap between events is less than the 10s threshold).
Session 2: 2 events (new session opens after 20s of inactivity).
Use case: User session duration, burst activity detection.
Properties:
- Closes after N seconds of observed inactivity (gap threshold).
- Dynamic size: brief interactions produce short windows, while continuous activity expands window duration.
- Cannot pre-allocate fixed memory due to unbounded session durations.
Late data handling across window types:
Watermark: a monotonically advancing event-time progress signal indicating that the pipeline considers events up to the watermark threshold complete enough to trigger window evaluation.
Allowed lateness: retains window state for an additional grace period after the watermark passes.
Side output: routes late events arriving after allowed lateness into a separate stream for auditing and batch reconciliation.
Production setting: choose a bounded-out-of-orderness tolerance such as 5 minutes based on the observed event-delay distribution, where the watermark advances approximately as max observed event time minus that tolerance (exact tolerance is workload-dependent and derived from observed lateness percentiles vs latency SLO requirements).
Trade-off: results are delayed by the chosen tolerance to prioritize higher data completeness over immediacy, though records exceeding that bound still require late-data handling.Exactly-Once in Stream Processing: Cost vs. Necessity
Performance overhead, failure recovery guarantees, and business justification across processing semantics:
Exactly-once guarantees introduce measurable operational trade-offs:
At-Least-Once (default, lower overhead):
- Checkpoints persist source partition offsets and whatever operator state is required for failure recovery, but omit barrier alignment coordination. Operators process records immediately without waiting for alignment across input channels.
- On failure, replay from checkpoint and unaligned channel state can cause in-flight records or downstream side effects to be processed more than once.
- Sinks receive duplicate writes and must implement idempotent deduplication.
- Throughput runs at baseline performance with near-zero barrier alignment wait time.
Exactly-Once (Flink barrier snapshot + transactional sinks):
- Two-phase commit protocol coordinates Flink checkpoints with transactional sinks.
- Flink checkpoints align barriers and flush in-flight buffers to persistent storage.
- On recovery, the pipeline replays only uncommitted records within the failed checkpoint epoch.
- Introduces measurable overhead that depends on checkpoint frequency, state size, transaction and sink behavior, and workload characteristics.
- Checkpoint interval trade-off:
Short (1 second): minimal reprocessing on failure, but higher CPU and storage coordination overhead.
Long (10 seconds): lower background overhead, but up to 10 seconds of records must replay during recovery.
End-to-end exactly-once prerequisites:
1. Source must be replayable (such as Apache Kafka or AWS Kinesis, whereas raw HTTP webhooks are non-replayable).
2. Stream transformations must be deterministic (avoid random generators or wall-clock timestamps in business logic).
3. Sink must support two-phase transactions or native idempotency:
- Idempotent sinks: Elasticsearch (deterministic document ID), Redis (SET key), PostgreSQL (idempotent upserts via INSERT ... ON CONFLICT DO UPDATE).
- Transactional sinks: Kafka (transactional producer with two-phase commit), Amazon S3 (checkpoint-coordinated committer/FileSink protocol), PostgreSQL (two-phase commit via XA/JDBC where supported).
- Non-idempotent sinks: raw SMTP email dispatch, unkeyed HTTP webhooks, blind counter increments.
When to avoid exactly-once overhead:
- Real-time operational dashboards where a ±0.1% variance is acceptable.
- High-volume log analytics where occasional duplicate lines cause no operational harm.
- Machine learning feature extraction where minor data drift does not degrade model inferences.
When exactly-once is mandatory:
- Financial ledger aggregation and account balance calculations.
- Inventory reservation and decrement streams.
- Telecom billing and cloud resource metering pipelines.Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.