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.
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\).
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.
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.
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.
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.
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.
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)$.