Papers for

distributed systems developers

Papers whose findings have a practical use for this group, as judged from the abstract. Open a paper to read what it means in practice.

Simple shortest-path trees give good low-stretch spanning trees

Low-Stretch Spanning Trees via Smoothed Analysis of Dijkstra's Algorithm

Abstract: Given an undirected weighted graph $G$, a $γ$-approximate low-stretch spanning tree (LSST) $T \subseteq G$ is a tree that approximates the distance metric of $G$ up to a $γ$-factor in expectation. Currently, existing algorithms to find a provably good LSST carefully construct an approximate shortest-path tree from an arbitrary source. The resulting algorithms are intricate. In contrast, practitioners observed that a much simpler heuristic performs surprisingly well: choose an arbitrary root, run Dijkstra's algorithm, and use the resulting shortest-path tree as an LSST. In this paper, we give a smoothed analysis of shortest-path tree algorithms, such as Dijkstra's algorithm, that explains this behavior. We show that adding a small perturbation to the weights of the input graph suffices to turn the shortest path tree rooted at an arbitrary node in the resulting graph into an $\tilde{O}(1)$-approximate LSST. We further show that the set of perturbations can be computed efficiently from few low-diameter decompositions (LDDs). Thus, our proof is also constructive in the sense of giving a novel approach to computing LSSTs.

Mon 28 SeptData Structures and Algorithms
The gist
Finding a tree inside a network that keeps distances roughly the same is usually complicated. The authors explain why simply picking any starting point and running a shortest-path algorithm like Dijkstra’s often works well in practice. They show that if you slightly adjust the network's edge weights, this simple method reliably finds a good approximation. Their approach also gives a new way to efficiently compute such trees.
Open → 2609.35136v1

Signature-free blockchain methods cut consensus delays significantly

Simple and Fast Signature-Free Blockchain Consensus

Abstract: Signature-free protocols avoid the cost of post-quantum signatures. We present two simple signature-free blockchain consensus protocols for eventual synchrony with optimal good-case commit latency (three message delays for \(f<n/3\) and two for \(f<n/5\)), optimistic responsiveness, a block time of only two message delays without speculation, and \(O(n^2)\) communication per view. They instantiate Generic Simplex, a blockchain consensus construction parameterized by a new abstraction called view agreement. The same construction also captures Simplex, Minimmit, and a new synchronous signature-free protocol with optimal good-case commit latency of two message delays for \(f<n/4\).

Sat 26 SeptDistributed, Parallel, and Cluster Computing
The gist
Blockchain systems need to agree on the order of transactions, which often requires complex digital signatures that could be vulnerable to quantum computers. The authors present new ways to reach consensus without these signatures, making the process faster while keeping security guarantees under certain network conditions. Their methods reduce the number of communication steps needed to confirm transactions and work efficiently even when some participants behave badly. These protocols use a new concept called view agreement to coordinate participants reliably.
Open → 2609.32985v1

HyperParallel-FSDP speeds up training of massive AI models on Ascend SuperPods

HyperParallel-FSDP: Topology-Aware Fully Sharded Training with Layout-Driven Muon on Ascend SuperPods

Abstract: Declarative SPMD programming uses tensor sharding descriptions to drive distributed execution, separating parallelization from model code. However, the evaluated PyTorch DTensor stack dispatches every operator below autograd, incurring repeated dispatch and metadata costs, while lacking an inexpensive end-to-end validation path. Existing FSDP and distributed Muon implementations also mismatch two-tier supernode topologies: FSDP relies on explicit parameter packing and unpacking, and Muon's whole-matrix orthogonalization conflicts with parameter sharding. We observe that distributed tensors need only express sharding semantics at the tensor API boundary above autograd, allowing differentiation and kernels to operate on plain tensors. Based on this insight, we present HyperParallel-FSDP, featuring: (1) dual-mode DTensor execution, using one sharding plan for both a production mode with one-time layout resolution and no steady-state dispatch overhead, and a validation mode with end-to-end metadata propagation, fail-fast checks, and gradient-equivalence testing; (2) topology-aware FSDP, with zero-copy intra-supernode collectives, fused inter-supernode reduction, and a cross-layer backward pipeline that avoids waits on slow links; and (3) layout-driven distributed Muon, with sharding-derived communication groups, deduplicated orthogonalization, and shape-fused Newton-Schulz iterations. On Atlas 900 A3 SuperPoD, HyperParallel-FSDP scales from 16 dies to 384 cards (768 ranks), sustaining 421k tokens/s for a 505B-parameter MoE while FSDP communication uses 2.9% of step time. It reduces mean step time by 29.7% versus PyTorch FSDP2 and 25.5% versus Megatron DDP, with Pearson correlation above 0.999997 over 1,000 steps. Distributed Muon improves profiler step time by 5.4-16.0% over competing systems. Source code is available at https://atomgit.com/mindspore/hyper-parallel.

