Skip to content

Fugue: Online Elasticity for Distributed Stateful Stream Processing

Jul 2026 · Proceedings of the VLDB Endowment · Vol 19, pp. 3288-3301 · 0 citations · 58 references

TL;DR

Fugue is a novel, self-contained reactive protocol that provides seamless and resource-efficient elasticity and reaches comparable handover performance while avoiding continuous replication overhead, and is implemented in Apache Flink.

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.

View source

Similar papers

Open access Sep 2026

Where Does Streaming State Cost Go? A Reproducible Comparison of Flink and Kafka Streams on Kafka

This paper presents a controlled comparison of exactly-once Kafka pipelines implemented with Apache Flink and Kafka Streams, two engines with different state-management architectures. The state management in Flink occurs through checkpoints in external storage, whereas Kafka Streams restores local state by replaying br...

Kiran N. Kumar, Santhoshkumar Saminathan · 0 citations
Preprint Sep 2026

Ermes: a Stateful Serverless Platform for the Edge-to-Cloud Continuum

Function-as-a-Service (FaaS) is a widely adopted paradigm to simplify application deployment across the edge-to-cloud continuum. However, its stateless nature forces functions to retrieve their state from external, typically cloud-centric, data stores, reintroducing the very latency that edge computing aims to eliminat...

Matteo Cenzato, Dario d'Abate, Arianna Dragoni et al. · 0 citations
Preprint Sep 2026

No-Restart Elasticity in an Adaptive Runtime System for Cloud-Native HPC

Exploiting discounted spot instances for HPC requires an application to change its resource allocation at runtime, shrinking ahead of an interruption and expanding onto replacement capacity. Existing elasticity mechanisms implement rescaling as a full process teardown followed by a cold restart at the new processor cou...

Aditya Bhosale, Laxmikant V. Kalé · 0 citations
Book Open access Sep 2026

TuxBot: Semantic-Aware Online OS Tuning with LLMs

Online OS tuning can improve long-running services, but existing tuners are not well suited for live hosts. They treat scheduler, power, memory, and I/O controls as black-box variables and optimize a scalar reward. This approach ignores cross-knob policy structure, breaks down when application metrics are unavailable,...

Georgios Liargkovas, M. Joshi, Hubertus Franke et al. · 0 citations
Open access Sep 2026

Serverless Data Engineering: Innovations in Python-Driven ETL Automation on AWS

Traditional cluster-based ETL architectures impose a structural tax on data engineering organisations: fixed compute resources provisioned for peak demand, scheduled batch cycles that introduce latency regardless of downstream urgency, and operational overhead that redirects engineering capacity from pipeline design to...

Rambabu Bolineni · 0 citations

We use cookies to run the site and, with your consent, for analytics and to show ads. See our Cookie Policy.