Back to results

University of Illinois at Urbana-Champaign

FastRecover: simple and effective fault recovery in a distributed operator-based stream processing engine

Abstract

dc:description

Fault tolerance is a key requirement in large-scale distributed stream processing engines (SPEs), especially those that run atop commodity hardware. Currently, fault tolerance in popular distributed SPEs is either inadequate (e.g., those without automatic recovery of operator states) or complex and inefficient (e.g., those with transactional semantics). There are two major considerations in the design of an effective fault tolerance mechanism: the overhead of additional checkpointing operations during normal processing, and the time required to recover and return to normal processing when a failure happens. The main challenge lies in that faster recovery requires higher checkpointing overhead, and vice versa. This thesis presents FastRecover, a novel fault tolerance mechanism for distributed SPEs that strikes a balance between recovery time and checkpointing overhead. Specifically, given an application topology consisting of interconnected operators, and an upper bound on checkpoint overhead, FastRecover computes the optimal expected recovery time, as well as the strategy used for checkpointing and recovery in each operator. The main idea of FastRecover is to compute an optimal partitioning of the streaming operator topology into independent segments; for each segment, FastRecover backs up its input tuples and periodically checkpoints the states of operators therein. During recovery for a particular segment, FastRecover restores each affected operator state in the segment to the latest checkpoint, and replays the inputs of the segment since then. Both checkpointing and recovery utilize the parallel processing capabilities of the distributed SPE. Extensive experiments demonstrate that FastRecover achieves an average of 50% reduction in expected recovery time compared to simple solutions. The experiments also show that the total expected recovery time varies proportionally to the total computational recovery time and recovery latency in tests with simulated failures, and hence is a good measure to optimize.

Degree

thesis:*
Name thesis:degree_name
M.S.
Level thesis:degree_level
Thesis
Discipline thesis:degree_discipline
Computer Science
Grantor
University of Illinois at Urbana-Champaign
Year dc:date
2016

Author and committee

dc:creator, dc:contributor.*
Author dc:creator
  • Yaduvanshi, Shashank
Contributors dc:contributor
  • Winslett, Marianne

Subjects

dc:subject × 2

Rights

dc:rights
Statement dc:rights
  • Copyright 2016 Shashank Yaduvanshi
Language dc:language
en

Identifiers

dc:identifier.*
Handle dc:identifier
http://hdl.handle.net/2142/90832
OAI identifier oai:identifier
oai:www.ideals.illinois.edu:2142/90832

Chain of custody

source
Harvested from
University of Illinois - Urbana-Champaign
Base URL
www.ideals.illinois.edu/oai-pmh
Last updated
2026-07-22
Source record
OAI-PMH GetRecord
citation

Yaduvanshi, Shashank. FastRecover: simple and effective fault recovery in a distributed operator-based stream processing engine. Thesis thesis, University of Illinois at Urbana-Champaign, 2016. http://hdl.handle.net/2142/90832