Fri 18 SeptDistributed, Parallel, and Cluster Computing
The gist
Training very large AI models requires splitting them across many processors efficiently. The authors noticed inefficiencies in existing tools that slow down communication and checks during training on specialized Ascend hardware. They designed HyperParallel-FSDP, which organizes data and computation to reduce overhead and improve validation without repeating costly steps. Their system runs much faster than other approaches, handling models with over 500 billion parameters while using fewer communication resources.
Open → 2609.21594v1

Cloud native MPI enables selective process relocation for resilient computing

XMPIaaS: Towards Cloud Native MPI via Cooperative Process Migration

Abstract: Message Passing Interface (MPI) has been the dominant programming model for High Performance Computing (HPC) for three decades, and as HPC workloads increasingly migrate to cloud infrastructure for scalability and cost efficiency, MPI applications must contend with an execution environment fundamentally unlike traditional supercomputers: ephemeral resources, dynamic pricing and preemptable instances. In such a volatile setting, the ability to relocate running MPI processes between nodes without restarting the job is a necessity for cost-effective, resilient execution. Existing approaches either require restarting the entire job from a global checkpoint, or transparently intercepting the full MPI stack at prohibitive complexity. To address these challenges, we propose \name, a cooperative migration system for MPI that enables selective process group migration on-the-fly. When a cloud instance is scheduled for preemption, only the affected ranks are relocated while the remaining processes briefly quiesce and resume in place, avoiding the cost of a full-job checkpoint. \name tackles this through a cooperative protocol between the MPI process management runtime and rank processes. We expose an \texttt{XMPI\_quiesce} interface built atop the MPI Sessions API that allows applications to mark safe migration points, and we extend the Hydra process manager to orchestrate the full migration lifecycle: rank quiescence, CRIU checkpoint/restore, proxy relaunch on the target node, and seamless rank reconnection. We evaluate and show that the cooperative quiesce phase accounts for less than 1.4\% of total migration downtime, and that this downtime is governed by the migrating node's rank count alone, independent of job size, and the instrumentation introduces no measurable overhead during normal execution.

Tue 15 SeptDistributed, Parallel, and Cluster Computing
The gist
High Performance Computing (HPC) programs use MPI to communicate between many computers working together. Moving these programs to the cloud is tricky because cloud computers can stop suddenly or change often. The authors created a system called XMPIaaS that lets parts of these programs move between computers without stopping the whole program. This helps programs keep running smoothly even if some cloud computers are taken away. Their system works efficiently and doesn’t slow down normal program operation.
Open → 2609.16531v1

Sublinear size coalitions can control outputs in complex systems

A Resolution of Friedgut's Conjecture on Influential Coalitions

Abstract: We prove that, for every constant $\varepsilon>0$ and every function $f:Σ^n\to\{0, 1\}$, there is a coalition of $O(n/\sqrt{\log n})$ coordinates and a target output $b\in\{0, 1\}$ such that, after the remaining coordinates are sampled uniformly and independently, the coalition can choose its values to make the output equal to $b$ with probability at least $1-\varepsilon$. The bound is independent of the alphabet size and also holds for monotone Boolean functions on $[0,1]^n$, resolving a conjecture of Friedgut (Combinatorics, Probability and Computing, 2004). Unlike the Boolean cube setting, where Kahn, Kalai, and Linial (FOCS, 1988) give a coalition bound of $O(n/\log n)$, no sublinear bound independent of the alphabet size was previously known. In collective coin flipping, our result gives the first sublinear bound on the number of bad players needed to force a fixed output with probability at least $1-\varepsilon$ in any one-round protocol with independent uniform messages, regardless of the message length. A key ingredient in our proof is an encoding that lets us relate the influence of a function on a product space to the $p$-biased influence of the encoded function. We then rely on a structure theorem of Hatami (Annals of Mathematics, 2012) for functions with small $p$-biased influence to bias the encoded function.

Mon 14 SeptComputational ComplexityDiscrete Mathematics
The gist
Figuring out which small groups of inputs can strongly influence the outcome of a complex function has been a tricky problem, especially when those inputs come from large sets. The authors proved that for any function taking many inputs, there's a relatively small group of them that can be set to force the output to a certain value with very high chance. This resolves a long-standing question and works even when inputs come from large or continuous sets. Their approach used clever encoding techniques and previously known structural insights to make this result possible.
Open → 2609.16401v1

Faster federated learning with first order bilevel optimization

Federated stochastic bilevel optimization with fully first-order gradients

