Fugue: Online Elasticity for Distributed Stateful Stream Processing
Yuqiu Zhang, Yunhao Mao, Hans-Arno Jacobsen
Abstract
Stateful stream processing engines are critical for real-time analytics but lack efficient mechanisms for runtime elasticity. The dominant "stop-the-world" model, used by systems like Apache Flink, requires halting applications globally for a long time, while recent on-the-fly protocols introduce severe trade-offs: proactive approaches impose a continuous resource tax by constantly replicating state, and existing reactive solutions suffer from architectural complexity and external dependencies. This paper introduces Fugue, a novel, self-contained reactive protocol that provides seamless and resource-efficient elasticity. The core of Fugue is a two-phase design that combines a pre-emptive background state transfer with an atomic, lightweight barrier-based cutover. By moving the bulk of an operator's state off the critical path and unifying the final ownership transfer with the system's native exactly-once synchronization mechanism, Fugue guarantees correctness with minimal disruption and steady-state overhead. We implemented Fugue in Apache Flink and our evaluation on realistic benchmarks shows it reduces tail reconfiguration latency by up to 98.6% relative to native Flink while maintaining over 90% of peak throughput. Compared to reactive pull-based baselines, Fugue reduces end-to-end migration latency by up to 93.7%. Compared to proactive replication, it reaches comparable handover performance while avoiding continuous replication overhead. Together, these results demonstrate a strong combination of robustness, performance, and operational simplicity.
Ask about this paper
Your agent reads all of it.
Lune indexed this paper to the last equation, along with the top-tier papers that cite it. Ask a question and the answer quotes them.
Your agent calls
Luneget_paper_fulltext
Free to start. No credit card required.
Terminal
Install the CLIlune papers fulltext 13c99751-2135-4b27-b5cb-de41d86ac8abBuilds on3
- Rhino: Efficient Management of Very Large Distributed State for Stream Processing EnginesBonaventura Del Monte, Steffen Zeuch, Tilmann Rabl, Volker MarklSIGMOD 2020 · 56 citations
- Meces: Latency-efficient Rescaling via Prioritized State Migration for Stateful Distributed Stream Processing SystemsRong Gu, Han Yin, Weichang Zhong, Chunfeng Yuan et al.USENIX ATC 2022 · 22 citations
- Towards Fine-Grained Scalability for Stateful Stream Processing SystemsYunfan Qing, Wenli ZhengICDE 2025 · 2 citations
Related papers
- StreamSwitch: Fulfilling Latency Service-Layer Agreement for Stateful StreamingZhaochen She, Yancan Mao, Hailin Xiang, Xin Wang et al.INFOCOM 2023 · 5 citations
- Scabbard: Single-Node Fault-Tolerant Stream ProcessingGeorgios Theodorakis, Fotios Kounelis, Peter R. Pietzuch, Holger PirkVLDB 2022 · 21 citations
- Latency-Oriented Elastic Memory Management at Task-Granularity for Stateful Streaming ProcessingRengan Dou, Richard T. B. MaINFOCOM 2023 · 2 citations
- Learning from the Past: Adaptive Parallelism Tuning for Stream Processing SystemsYuxing Han, Lixiang Chen, Haoyu Wang, Zhanghao Chen et al.ICDE 2025 · 2 citations
- Sluice: End-to-End Latency Guarantee for Long-running Dataflow SystemsZhaochen She, Yancan Mao, Richard T. B. MaINFOCOM 2026
