Learn Labs
12. Stream Processing

12.5 Fault Tolerance

Why batch's approach doesn't transfer:

Batch fault tolerance works because INPUT FILES ARE IMMUTABLE, EACH TASK WRITES TO A SEPARATE FILE, AND OUTPUT IS MADE VISIBLE ONLY WHEN A TASK COMPLETES SUCCESSFULLY. "IT APPEARS AS THOUGH EVERY INPUT RECORD WAS PROCESSED EXACTLY ONCE. Although restarting tasks means records MAY BE PROCESSED MULTIPLE TIMES, THE VISIBLE EFFECT IN THE OUTPUT IS AS IF THEY HAD BEEN PROCESSED ONLY ONCE."

This is EXACTLY-ONCE SEMANTICS — "although EFFECTIVELY-ONCE WOULD BE A MORE DESCRIPTIVE TERM."

"WAITING UNTIL A TASK IS FINISHED BEFORE MAKING ITS OUTPUT VISIBLE IS NOT AN OPTION, BECAUSE A STREAM IS INFINITE, SO YOU CAN NEVER FINISH PROCESSING IT."

5.1 Microbatching and checkpointing

ApproachSystemMechanism
MicrobatchingSpark StreamingBreak the stream into small blocks and treat each like a miniature batch process. Batch size typically ~1 second — "a performance compromise: SMALLER batches incur greater SCHEDULING AND COORDINATION OVERHEAD, while LARGER batches mean A LONGER DELAY before results become visible." ⚠️ "Microbatching IMPLICITLY PROVIDES A TUMBLING WINDOW EQUAL TO THE BATCH SIZE (windowed by PROCESSING time, not event timestamps); jobs requiring larger windows must EXPLICITLY CARRY OVER STATE from one microbatch to the next"
CheckpointingApache FlinkPeriodically generate ROLLING CHECKPOINTS of state and write them to durable storage. If an operator crashes, RESTART FROM THE MOST RECENT CHECKPOINT and DISCARD ANY OUTPUT generated between the checkpoint and the crash. Checkpoints are triggered by BARRIERS in the message stream — similar to microbatch boundaries, BUT WITHOUT FORCING A PARTICULAR WINDOW SIZE

⚠️ THE CRITICAL LIMIT: "Within the confines of the framework, these provide the same exactly-once semantics as batch processing. HOWEVER, AS SOON AS OUTPUT LEAVES THE STREAM PROCESSOR — when it writes to a database, publishes to an external broker, or TRIGGERS THE SENDING OF EMAILS — THE FRAMEWORK IS NO LONGER ABLE TO DISCARD THE OUTPUT OF A FAILED MICROBATCH. Restarting causes the EXTERNAL SIDE EFFECT TO HAPPEN TWICE."

5.2 Atomic commit revisited

To give the appearance of exactly-once processing, ALL OUTPUTS AND SIDE EFFECTS MUST PERSIST IF AND ONLY IF THE PROCESSING IS SUCCESSFUL. That includes:

  • messages sent to downstream operators or external messaging systems (including EMAIL OR PUSH NOTIFICATIONS)
  • database writes
  • changes to operator state
  • acknowledgments of input messages (INCLUDING MOVING THE CONSUMER OFFSET FORWARD)

These must all happen ATOMICALLY.

Why this works where XA didn't (Ch 8 §5.4):

"In more RESTRICTED environments it is possible to implement such an atomic commit facility EFFICIENTLY. Used in Google Cloud Dataflow, VoltDB, and Apache Kafka. UNLIKE XA, THESE IMPLEMENTATIONS DO NOT ATTEMPT TO PROVIDE TRANSACTIONS ACROSS HETEROGENEOUS TECHNOLOGIES, but instead KEEP THE TRANSACTIONS INTERNAL by managing BOTH STATE CHANGES AND MESSAGING WITHIN THE FRAMEWORK. THE OVERHEAD OF THE TRANSACTION PROTOCOL CAN BE AMORTIZED BY PROCESSING SEVERAL INPUT MESSAGES WITHIN A SINGLE TRANSACTION."

5.3 Idempotence — the cheaper route

"An IDEMPOTENT operation is one you can perform MULTIPLE TIMES and it has THE SAME EFFECT as if you performed it ONCE. Deleting a key is idempotent; INCREMENTING A COUNTER IS NOT."

"Even if an operation is NOT NATURALLY idempotent, IT CAN OFTEN BE MADE IDEMPOTENT WITH A BIT OF EXTRA METADATA. When consuming from Kafka, every message has a persistent, monotonically increasing OFFSET. WHEN WRITING TO AN EXTERNAL DATABASE, INCLUDE THE OFFSET OF THE MESSAGE THAT TRIGGERED THE LAST WRITE WITH THE VALUE. Thus you can tell whether an update HAS ALREADY BEEN APPLIED and avoid performing it again."

Four assumptions this relies on — all of them load-bearing:

  1. Restarting a failed task must Replay the same messages in the same order (a Log-based broker does this — an AMQP/JMS one does Not, per §1.4)
  2. The processing must be Deterministic
  3. No other node may concurrently update the same value
  4. When failing over between nodes, Fencing may be required (Ch 9 §5.2) to prevent interference from a node Thought to be dead but actually alive

“Despite all those caveats, Idempotent operations can be an effective way of achieving exactly-once semantics with only A small overhead.”

5.4 Rebuilding state after a failure

OptionDetail
Remote datastore, replicated"Having to QUERY A REMOTE DATABASE FOR EACH INDIVIDUAL MESSAGE CAN BE SLOW"
Local state, replicated periodically ✔On recovery, "the new task can READ THE REPLICATED STATE AND RESUME PROCESSING WITHOUT DATA LOSS"
FlinkPeriodic SNAPSHOTS of operator state written to durable storage (a DFS)
Kafka StreamsReplicates state changes by sending them to A DEDICATED KAFKA TOPIC WITH LOG COMPACTION — similar to CDC
VoltDBReplicates state by REDUNDANTLY PROCESSING EACH INPUT MESSAGE ON SEVERAL NODES (Ch 8's serial execution)
Rebuild from the input"Replicating the state may not even be NECESSARY. If the state is aggregations over a FAIRLY SHORT WINDOW, it may be fast enough to simply REPLAY THE INPUT EVENTS. If the state is a local replica of a database maintained by CDC, the database can be REBUILT FROM THE LOG-COMPACTED CHANGE STREAM"

"All this depends on the performance characteristics of the underlying infrastructure. In some systems, NETWORK DELAY MAY BE LOWER THAN DISK ACCESS LATENCY, and network bandwidth may be COMPARABLE TO DISK BANDWIDTH. NO SOLUTION IS UNIVERSALLY IDEAL, and the merits of LOCAL VERSUS REMOTE state MAY ALSO SHIFT AS STORAGE AND NETWORKING TECHNOLOGIES EVOLVE."


On this page