Abstract: Federated stochastic bilevel optimization has been actively studied in recent years due to its widespread applications in machine learning. However, most existing federated stochastic bilevel optimization algorithms require the computation of second-order Hessian and Jacobian matrices, which leads to longer running times in practice. To address these challenges, we propose a novel federated stochastic variance-reduced bilevel gradient descent algorithm that relies solely on first-order oracles. Specifically, our approach does not require the computation of second-order Hessian and Jacobian matrices, significantly reducing running time. Furthermore, we introduce a novel learning rate mechanism, i.e., a constant single-timescale learning rate, to coordinate the update of different variables. We also present a new strategy to establish the convergence rate of our algorithm. Finally, the extensive experimental results confirm the efficacy of our proposed algorithm.

Mon 14 SeptMachine Learning
The gist
Training complex machine learning models often requires updates that involve heavy math, like second-order derivatives, which slow things down. The authors introduce a new method that avoids these slow calculations by using only simpler first-order information, making training quicker in federated (distributed) setups. They also design a stable learning rate approach that improves how different parts of the model update together. Their experiments show this method works well while saving time.
Open → 2609.16350v1

Deterministic labeling improves fault-tolerant connectivity checks in networks

Deterministic Edge-Fault-Tolerant Connectivity Labeling Schemes with Nearly Optimal Label Size

Abstract: For an undirected graph $G = (V,E)$ and a fault bound $f$, an edge-fault-tolerant connectivity labeling scheme assigns short labels to vertices and edges, so that for any vertex pair $(s,t)$ and failed edge set $F\subseteq E$ with $|F|\leq f$, the connectivity between $s$ and $t$ in $G-F$ can be answered by inspecting only the labels of $s$, $t$ and edges in $F$. In this paper, we present a labeling scheme that uses $O(\log^{2}n)$-bit labels that can be computed in deterministic polynomial time. This improves upon the previous $\tilde{O}(\sqrt{f})$ deterministic bound of [Long, Pettie, Saranurak'25], and even slightly improves the $O(\min\{f+\log n,\log^{2}n\log f\})$ randomized bound of [Dory, Parter'21] and [Long, Pettie, Saranurak'25] when $f = Ω(\log^{2}n)$. Moreover, for a general $f$, this is the first labeling scheme that produces an $\tilde{O}(1)$-size labeling which is simultaneously correct across all queries. Our approach combines the cycle-space-based labeling scheme from Dory and Parter with a recent result by [Knauer'26] on sparse cycle bases.

Tue 8 SeptData Structures and Algorithms
The gist
Finding out if two points in a network stay connected after some connections fail can be tricky. The authors show a way to give each point and connection a short label so you can quickly tell if two points remain connected after some failures. Their method is both efficient and guaranteed to work for all possible failures up to a given size. This improves on earlier methods by making the labels smaller and by giving stronger guarantees without randomness.
Open → 2609.09031v1

Maximum mutual visibility sets found and built on cactus graphs

The Maximum Mutual Visibility Set on a Cactus Graph and the Self-stabilizing Constructions

Abstract: Given a graph $G=(V,E)$, let $S$ ($\subseteq V$) be a set of vertices. Two vertices are \emph{mutually visible} if there exists a shortest path in $G$ between them that does not contain any other vertex of $S$. A set $S$ is a \emph{Mutual Visibility Set} (\MVS) if every pair of vertices in $S$ is mutually visible. The concept of \MVS s in graphs has attracted significant attention since its introduction, as it provides an important structural property of graphs. However, determining a maximum \MVS\ in general graphs is computationally intractable; the decision problem of whether a graph admits an \MVS\ of size at least $k$ has been shown to be \emph{NP-complete}. Thus, prior work has focused on finding maximal \MVS s or restricting attention to specific graph classes. Cactus graphs form a fundamental low-treewidth class, yet the maximum \MVS\ problem for this class remains open. In this paper, we first determine the size of maximum \MVS~in cactus graphs, and introduce two self-stabilizing algorithms that construct such sets. The first algorithm uses a single BFS tree and stabilizes in $O(D)$ rounds with $O(\log n)$ bits per process on average; the second one uses parallel BFS trees and stabilizes in $O(|C_{\max}|+|T_{\max}|)$ rounds, which we show to be asymptotically tight as a function of these two parameters, even on graphs where $|C_{\max}|+|T_{\max}| = o(D)$.

Mon 7 SeptDistributed, Parallel, and Cluster ComputingDiscrete MathematicsData Structures and Algorithms
The gist
The paper studies a way to pick special points (vertices) on a network (graph) so that each pair can see each other clearly along shortest paths. Finding the largest such set is known to be very hard for general graphs. The authors focus on cactus graphs, a simpler kind of network, and figure out the exact size of the largest set. They also provide two algorithms that can run independently and eventually build these sets automatically in a distributed way.
Open → 2609.07253v1