arrow
Return

Adaptive Fragment-Based Parallel State Recovery for Stream Processing Systems

delete2023-08-01
delete1
PRE
AI
H
Hailu Xu
P
Pinchao Liu
D
Dilma Da Silva
DOI:10.1109/TPDS.2023.3251997delete
deleteOriginal
deleteOriginal request for help
deleteShare
deleteSave
Abstract

Abstract

En 中文
Today, large-scale cloud organizations are deploying datacenters and edge clusters globally to provide low-latency access to services. Running stream applications across geo-distributed sites are emerging as a daily requirement. However, existing efforts have dominantly centered around stateless stream processing, leaving another urgent trend-stateful stream processing-much less explored. A driving need is to store and update states during processing, and most importantly, successfully recover large distributed states when faults and failures happen. Existing studies exhibit major limitations including: (1) they mostly inherit MapReduce's single master/many workers architecture, where the central master can easily become ascalability bottleneck; (2) they offer state recovery mainly through three approaches: replication recovery, checkpointing recovery, and DStream-based lineage recovery, which are either slow, resource-expensive or failing to handle multiple failures; and (3) they are not adaptive to heterogeneous hardware settings. We present A-FP4S, a novel adaptive fragments-based parallel state recovery mechanism for stream processing systems. A-FP4S organizes stream operators into a distributed hash table based peer-to-peer overlay and divides each node's local state into many fragments. These fragments are periodically stored in node's multiple neighbors, ensuring different sets of available fragments can reconstruct failed states in parallel. This mechanism is extremely scalable to the lost state, significantly reduces failure recovery time, and can tolerate multiple node failures. A-FP4S is adaptive to heterogeneous hardware settings by automatic parameter tuning over phases. Compared to Apache Storm, A-FP4S achieves 31.8% to 50.5% reduction in recovery latency. Large-scale experiments using real-world datasets demonstrate A-FP4S's attractive scalability and adaptivity properties.
Keywords:
Distributed hash tables
state recovery
stream processing

Journal

IEEE Transactions on Parallel and Distributed Systems cover
IEEE Transactions on Parallel and Distributed Systems
IF:
6
Papers:
5.2K
Citations:
1.1W

Organization

California State University System cover
California State University System
Scholars:
2.8W
Papers: 2.4W
Citations: 457
State University System of Florida cover
State University System of Florida
Scholars:
12.7W
Papers: 10.9W
Citations: 130
California State University, Long Beach cover
California State University, Long Beach
Scholars:
1.2K
Papers: 968
Citations: 2.1K
F
Florida International University
Scholars:
7.3K
Papers: 5.9K
Citations: 1.1W
researcher View more organizations