EDBT 2026 Demo / reviewers in the wild / expert
Peter Robinson 0002
dblp:06/5220-2
· DBLP profile ↗
65ranked-venue papers
10as first author
21since 2021 · last 2026
0000-0002-7442-7002ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 33 · 6 first-author · 11 since 2021Theory of computation · 18 · 2 first-author · 6 since 2021Security and privacy · 1 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Near-Optimal Bounds for Adversarial Wake-up in Distributed Networks
Peter Robinson 0002, Ming Ming Tan |
SPAA | 1 |
| 2025 | Time-Optimal and Energy-Efficient Deterministic ConsensusabstractWe study fault-tolerant consensus in a variant of the synchronous message passing model, where, in each round, every node can choose to be awake or asleep. This is known as the sleeping model (Chatterjee, Gmyr, Pandurangan PODC 2020) and defines the awake complexity (also called energy complexity), which measures the maximum number of rounds that any node is awake throughout the execution. Only awake nodes can send and receive messages in a given round and all messages sent to sleeping nodes are lost. We present new deterministic consensus algorithms that tolerate up to f < n crash failures, where n is the number of nodes. Our algorithms match the optimal time complexity lower bound of f+1 rounds. For multi-value consensus, where the input values are chosen from some possibly large set, we achieve an energy complexity of 𝒪(⌈ f² / n ⌉) rounds, whereas for binary consensus, we show an algorithm to achieve 𝒪(⌈ f / √n ⌉) energy complexity. Shachar Meir, Hugo Mirault, David Peleg, Peter Robinson 0002 |
OPODIS | 4 |
| 2025 | Brief Announcement: Rise and Shine Efficiently! The Complexity of Adversarial Wake-upabstractWe study the wake-up problem in distributed networks, where an adversary awakens a subset of nodes at arbitrary times, and the goal is to wake up all other nodes as quickly as possible by sending only few messages. We prove the following lower bounds: Peter Robinson 0002, Ming Ming Tan |
PODC | 1 |
| 2025 | Message Optimality and Message-Time Trade-offs for APSP and BeyondabstractRound complexity is an extensively studied metric of distributed algorithms. In contrast, our knowledge of the message complexity of distributed computing problems and its relationship (if any) with round complexity is still quite limited. To illustrate, for many fundamental distributed graph optimization problems such as (exact) diameter computation, All-Pairs Shortest Paths (APSP), Maximum Matching etc., while (near) round-optimal algorithms are known, message-optimal algorithms are hitherto unknown. More importantly, the existing round-optimal algorithms are not message-optimal. This raises two important questions: (1) Can we design message-optimal algorithms for these problems? (2) Can we give message-time tradeoffs for these problems in case the message-optimal algorithms are not round-optimal? Fabien Dufoulon, Shreyas Pai, Gopal Pandurangan, Sriram V. Pemmaraju, Peter Robinson 0002 |
PODC | 5 |
| 2025 | Brief Announcement: Towards Energy-Efficient Distributed AgreementabstractWe study fault-tolerant consensus in a variant of the synchronous message passing model, where, in each round, every node can choose to be awake or asleep. This is known as the sleeping model (Chatterjee, Gmyr, Pandurangan PODC 2020) and defines the awake complexity (also called energy complexity), which measures the maximum number of rounds that any node is awake throughout the execution. Only awake nodes can send and receive messages in a given round and all messages sent to sleeping nodes are lost. We present new deterministic consensus algorithms that tolerate up to f < n crash failures, where n is the number of nodes. Our algorithms match the optimal time complexity lower bound of f + 1 rounds. For multi-value consensus, where the input values are chosen from some possibly large set, we achieve an energy complexity of [EQUATION] rounds, whereas for binary consensus, we show that [EQUATION] rounds are possible. Hugo Mirault, Peter Robinson 0002 |
PODC | 2 |
| 2025 | Brief Announcement: Perfect Matching with Few Link Activations
Hugo Mirault, Peter Robinson 0002, Ming Ming Tan, Xianbin Zhu 0002 |
SIROCCO | 2 |
| 2025 | Tight bounds on the message complexity of distributed tree verification
Shay Kutten, Peter Robinson 0002, Ming Ming Tan |
Distributed Comput. | 2 |
| 2024 | The Message Complexity of Distributed Graph OptimizationabstractThe message complexity of a distributed algorithm is the total number of messages sent by all nodes over the course of the algorithm. This paper studies the message complexity of distributed algorithms for fundamental graph optimization problems. We focus on four classical graph optimization problems: Maximum Matching (MaxM), Minimum Vertex Cover (MVC), Minimum Dominating Set (MDS), and Maximum Independent Set (MaxIS). In the sequential setting, these problems are representative of a wide spectrum of hardness of approximation. While there has been some progress in understanding the round complexity of distributed algorithms (for both exact and approximate versions) for these problems, much less is known about their message complexity and its relation with the quality of approximation. We almost fully quantify the message complexity of distributed graph optimization by showing the following results: 1) Cubic regime: Our first main contribution is showing essentially cubic, i.e., Ω̃(n³) lower bounds (where n is the number of nodes in the graph) on the message complexity of distributed exact computation of Minimum Vertex Cover (MVC), Minimum Dominating Set (MDS), and Maximum Independent Set (MaxIS). Our lower bounds apply to any distributed algorithm that runs in polynomial number of rounds (a mild and necessary restriction). Our result is significant since, to the best of our knowledge, this are the first ω(m) (where m is the number of edges in the graph) message lower bound known for distributed computation of such classical graph optimization problems. Our bounds are essentially tight, as all these problems can be solved trivially using O(n³) messages in polynomial rounds. All these bounds hold in the standard CONGEST model of distributed computation in which messages are of O(log n) size. 2) Quadratic regime: In contrast, we show that if we allow approximate computation then Θ̃(n²) messages are both necessary and sufficient. Specifically, we show that Ω̃(n²) messages are required for constant-factor approximation algorithms for all four problems. For MaxM and MVC, these bounds hold for any constant-factor approximation, whereas for MDS and MaxIS they hold for any approximation factor better than some specific constants. These lower bounds hold even in the LOCAL model (in which messages can be arbitrarily large) and they even apply to algorithms that take arbitrarily many rounds. We show that our lower bounds are essentially tight, by showing that if we allow approximation to within an arbitrarily small constant factor, then all these problems can be solved using Õ(n²) messages even in the CONGEST model. 3) Linear regime: We complement the above lower bounds by showing distributed algorithms with Õ(n) message complexity that run in polylogarithmic rounds and give constant-factor approximations for all four problems on random graphs. These results imply that almost linear (in n) message complexity is achievable on almost all (connected) graphs of every edge density. Fabien Dufoulon, Shreyas Pai, Gopal Pandurangan, Sriram V. Pemmaraju, Peter Robinson 0002 |
ITCS | 5 |
| 2024 | Dynamic Maximal Matching in Clique Networks
Minming Li, Peter Robinson 0002, Xianbin Zhu 0002 |
ITCS | 2 |
| 2024 | The Singular Optimality of Distributed Computation in LOCALabstractIt has been shown that one can design distributed algorithms that are (nearly) singularly optimal, meaning they simultaneously achieve optimal time and message complexity (within polylogarithmic factors), for several fundamental global problems such as broadcast, leader election, and spanning tree construction, under the KT₀ assumption. With this assumption, nodes have initial knowledge only of themselves, not their neighbors. In this case the time and message lower bounds are Ω(D) and Ω(m), respectively, where D is the diameter of the network and m is the number of edges, and there exist (even) deterministic algorithms that simultaneously match these bounds. On the other hand, under the KT₁ assumption, whereby each node has initial knowledge of itself and the identifiers of its neighbors, the situation is not clear. For the KT₁ CONGEST model (where messages are of small size), King, Kutten, and Thorup (KKT) showed that one can solve several fundamental global problems (with the notable exception of BFS tree construction) such as broadcast, leader election, and spanning tree construction with Õ(n) message complexity (n is the network size), which can be significantly smaller than m. Randomization is crucial in obtaining this result. While the message complexity of the KKT result is near-optimal, its time complexity is Õ(n) rounds, which is far from the standard lower bound of Ω(D). An important open question is whether one can achieve singular optimality for the above problems in the KT₁ CONGEST model, i.e., whether there exists an algorithm running in Õ(D) rounds and Õ(n) messages. Another important and related question is whether the fundamental BFS tree construction can be solved with Õ(n) messages (regardless of the number of rounds as long as it is polynomial in n) in KT₁. In this paper, we show that in the KT₁ LOCAL model (where message sizes are not restricted), singular optimality is achievable. Our main result is that all global problems, including BFS tree construction, can be solved in Õ(D) rounds and Õ(n) messages, where both bounds are optimal up to polylogarithmic factors. Moreover, we show that this can be achieved deterministically. Fabien Dufoulon, Gopal Pandurangan, Peter Robinson 0002, Michele Scquizzato |
OPODIS | 3 |
| 2024 | Sparse Spanners with Small Distance and Congestion StretchesabstractGiven a graph G, a classical problem in graph theory is the construction of a spanner H -- a sparse subgraph of G that closely approximates the distances between nodes in G. The distance stretch~α of H is the factor of how much the distances in H increase versus G. Here, we consider sparse spanner constructions that can also preserve the node congestion of routing problems in G. The congestion stretch β of H is the factor of how much the (smallest) congestion of a routing problem increases in H versus G. We introduce the notion of (α, β)-DC-spanner (i.e., a Distance-Congestion-spanner) that simultaneously controls the stretches for distance and congestion. We show that for expander graphs with n nodes, there is a (3, O(log n))-DC-spanner with O(n5/3) edges. We also examine Δ-regular graphs with Δ ≥ n2/3, where we show how to obtain a (3, O(√Δ ⋅ log n))-DC-spanner with O(n5/3 log2n) edges. Finally, we show that there is a graph such that any optimal size 3-distance spanner has Ω(n7/6) edges and is a (3, Ω(n1/6))-DC-spanner. Costas Busch, Dariusz R. Kowalski, Peter Robinson 0002 |
SPAA | 3 |
| 2023 | Tight Bounds on the Message Complexity of Distributed Tree Verification
Shay Kutten, Peter Robinson 0002, Ming Ming Tan |
OPODIS | 2 |
| 2023 | Improved Tradeoffs for Leader ElectionabstractWe consider leader election in clique networks, where n nodes are connected by point-to-point communication links. For the synchronous clique under simultaneous wake-up, i.e., where all nodes start executing the algorithm in round 1, we show a tradeoff between the number of messages and the amount of time. The previous lower bound side of such a tradeoff, in the seminal paper of Afek and Gafni (1991), was shown only assuming adversarial wake-up. Interestingly, our new tradeoff also improves the previous lower bounds for a large part of the spectrum, even under simultaneous wake-up. More specifically, we show that any deterministic algorithm with a message complexity of n f(n) requires Ω((log n) / (log f(n)+1)) rounds, for f(n) > 1. Our result holds even if the node IDs are chosen from a relatively small set of size Θ(n log n), as we are able to avoid using Ramsey's theorem, in contrast to many existing lower bounds for deterministic algorithms. We also give an upper bound that improves over the previously-best tradeoff achieved by the algorithm of Afek and Gafni. Our second contribution for the synchronous clique under simultaneous wake-up is to show that Ω (n log n) is in fact a lower bound on the message complexity that holds for any deterministic algorithm with a termination time T(n) (i.e., any function of n), for a sufficiently large ID space. We complement this result by giving a simple deterministic algorithm that achieves leader election in sublinear time while sending only o(n log n) messages, if the ID space is of at most linear size. We also show that Las Vegas algorithms (that never fail) require Θ(n) messages. This exhibits a gap between Las Vegas and Monte Carlo algorithms. Shay Kutten, Peter Robinson 0002, Ming Ming Tan, Xianbin Zhu 0002 |
PODC | 2 |
| 2023 | Brief Announcement: What Can We Compute in a Single Round of the Congested Clique?abstractWe show that any one-round algorithm that computes a minimum spanning tree (MST) in the unicast congested clique must use a link bandwidth of Ω(log3 n) bits in the worst case. Consequently, computing an MST under the standard assumption of O(log n)-size messages requires at least 2 rounds. This is the first round complexity lower bound in the unicast congested clique for a problem where the output size is small, i.e., O(n log n) bits. Our lower bound holds as long as every edge of the MST is output by an incident node. To the best of our knowledge, all prior lower bounds for the unicast congested clique either considered problems with large output sizes (e.g., subgraph enumeration) or required every node to learn the entire output. Peter Robinson 0002 |
PODC | 1 |
| 2023 | Distributed Sketching Lower Bounds for k-Edge Connected Spanning Subgraphs, BFS Trees, and LCL Problems
Peter Robinson 0002 |
DISC | 1 |
| 2023 | Leader Election in Well-Connected Graphs
Seth Gilbert, Peter Robinson 0002, Suman Sourav |
Algorithmica | 2 |
| 2022 | Byzantine-Resilient Counting in NetworksabstractWe present two distributed algorithms for the Byzantine counting problem, which is concerned with estimating the size of a network in the presence of a large number of Byzantine nodes.In an n-node network (n is unknown), our first algorithm, which is deterministic, finishes in O(log n) rounds and is time-optimal. This algorithm can tolerate up to O(n1−γ) arbitrarily (adversarially) placed Byzantine nodes for any arbitrarily small (but fixed) positive constant γ. It outputs a (fixed) constant factor estimate of log n that would be known to all but o(1) fraction of the good nodes. This algorithm works for any bounded degree expander network. However, this algorithms assumes that good nodes can send arbitrarily large-sized messages in a round.Our second algorithm is randomized and most good nodes send only small-sized messages.1This algorithm works in almost all d-regular graphs. It tolerates up to $B(n) = {n^{\frac{1}{2} - \xi }}$ (note that n and B(n) are unknown to the algorithm) arbitrarily (adversarially) placed Byzantine nodes, where ξ is any arbitrarily small (but fixed) positive constant. This algorithm takes O(B(n) log2n) rounds and outputs a constant factor estimate of log n with probability at least 1−o(1). The said estimate is known to most nodes, i.e., ≥ (1 − β)n nodes for any arbitrarily small (but fixed) positive constant β.To complement our algorithms, we also present an impossibility result that shows that it is impossible to estimate the network size with any reasonable approximation with any non-trivial probability of success if the network does not have sufficient vertex expansion.Both algorithms are the first such algorithms that solve Byzantine counting in sparse, bounded degree networks under very general assumptions. Both algorithms are fully local and need no global knowledge.Our algorithms can be used for the design of efficient distributed algorithms resilient against Byzantine failures, where the knowledge of the network size — a global parameter — may not be known a priori. Soumyottam Chatterjee, Gopal Pandurangan, Peter Robinson 0002 |
ICDCS | 3 |
| 2022 | Awake-Efficient Distributed Algorithms for Maximal Independent SetabstractWe present a simple algorithmic framework for designing efficient distributed algorithms for the fundamental symmetry breaking problem of Maximal Independent Set (MIS) in the sleeping model [Chatterjee et al, PODC 2020]. In the sleeping model, only the rounds in which a node is awake are counted for the awake complexity, while sleeping rounds are ignored. This is motivated by the fact that a node spends resources only in its awake rounds and hence the goal is to minimize the awake complexity.Our framework allows us to design distributed MIS algorithms that have ${\mathcal{O}}({\text{polyloglog }}n)$ (worst-case) awake complexity in certain important graph classes which satisfy the so-called adjacency property. Informally, the adjacency property guarantees that the graph can be partitioned into an appropriate number of classes so that each node has at least one neighbor belonging to every class. Graphs that can satisfy the adjacency property are random graphs with large clustering coefficient such as random geometric graphs as well as line graphs of regular (or near regular) graphs.We first apply our framework to design two randomized distributed MIS algorithms for random geometric graphs of arbitrary dimension d (even non-constant). The first algorithm has ${\mathcal{O}}({\text{polyloglog }}n)$ (worst-case) awake complexity with high probability, where n is the number of nodes in the graph.1This means that any node in the network spends only ${\mathcal{O}}({\text{polyloglog }}n)$ awake rounds; this is almost exponentially better than the (traditional) time complexity of ${\mathcal{O}}({\text{log }}n)$ rounds (where there is no distinction between awake and sleeping rounds) known for distributed MIS algorithms on general graphs or even the faster ${\mathcal{O}}\left({\sqrt {\frac{{{\text{log }}n}}{{{\text{loglog }}n}}} }\right)$ rounds known for Erdos-Renyi random graphs. However, the (traditional) time complexity of our first algorithm is quite large—essentially proportional to the degree of the graph. Our second algorithm has a slightly worse awake complexity of ${\mathcal{O}}(d\,{\text{polyloglog }}n)$, but achieves a significantly better time complexity of ${\mathcal{O}}(d\,\log n\,{\text{polyloglog }}n)$ rounds whp.We also show that our framework can be used to design ${\mathcal{O}}({\text{polyloglog }}n)$ awake complexity MIS algorithms in other types of random graphs, namely an augmented Erdos-Renyi random graph that has a large clustering coefficient. Khalid Hourani, Gopal Pandurangan, Peter Robinson 0002 |
ICDCS | 3 |
| 2022 | Latency, capacity, and distributed minimum spanning trees
John Augustine 0001, Seth Gilbert, Fabian Kuhn, Peter Robinson 0002, Suman Sourav |
J. Comput. Syst. Sci. | 4 |
| 2021 | Can We Break Symmetry with o(m) Communication?abstractWe study the communication cost (or message complexity) of fundamental distributed symmetry breaking problems, namely, coloring and MIS. While significant progress has been made in understanding and improving the running time of such problems, much less is known about the message complexity of these problems. In fact, all known algorithms need at least Ω(m) communication for these problems, where m is the number of edges in the graph. We addressthe following question in this paper: can we solve problems such as coloring and MIS using sublinear, i.e., o(m) communication, and if sounder what conditions? Shreyas Pai, Gopal Pandurangan, Sriram V. Pemmaraju, Peter Robinson 0002 |
PODC | 4 |
| 2021 | Being Fast Means Being Chatty: The Local Information Cost of Graph SpannersabstractWe introduce a new measure for quantifying the amount of information that the nodes in a network need to learn to solve a graph problem. We show that the local information cost (LIC) presents a natural lower bound on the communication complexity of distributed algorithms. For the synchronous CONGEST-KT1 model, where each node has initial knowledge of its neighbors' IDs, we prove that bits are required for solving a graph problem P with a τ-round algorithm that errs with probability at most γ. Our result is the first lower bound that yields a general trade-off between communication and time for graph problems in the CONGEST-KT1 model. We demonstrate how to apply the local information cost by deriving a lower bound on the communication complexity of computing a multiplicative spanner with stretch 2t – 1 that consists of at most edges, where ∊ = O(1/t2). Our main result is that any O(poly(n))-time algorithm must send at least bits in the CONGEST model under the KT1 assumption. Previously, only a trivial lower bound of bits was known for this problem; in fact, this is the first nontrivial lower bound on the communication complexity of a sparse subgraph problem in this setting. A consequence of our lower bound is that achieving both time- and communication-optimality is impossible when designing a distributed spanner algorithm. In light of the work of King, Kutten, and Thorup (2015), this shows that computing a minimum spanning tree can be done significantly faster than finding a spanner when considering algorithms with Õ(n) communication complexity. Our result also implies time complexity lower bounds for constructing a spanner in the node-congested clique of Augustine et al. (2019) and in the push-pull gossip model with limited bandwidth. Peter Robinson 0002 |
SODA | 1 |
| 2020 | Latency, Capacity, and Distributed Minimum Spanning Tree†abstractWe study the cost of distributed MST construction in the setting where each edge has a latency and a capacity, along with the weight. Edge latencies capture the delay on the links of the communication network, while capacity captures their throughput (the rate at which messages can be sent). Depending on how the edge latencies relate to the edge weights, we provide several tight bounds on the time and messages required to construct an MST.When edge weights exactly correspond with the latencies, we show that, perhaps interestingly, the bottleneck parameter in determining the running time of an algorithm is the total weight W of the MST (rather than the total number of nodes n, as in the standard CONGEST model). That is, we show a tight bound of $\tilde \Theta $ (D + $\sqrt {W/c} $) rounds, where D refers to the latency diameter of the graph, W refers to the total weight of the constructed MST and edges have capacity c. The proposed algorithm sends Õ (m + W) messages, where m, the total number of edges in the network graph under consideration, is a known lower bound on message complexity for MST construction. We also show that Ω(W) is a lower bound for fast MST constructions.When the edge latencies and the corresponding edge weights are unrelated, and either can take arbitrary values, we show that (unlike the sub-linear time algorithms in the standard CONGEST model, on small diameter graphs), the best time complexity that can be achieved is Θ(D + n/c). However, if we restrict all edges to have equal latency ℓ and capacity c while having possibly different weights (weights could deviate arbitrarily from ℓ), we give an algorithm that constructs an MST in Õ (D + $\sqrt {n\ell /c} $) time. In each case, we provide nearly matching upper and lower bounds. John Augustine 0001, Seth Gilbert, Fabian Kuhn, Peter Robinson 0002, Suman Sourav |
ICDCS | 4 |
| 2020 | DConstructor: Efficient and Robust Network Construction with Polylogarithmic OverheadabstractWith the rise of dynamic reconfigurable networks such as Peer-to-Peer (P2P) networks, overlay networks, ad hoc wireless and mesh networks, it has become important to construct and maintain topologies with various desirable properties (such as connectivity, low diameter, expansion, low degree etc.) in an efficient decentralized manner. The main result of this paper is a distributed protocol called DConstructor that given any (connected) network topology will "converge" to a given (desired) target topology such as an expander, hypercube, or Chord, with high probability. Our protocol is efficient, lightweight, and scalable, and it incurs only O(polylog(n)) overhead (where n is the network size) for topology construction and maintenance: only polylogarithmic (in n) bits need to be processed and sent by each node per round, the convergence time is polylogarithmic rounds and any node's computation cost per round is also polylogarithmic. Our protocol is robust and self-repairing in the sense that it will converge to the desired topology in polylogarithmic rounds and polylogarithmic communication cost under dynamic topology changes and arbitrary insertions and deletions of nodes. Seth Gilbert, Gopal Pandurangan, Peter Robinson 0002, Amitabh Trehan |
PODC | 3 |
| 2020 | The complexity of leader election in diameter-two networks
Soumyottam Chatterjee, Gopal Pandurangan, Peter Robinson 0002 |
Distributed Comput. | 3 |
| 2020 | A Time- and Message-Optimal Distributed Algorithm for Minimum Spanning TreesabstractThis paper presents a randomized (Las Vegas) distributed algorithm that constructs a minimum spanning tree (MST) in weighted networks with optimal (up to polylogarithmic factors) time and message complexity.This algorithm runs in Õ(D + √ n) time and exchanges Õ(m) messages (both with high probability), where n is the number of nodes of the network, D is the hop-diameter, and m is the number of edges.This is the first distributed MST algorithm that matches simultaneously the time lower bound of Ω(D + √ n) [Elkin, SIAM J. Comput.2006] and the message lower bound of Ω(m) [Kutten et al., J. ACM 2015], which both apply to randomized Monte Carlo algorithms.The prior time and message lower bounds are derived using two completely different graph constructions; the existing lower bound construction that shows one lower bound does not work for the other.To complement our algorithm, we present a new lower bound graph construction for which any distributed MST algorithm requires both Ω(D + √ n) rounds and Ω(m) messages. Gopal Pandurangan, Peter Robinson 0002, Michele Scquizzato |
ACM Trans. Algorithms | 2 |
| 2019 | Network Size Estimation in Small-World Networks Under Byzantine FaultsabstractWe study the fundamental problem of counting the number of nodes in a sparse network (of unknown size) under the presence of a large number of Byzantine nodes. We assume the full information model where the Byzantine nodes have complete knowledge about the entire state of the network at every round (including random choices made by all the nodes), have unbounded computational power, and can deviate arbitrarily from the protocol. Essentially all known algorithms for fundamental Byzantine problems (e.g., agreement, leader election, sampling) studied in the literature assume the knowledge (or at least an estimate) of the size of the network. It is nontrivial to design algorithms for Byzantine problems that work without knowledge of the network size, especially in boundeddegree (expander) networks where the local views of all nodes are (essentially) the same and limited, and Byzantine nodes can quite easily fake the presence/absence of non-existing nodes. To design truly local algorithms that do not rely on any global knowledge (including network size), estimating the size of the network under Byzantine nodes is an important first step. Our main contribution is a randomized distributed algorithm that estimates the size of a network under the presence of a large number of Byzantine nodes. In particular, our algorithm estimates the size of a sparse, “small-world”, expander network with up to O(n1-δ) Byzantine nodes, where n is the (unknown) network size and δ > 0 can be be any arbitrarily small (but fixed) constant. Our algorithm outputs a (fixed) constant factor estimate of log(n) with high probability; the correct estimate of the network size will be known to a large fraction ((1 - ϵ)fraction, for any fixed positive constant ϵ) of the honest nodes. Our algorithm is fully distributed, lightweight, and simple to implement, runs in O(log3n) rounds, and requires nodes to send and receive messages of only small-sized messages per round; any node's local computation cost per round is also small. Soumyottam Chatterjee, Gopal Pandurangan, Peter Robinson 0002 |
IPDPS | 3 |
| 2019 | The Complexity of Symmetry Breaking in Massive GraphsabstractThe goal of this paper is to understand the complexity of symmetry breaking problems, specifically maximal independent set (MIS) and the closely related $β$-ruling set problem, in two computational models suited for large-scale graph processing, namely the $k$-machine model and the graph streaming model. We present a number of results. For MIS in the $k$-machine model, we improve the $\tilde{O}(m/k^2 + Δ/k)$-round upper bound of Klauck et al. (SODA 2015) by presenting an $\tilde{O}(m/k^2)$-round algorithm. We also present an $\tildeΩ(n/k^2)$ round lower bound for MIS, the first lower bound for a symmetry breaking problem in the $k$-machine model. For $β$-ruling sets, we use hierarchical sampling to obtain more efficient algorithms in the $k$-machine model and also in the graph streaming model. More specifically, we obtain a $k$-machine algorithm that runs in $\tilde{O}(βnΔ^{1/β}/k^2)$ rounds and, by using a similar hierarchical sampling technique, we obtain one-pass algorithms for both insertion-only and insertion-deletion streams that use $O(β\cdot n^{1+1/2^{β-1}})$ space. The latter result establishes a clear separation between MIS, which is known to require $Ω(n^2)$ space (Cormode et al., ICALP 2019), and $β$-ruling sets, even for $β= 2$. Finally, we present an even faster 2-ruling set algorithm in the $k$-machine model, one that runs in $\tilde{O}(n/k^{2-ε} + k^{1-ε})$ rounds for any $ε$, $0 \le ε\le 1$. Christian Konrad 0001, Sriram V. Pemmaraju, Talal Riaz, Peter Robinson 0002 |
DISC | 4 |
| 2019 | Slow Links, Fast Links, and the Cost of GossipabstractConsider the classical problem of information dissemination: one (or more) nodes in a network have some information that they want to distribute to the remainder of the network. In this paper, we study the cost of information dissemination in networks where edges have latencies, i.e., sending a message from one node to another takes some amount of time. We first generalize the idea of conductance to weighted graphs by defining φ*to be the “critical weighted conductance” and ℓ*to be the “critical latency”. One goal of this paper is to argue that φ*characterizes the connectivity of a weighted graph with latencies in much the same way that conductance characterizes the connectivity of unweighted graphs. We give near tight lower and upper bounds on the problem of information dissemination, up to polylogarithmic factors. Specifically, we show that in a graph with (weighted) diameter D (with latencies as weights) and maximum degree Δ, any information dissemination algorithm requires at least Ω(min(D+Δ,ℓ*/φ*)) time in the worst case. We show several variants of the lower bound (e.g., for graphs with small diameter, graphs with small max-degree, etc.) by reduction to a simple combinatorial game. We then give nearly matching algorithms, showing that information dissemination can be solved in O(min((D+Δ)log3n,(ℓ*/φ*)log n) time. This is achieved by combining two cases. We show that the classical push-pull algorithm is (near) optimal when the diameter or the maximum degree is large. For the case where the diameter and the maximum degree are small, we give an alternative strategy in which we first discover the latencies and then use an algorithm for known latencies based on a weighted spanner construction. (Our algorithms are within polylogarithmic factors of being tight both for known and unknown latencies.) While it is easiest to express our bounds in terms of φ*and ℓ*, in some cases they do not provide the most convenient definition of conductance in weighted graphs. Therefore we give a second (nearly) equivalent characterization, namely the average weighted conductance φavg. Suman Sourav, Peter Robinson 0002, Seth Gilbert |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2018 | Slow Links, Fast Links, and the Cost of GossipabstractConsider the classical problem of information dissemination: one (or more) nodes in a network have some information that they want to distribute to the remainder of the network. In this paper, we study the cost of information dissemination in networks where edges have latencies, i.e., sending a message from one node to another takes some amount of time. We first generalize the idea of conductance to weighted graphs by defining φ*to be the "critical conductance" and ℓ*to be the "critical latency". One goal of this paper is to argue that φ*characterizes the connectivity of a weighted graph with latencies in much the same way that conductance characterizes the connectivity of unweighted graphs. We give near tight lower and upper bounds on the problem of information dissemination, up to polylogarithmic factors. Specifically, we show that in a graph with (weighted) diameter d (with latencies as weights) and maximum degree Δ, any information dissemination algorithm requires at least Δ(min(D+Δ, ℓ*/φ*)) time in the worst case. We show several variants of the lower bound (e.g., for graphs with small diameter, graphs with small max-degree, etc.) by reduction to a simple combinatorial game. We then give nearly matching algorithms, showing that information dissemination can be solved in O(min((D+Δ)log3n, (ℓ*/φ;*)\log n) time. This is achieved by combining two cases. We show that the classical push-pull algorithm is (near) optimal when the diameter or the maximum degree is large. For the case where the diameter and the maximum degree are small, we give an alternative strategy in which we first discover the latencies and then use an algorithm for known latencies based on a weighted spanner construction. (Our algorithms are within polylogarithmic factors of being tight both for known and unknown latencies.) While it is easiest to express our bounds in terms of φ*and ℓ*, in some cases they do not provide the most convenient definition of conductance in weighted graphs. Therefore, we give a second (nearly) equivalent characterization, namely the average conductance φavg. Suman Sourav, Peter Robinson 0002, Seth Gilbert |
ICDCS | 2 |
| 2018 | Leader Election in Well-Connected GraphsabstractIn this paper, we look at the problem of randomized leader election in synchronous distributed networks with a special focus on the message complexity. We provide an algorithm that solves the implicit version of leader election (where non-leader nodes need not be aware of the identity of the leader) in any general network with O( √ n log7/2 n tmix ) messages and in O(tmix log2 n) time, where n is the number of nodes and tmix refers to the mixing time of a random walk in the network graph G. For several classes of wellconnected networks (that have a large conductance or alternatively small mixing times e.g. expanders, hypercubes, etc), the above result implies extremely efficient (sublinear running time and messages) leader election algorithms. Correspondingly, we show that any substantial improvement is not possible over our algorithm, by presenting an almost matching lower bound for randomized leader election. We show that Ω( √ n/Φ3/4) messages are needed for any leader election algorithm that succeeds with probability at least 1 - o(1), where Φ refers to the conductance of a graph. To the best of our knowledge, this is the first work that shows a dependence between the time and message complexity to solve leader election and the connectivity of the graph G, which is often characterized by the graph's conductance Φ. Apart from the Ω(m) bound in [23] (where m denotes the number of edges of the graph), this work also provides one of the first non-trivial lower bounds for leader election in general networks. Seth Gilbert, Peter Robinson 0002, Suman Sourav |
PODC | 2 |
| 2018 | Session details: Session 1D: Graph Algorithms
Peter Robinson 0002 |
PODC | 1 |
| 2018 | On the Distributed Complexity of Large-Scale Graph ComputationsabstractMotivated by the increasing need to understand the distributed algorithmic foundations of large-scale graph computations, we study some fundamental graph problems in a message-passing model for distributed computing where k ≥ 2 machines jointly perform computations on graphs with n nodes (typically, n >> k). The input graph is assumed to be initially randomly partitioned among the k machines, a common implementation in many real-world systems. Communication is point-to-point, and the goal is to minimize the number of communication rounds of the computation. Our main contribution is the General Lower Bound Theorem , a theorem that can be used to show non-trivial lower bounds on the round complexity of distributed large-scale data computations. This result is established via an information-theoretic approach that relates the round complexity to the minimal amount of information required by machines to solve the problem. Our approach is generic, and this theorem can be used in a “cookbook” fashion to show distributed lower bounds for several problems, including non-graph problems. We present two applications by showing (almost) tight lower bounds on the round complexity of two fundamental graph problems, namely, PageRank computation and triangle enumeration . These applications show that our approach can yield lower bounds for problems where the application of communication complexity techniques seems not obvious or gives weak bounds, including and especially under a stochastic partition of the input. We then present distributed algorithms for PageRank and triangle enumeration with a round complexity that (almost) matches the respective lower bounds; these algorithms exhibit a round complexity that scales superlinearly in k , improving significantly over previous results [Klauck et al., SODA 2015]. Specifically, we show the following results: PageRank: We show a lower bound of Ὼ(n/k 2 ) rounds and present a distributed algorithm that computes an approximation of the PageRank of all the nodes of a graph in Õ(n/k 2 ) rounds. Triangle enumeration: We show that there exist graphs with m edges where any distributed algorithm requires Ὼ(m/k 5/3 ) rounds. This result also implies the first non-trivial lower bound of Ὼ(n 1/3 ) rounds for the congested clique model, which is tight up to logarithmic factors. We then present a distributed algorithm that enumerates all the triangles of a graph in Õ(m/k 5/3 + n/k 4/3 ) rounds. Gopal Pandurangan, Peter Robinson 0002, Michele Scquizzato |
SPAA | 2 |
| 2018 | Breaking the $ilde$Omega($sqrt{n})$ Barrier: Fast Consensus under a Late AdversaryabstractWe study the consensus problem in a synchronous distributed system of n nodes under an adaptive adversary that has a slightly outdated view of the system and can block all incoming and outgoing communication of a constant fraction of the nodes in each round. Motivated by a result of Ben-Or and Bar-Joseph (1998), showing that any consensus algorithm that is resilient against a linear number of crash faults requires $\tilde Ømega(\sqrtn )$ rounds in an n -node network against an adaptive adversary, we consider a late adaptive adversary, who has full knowledge of the network state at the beginning of the previous round and unlimited computational power, but is oblivious to the current state of the nodes. % Our main contributions are randomized distributed algorithms that achieve consensus with high probability among all except a small constant fraction of the nodes (i.e.,\ "almost-everywhere'') against a late adaptive adversary who can block up to ε n$ nodes in each round, for a small constant ε >0$. Our first protocol achieves binary almost-everywhere consensus and also guarantees a decision on the majority input value, thus ensuring plurality consensus. We also present an algorithm that achieves the same time complexity for multi-value consensus. Both of our algorithms succeed in $O(łog n)$ rounds with high probability, thus showing an exponential gap to the $\tildeØmega(\sqrtn )$ lower bound of Ben-Or and Bar-Joseph for strongly adaptive crash-failure adversaries, which can be strengthened to $Ømega(n)$ when allowing the adversary to block nodes instead of permanently crashing them. Our algorithms are scalable to large systems as each node contacts only an (amortized) constant number of peers in each communication round. We show that our algorithms are optimal up to constant (resp.\ sub-logarithmic) factors by proving that every almost-everywhere consensus protocol takes $Ømega(łog_d n)$ rounds in the worst case, where d is an upper bound on the number of communication requests initiated per node in each round. We complement our theoretical results with an experimental evaluation of the binary almost-everywhere consensus protocol revealing a short convergence time even against an adversary blocking a large fraction of nodes. Peter Robinson 0002, Christian Scheideler, Alexander Setzer |
SPAA | 1 |
| 2018 | Fault-Tolerant Consensus with an Abstract MAC LayerabstractIn this paper, we study fault-tolerant distributed consensus in wireless systems. In more detail, we produce two new randomized algorithms that solve this problem in the abstract MAC layer model, which captures the basic interface and communication guarantees provided by most wireless MAC layers. Our algorithms work for any number of failures, require no advance knowledge of the network participants or network size, and guarantee termination with high probability after a number of broadcasts that are polynomial in the network size. Our first algorithm satisfies the standard agreement property, while our second trades a faster termination guarantee in exchange for a looser agreement property in which most nodes agree on the same value. These are the first known fault-tolerant consensus algorithms for this model. In addition to our main upper bound results, we explore the gap between the abstract MAC layer and the standard asynchronous message passing model by proving fault-tolerant consensus is impossible in the latter in the absence of information regarding the network participants, even if we assume no faults, allow randomized solutions, and provide the algorithm a constant-factor approximation of the network size. Calvin C. Newport, Peter Robinson 0002 |
DISC | 2 |
| 2018 | Gracefully degrading consensus and k-set agreement in directed dynamic networksabstractWe study distributed agreement in synchronous directed dynamic networks, where an omniscient message adversary controls the presence/absence of communication links. We prove that consensus is impossible under a message adversary that guarantees weak connectivity only, and introduce eventually vertex-stable source components (VSSCs) as a means for circumventing this impossibility: A VSSC ( k , d ) message adversary guarantees that, eventually, there is an interval of d consecutive rounds where every communication graph contains at most k strongly connected components consisting of the same processes (with possibly varying interconnect topology), which have no incoming links from outside processes. We present a consensus algorithm that works correctly under a VSSC ( 1 , 4 E + 2 ) message adversary, where E is the dynamic network depth. Our algorithm maintains local estimates of the communication graphs, and applies techniques for detecting network stability and univalent system configurations. Several related impossibility results and lower bounds, in particular, that neither a VSSC ( 1 , E − 1 ) message adversary nor a VSSC ( 2 , ∞ ) one allow to solve consensus, reveal that there is not much hope to deal with (much) stronger message adversaries here. However, we show that gracefully degrading consensus, which degrades to general k -set agreement in case of unfavorable network conditions, allows to cope with stronger message adversaries: We provide a k -universal k -set agreement algorithm, where the number of system-wide decision values k is not encoded in the algorithm, but rather determined by the actual power of the message adversary in a run: Our algorithm guarantees at most k decision values under a VSSC ( n , d ) + MAJINF ( k ) message adversary, which combines VSSC ( n , d ) (with some small value of d , ensuring termination) with some information flow guarantee MAJINF ( k ) between certain VSSCs (ensuring k -agreement). Since related impossibility results reveal that a VSSC ( k , d ) message adversary is too strong for solving k -set agreement and that some information flow between VSSCs is mandatory for this purpose as well, our results provide a significant step towards the exact solvability/impossibility border of general k -set agreement in directed dynamic networks. Finally, we relate (the eventually-forever-variants of) our message adversaries to failure detectors. It turns out that even though VSSC ( 1 , ∞ ) allows to solve consensus and to implement the Ω failure detector, it does not allow to implement Σ. This contrasts the fact that, in asynchronous message-passing systems with a majority of process crashes, ( Σ , Ω ) is a weakest failure detector for solving consensus. Similarly, although the message adversary VSSC ( n , d ) + MAJINF ( k ) allows to solve k -set agreement, it does not allow to implement the failure detector Σ k , which is known to be necessary for k -set agreement in asynchronous message-passing systems with a majority of process crashes. Consequently, it is not possible to adapt failure-detector-based algorithms to work in conjunction with our message adversaries. Martin Biely, Peter Robinson 0002, Ulrich Schmid 0001, Manfred Schwarz, Kyrill Winkler |
Theor. Comput. Sci. | 2 |
| 2018 | Special Issue of ICDCN 2016 (Distributed Computing Track)
Gopal Pandurangan, Peter Robinson 0002 |
Theor. Comput. Sci. | 2 |
| 2017 | Brief Announcement: Gossiping with LatenciesabstractConsider the classical problem of information dissemination: one (or more) nodes in a network have some information that they want to distribute to the remainder of the network. In this paper, we study the cost of information dissemination in networks where edges have latencies, i.e., sending a message from one node to another takes some amount of time. We first generalize the idea of conductance to weighted graphs, defining φ* to be the "weighted conductance" and l* to be the "critical latency." One goal of this paper is to argue that φ* characterizes the connectivity of a weighted graph with latencies in much the same way that conductance characterizes the connectivity of unweighted graphs. We give near tight lower and upper bounds on the problem of information dissemination. Specifically, we show that in a graph with (weighted) diameter D (with latencies as weights), maximum degree Δ, weighted conductance φ* and critical latency l*, any information dissemination algorithm requires at least Ω(min(D+Δ, l*/φ*)) time. We then give nearly matching algorithms, showing that information dissemination can be solved in O(min((D + Δ)log3n), (l*/φ*)log(n)) time. Seth Gilbert, Peter Robinson 0002, Suman Sourav |
PODC | 2 |
| 2017 | Brief Announcement: Symmetry Breaking in the CONGEST Model: Time- and Message-Efficient Algorithms for Ruling SetsabstractWe study local symmetry breaking problems in the Congest model, focusing on ruling set problems, which generalize the fundamental Maximal Independent Set (MIS) problem. Our work is motivated by the following central question: can we break the long-standing Θ(log n) time-complexity barrier and the Θ(m) message-complexity barrier in the Congest model for MIS or closely-related symmetry breaking problems? Shreyas Pai, Gopal Pandurangan, Sriram V. Pemmaraju, Talal Riaz, Peter Robinson 0002 |
PODC | 5 |
| 2017 | A time- and message-optimal distributed algorithm for minimum spanning treesabstractThis paper presents a randomized (Las Vegas) distributed algorithm that constructs a minimum spanning tree (MST) in weighted networks with optimal (up to polylogarithmic factors) time and message complexity. This algorithm runs in Õ(D + √n) time and exchanges Õ(m) messages (both with high probability), where n is the number of nodes of the network, D is the diameter, and m is the number of edges. This is the first distributed MST algorithm that matches simultaneously the time lower bound of Ω(D + √n) [Elkin, SIAM J. Comput. 2006] and the message lower bound of Ω(m) [Kutten et al., J. ACM 2015], which both apply to randomized Monte Carlo algorithms. Gopal Pandurangan, Peter Robinson 0002, Michele Scquizzato |
STOC | 2 |
| 2017 | Symmetry Breaking in the Congest Model: Time- and Message-Efficient Algorithms for Ruling SetsabstractWe study local symmetry breaking problems in the Congest model, focusing on ruling set problems, which generalize the fundamental Maximal Independent Set (MIS) problem. The time (round) complexity of MIS (and ruling sets) have attracted much attention in the Local model. Indeed, recent results (Barenboim et al., FOCS 2012, Ghaffari SODA 2016) for the MIS problem have tried to break the long-standing O(log n)-round "barrier" achieved by Luby's algorithm, but these yield o(log n)-round complexity only when the maximum degree Delta is somewhat small relative to n. More importantly, these results apply only in the Local model. In fact, the best known time bound in the Congest model is still O(log n) (via Luby's algorithm) even for moderately small Delta (i.e., for Delta = Omega(log n) and Delta = o(n)). Furthermore, message complexity has been largely ignored in the context of local symmetry breaking. Luby's algorithm takes O(m) messages on m-edge graphs and this is the best known bound with respect to messages. Our work is motivated by the following central question: can we break the Theta(log n) time complexity barrier and the Theta(m) message complexity barrier in the Congest model for MIS or closely-related symmetry breaking problems? This paper presents progress towards this question for the distributed ruling set problem in the Congest model. A beta-ruling set is an independent set such that every node in the graph is at most beta hops from a node in the independent set. We present the following results: - Time Complexity: We show that we can break the O(log n) "barrier" for 2- and 3-ruling sets. We compute 3-ruling sets in O(log n/log log n) rounds with high probability (whp). More generally we show that 2-ruling sets can be computed in O(log Delta (log n)^(1/2 + epsilon) + log n/log log n) rounds for any epsilon > 0, which is o(log n) for a wide range of Delta values (e.g., Delta = 2^(log n)^(1/2-epsilon)). These are the first 2- and 3-ruling set algorithms to improve over the O(log n)-round complexity of Luby's algorithm in the Congest model. - Message Complexity: We show an Omega(n^2) lower bound on the message complexity of computing an MIS (i.e., 1-ruling set) which holds also for randomized algorithms and present a contrast to this by showing a randomized algorithm for 2-ruling sets that, whp, uses only O(n log^2 n) messages and runs in O(Delta log n) rounds. This is the first message-efficient algorithm known for ruling sets, which has message complexity nearly linear in n (which is optimal up to a polylogarithmic factor). Shreyas Pai, Gopal Pandurangan, Sriram V. Pemmaraju, Talal Riaz, Peter Robinson 0002 |
DISC | 5 |
| 2016 | Fast Distributed Algorithms for Connectivity and MST in Large GraphsabstractMotivated by the increasing need to understand the algorithmic foundations of distributed large-scale graph computations, we study a number of fundamental graph problems in a message-passing model for distributed computing where k ≥ 2 machines jointly perform computations on graphs with n nodes (typically, n gg k). The input graph is assumed to be initially randomly partitioned among the k machines, a common implementation in many real-world systems. Communication is point-to-point, and the goal is to minimize the number of communication rounds of the computation. Our main result is an (almost) optimal distributed randomized algorithm for graph connectivity. Our algorithm runs in ~O(n/k2) rounds (~O notation hides a polylog(n) factor and an additive polylog(n) term). This improves over the best previously known bound of ~O(n/k) [Klauck et al., SODA 2015], and is optimal (up to a polylogarithmic factor) in view of an existing lower bound of ~Ω(n/k2). Our improved algorithm uses a bunch of techniques, including linear graph sketching, that prove useful in the design of efficient distributed graph algorithms. We then present fast randomized algorithms for computing minimum spanning trees, (approximate) min-cuts, and for many graph verification problems. All these algorithms take ~O(n/k2) rounds, and are optimal up to polylogarithmic factors. We also show an almost matching lower bound of ~Ω(n/k2) for many graph verification problems using lower bounds in random-partition communication complexity. Gopal Pandurangan, Peter Robinson 0002, Michele Scquizzato |
SPAA | 2 |
| 2016 | DEX: self-healing expandersabstractWe present a fully-distributed self-healing algorithm dex that maintains a constant degree expander network in a dynamic setting. To the best of our knowledge, our algorithm provides the first efficient distributed construction of expanders—whose expansion properties hold deterministically—that works even under an all-powerful adaptive adversary that controls the dynamic changes to the network (the adversary has unlimited computational power and knowledge of the entire network state, can decide which nodes join and leave and at what time, and knows the past random choices made by the algorithm). Previous distributed expander constructions typically provide only probabilistic guarantees on the network expansion which rapidly degrade in a dynamic setting; in particular, the expansion properties can degrade even more rapidly under adversarial insertions and deletions. Our algorithm provides efficient maintenance and incurs a low overhead per insertion/deletion by an adaptive adversary: only $$O(\log n)$$ rounds and $$O(\log n)$$ messages are needed with high probability (n is the number of nodes currently in the network). The algorithm requires only a constant number of topology changes. Moreover, our algorithm allows for an efficient implementation and maintenance of a distributed hash table on top of dex with only a constant additional overhead. Our results are a step towards implementing efficient self-healing networks that have guaranteed properties (constant bounded degree and expansion) despite dynamic changes. Gopal Pandurangan, Peter Robinson 0002, Amitabh Trehan |
Distributed Comput. | 2 |
| 2015 | Enabling Robust and Efficient Distributed Computation in Dynamic Peer-to-Peer NetworksabstractMotivated by the need for designing efficient and robust fully-distributed computation in highly dynamic networks such as Peer-to-Peer (P2P) networks, we study distributed protocols for constructing and maintaining dynamic network topologies with good expansion properties. Our goal is to maintain a sparse (bounded degree) expander topology despite heavy churn (i.e., Nodes joining and leaving the network continuously over time). We assume that the churn is controlled by an adversary that has complete knowledge and control of what nodes join and leave and at what time and has unlimited computational power, but is oblivious to the random choices made by the algorithm. Our main contribution is a randomized distributed protocol that guarantees with high probability the maintenance of a constant degree graph with high expansion even under continuous high adversarial churn. Our protocol can tolerate a churn rate of up to O(n/polylog(n)) per round (where n is the stable network size). Our protocol is efficient, lightweight, and scalable, and it incurs only O(polylog(n)) overhead for topology maintenance: only polylogarithmic(in n) bits needs to be processed and sent by each node per round and any node's computation cost per round is also polylogarithmic. The given protocol is a fundamental ingredient that is needed for the design of efficient fully-distributed algorithms for solving fundamental distributed computing problems such as agreement, leader election, search, and storage in highly dynamic P2P networks and enables fast and scalable algorithms for these problems that can tolerate a large amount of churn. John Augustine 0001, Gopal Pandurangan, Peter Robinson 0002, Scott T. Roche, Eli Upfal |
FOCS | 3 |
| 2015 | Distributed Computation of Large-scale Graph ProblemsabstractMotivated by the increasing need for fast distributed processing of large-scale graphs such as the Web graph and various social networks, we study a number of fundamental graph problems in the message-passing model, where we have k machines that jointly perform computation on an arbitrary n-node (typically, n ≫ k) input graph. The graph is assumed to be randomly partitioned among the k ≥ 2 machines (a common implementation in many real world systems). The communication is point-to-point, and the goal is to minimize the time complexity, i.e., the number of communication rounds, of solving various fundamental graph problems. We present lower bounds that quantify the fundamental time limitations of distributively solving graph problems. We first show a lower bound of Ω(n/k) rounds for computing a spanning tree (ST) of the input graph. This result also implies the same bound for other fundamental problems such as computing a minimum spanning tree (MST), breadth-first tree (BFS), and shortest paths tree (SPT). We also show an Ω(n/k2) lower bound for connectivity, ST verification and other related problems. Our lower bounds develop and use new bounds in random-partition communication complexity. To complement our lower bounds, we also give algorithms for various fundamental graph problems, e.g., PageRank, MST, connectivity, ST verification, shortest paths, cuts, spanners, covering problems, densest subgraph, subgraph isomorphism, finding triangles, etc. We show that problems such as PageRank, MST, connectivity, and graph covering can be solved in Õ(n/k) time (the notation Õ hides polylog(n) factors and an additive polylog(n) term); this shows that one can achieve almost linear (in k) speedup, whereas for shortest paths, we present algorithms that run in time (for (1 + ε)-factor approximation) and in time (for O(log n)-factor approximation) respectively. Our results step towards understanding the complexity of distributively solving large-scale graph problems. Hartmut Klauck, Danupon Nanongkai, Gopal Pandurangan, Peter Robinson 0002 |
SODA | 4 |
| 2015 | Fast Byzantine Leader Election in Dynamic Networks
John Augustine 0001, Gopal Pandurangan, Peter Robinson 0002 |
DISC | 3 |
| 2015 | On the Complexity of Universal Leader ElectionabstractElecting a leader is a fundamental task in distributed computing. In its implicit version, only the leader must know who is the elected leader. This article focuses on studying the message and time complexity of randomized implicit leader election in synchronous distributed networks. Surprisingly, the most “obvious” complexity bounds have not been proven for randomized algorithms. In particular, the seemingly obvious lower bounds of Ω( m ) messages, where m is the number of edges in the network, and Ω( D ) time, where D is the network diameter, are nontrivial to show for randomized (Monte Carlo) algorithms. (Recent results, showing that even Ω( n ), where n is the number of nodes in the network, is not a lower bound on the messages in complete networks, make the above bounds somewhat less obvious). To the best of our knowledge, these basic lower bounds have not been established even for deterministic algorithms, except for the restricted case of comparison algorithms, where it was also required that nodes may not wake up spontaneously and that D and n were not known. We establish these fundamental lower bounds in this article for the general case, even for randomized Monte Carlo algorithms. Our lower bounds are universal in the sense that they hold for all universal algorithms (namely, algorithms that work for all graphs), apply to every D , m , and n , and hold even if D , m , and n are known, all the nodes wake up simultaneously, and the algorithms can make any use of node's identities. To show that these bounds are tight, we present an O ( m ) messages algorithm. An O ( D ) time leader election algorithm is known. A slight adaptation of our lower bound technique gives rise to an Ω( m ) message lower bound for randomized broadcast algorithms. An interesting fundamental problem is whether both upper bounds (messages and time) can be reached simultaneously in the randomized setting for all graphs. The answer is known to be negative in the deterministic setting. We answer this problem partially by presenting a randomized algorithm that matches both complexities in some cases. This already separates (for some cases) randomized algorithms from deterministic ones. As first steps towards the general case, we present several universal leader election algorithms with bounds that tradeoff messages versus time. We view our results as a step towards understanding the complexity of universal leader election in distributed networks. Shay Kutten, Gopal Pandurangan, David Peleg, Peter Robinson 0002, Amitabh Trehan |
J. ACM | 4 |
| 2015 | Distributed agreement in dynamic peer-to-peer networks
John Augustine 0001, Gopal Pandurangan, Peter Robinson 0002, Eli Upfal |
J. Comput. Syst. Sci. | 3 |
| 2015 | Sublinear bounds for randomized leader election
Shay Kutten, Gopal Pandurangan, David Peleg, Peter Robinson 0002, Amitabh Trehan |
Theor. Comput. Sci. | 4 |
| 2014 | DEX: Self-Healing ExpandersabstractWe present a fully-distributed self-healing algorithm DEX, that maintains a constant degree expander network in a dynamic setting. To the best of our knowledge, our algorithm provides the first efficient distributed construction of expanders - whose expansion properties hold deterministically - that works even under an all-powerful adaptive adversary that controls the dynamic changes to the network (the adversary has unlimited computational power and knowledge of the entire network state, can decide which nodes join and leave and at what time, and knows the past random choices made by the algorithm). Previous distributed expander constructions typically provide only probabilistic guarantees on the network expansion which rapidly degrade in a dynamic setting, in particular, the expansion properties can degrade even more rapidly under adversarial insertions and deletions. Our algorithm provides efficient maintenance and incurs a low overhead per insertion/deletion by an adaptive adversary: only O(log n) rounds and O(log n) messages are needed with high probability (n is the number of nodes currently in the network). The algorithm requires only a constant number of topology changes. Moreover, our algorithm allows for an efficient implementation and maintenance of a distributed hash table (DHT) on top of DEX, with only a constant additional overhead. Our results are a step towards implementing efficient self-healing networks that have guaranteed properties (constant bounded degree and expansion) despite dynamic changes. Gopal Pandurangan, Peter Robinson 0002, Amitabh Trehan |
IPDPS | 2 |
| 2014 | Brief announcement: gracefully degrading consensus and k-set agreement under dynamic link failuresabstractWe present a k-set agreement algorithm for synchronous dynamic distributed systems with unidirectional links controlled by an omniscient adversary. Our algorithm automatically adapts to the actual network properties: If the network is sufficiently well-connected, it solves consensus, while degrading gracefully to general k-set agreement in less well-behaved runs. The algorithm is oblivious to the maximum number of system-wide decision values k, which is bounded by the number of certain strongly connected components occurring in the dynamically changing network in a run. Related impossibility results reveal that this bound is close to the solvability border for k-set agreement. To the best of our knowledge, this is the first consensus algorithm that degrades in a graceful way in a dynamic network. Manfred Schwarz, Kyrill Winkler, Ulrich Schmid 0001, Martin Biely, Peter Robinson 0002 |
PODC | 5 |
| 2014 | Distributed Symmetry Breaking in Hypergraphs
Shay Kutten, Danupon Nanongkai, Gopal Pandurangan, Peter Robinson 0002 |
DISC | 4 |
| 2014 | The Generalized Loneliness Detector and Weak System Models for k-Set AgreementabstractThis paper presents two weak partially synchronous system models Manti(n-k)and Msink(n-k), which are just strong enough for solving k-set agreement: We introduce the generalized (n-k)-loneliness failure detector L(k), which we first prove to be sufficient for solving k-set agreement, and show that L(k) but not L(k-1) can be implemented in both models. Manti(n-k)and Msink(n-k)are hence the first message passing models that lie between models where Ω (and therefore consensus) can be implemented and the purely asynchronous model. We also address k-set agreement in anonymous systems, that is, in systems where (unique) process identifiers are not available. Since our novel k -set agreement algorithm using L(k) also works in anonymous systems, it turns out that the loneliness failure detector L=L(n-1) introduced by Delporte et al. is also the weakest failure detector for set agreement in anonymous systems. Finally, we analyze the relationship between L(k) and other failure detectors suitable for solving k-set agreement. Martin Biely, Peter Robinson 0002, Ulrich Schmid 0001 |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2013 | Efficient Computation of Balanced Structures
David G. Harris 0001, Ehab Morsy, Gopal Pandurangan, Peter Robinson 0002, Aravind Srinivasan |
ICALP (2) | 4 |
| 2013 | Fast byzantine agreement in dynamic networksabstractWe study Byzantine agreement in dynamic networks where topology can change from round to round and nodes can also experience heavy churn (i.e., nodes can join and leave the network continuously over time). Our main contributions are randomized distributed algorithms that achieve almost-everywhere Byzantine agreement with high probability even under a large number of adaptively chosen Byzantine nodes and continuous adversarial churn in a number of rounds that is polylogarithmic in n (where n is the stable network size). We show that our algorithms are essentially optimal (up to polylogarithmic factors) with respect to the amount of Byzantine nodes and churn rate that they can tolerate by showing a lower bound. In particular, we present the following results: John Augustine 0001, Gopal Pandurangan, Peter Robinson 0002 |
PODC | 3 |
| 2013 | On the complexity of universal leader electionabstractElecting a leader is a fundamental task in distributed computing. In its implicit version, only the leader must know who is the elected leader. This paper focuses on studying the message and time complexity of randomized implicit leader election in synchronous distributed networks. Surprisingly, the most "obvious" complexity bounds have not been proven for randomized algorithms. The "obvious" lower bounds of Ω(m) messages (m is the number of edges in the network) and Ω(D) time (D is the network diameter) are non-trivial to show for randomized (Monte Carlo) algorithms. (Recent results that show that even Ω(n) (n is the number of nodes in the network) is not a lower bound on the messages in complete networks, make the above bounds somewhat less obvious). To the best of our knowledge, these basic lower bounds have not been established even for deterministic algorithms (except for the limited case of comparison algorithms, where it was also required that some nodes may not wake up spontaneously, and that D and n were not known). Shay Kutten, Gopal Pandurangan, David Peleg, Peter Robinson 0002, Amitabh Trehan |
PODC | 4 |
| 2013 | Storage and search in dynamic peer-to-peer networksabstractWe study robust and efficient distributed algorithms for searching, storing, and maintaining data in dynamic Peer-to-Peer (P2P) networks. P2P networks are highly dynamic networks that experience heavy node churn (i.e., nodes join and leave the network continuously over time). Our goal is to guarantee, despite high node churn rate, that a large number of nodes in the network can store, retrieve, and maintain a large number of data items. Our main contributions are fast randomized distributed algorithms that guarantee the above with high probability even under high adversarial churn. In particular, we present the following main results: John Augustine 0001, Anisur Rahaman Molla, Ehab Morsy, Gopal Pandurangan, Peter Robinson 0002, Eli Upfal |
SPAA | 5 |
| 2012 | Agreement in Directed Dynamic Networks
Martin Biely, Peter Robinson 0002, Ulrich Schmid 0001 |
SIROCCO | 2 |
| 2012 | Towards robust and efficient computation in dynamic peer-to-peer networksabstractMotivated by the need for robust and fast distributed computation in highly dynamic Peer-to-Peer (P2P) networks, we study algorithms for the fundamental distributed agreement problem. P2P networks are highly dynamic networks that experience heavy node churn (i.e., nodes join and leave the network continuously over time). Our goal is to design fast algorithms (running in a small number of rounds) that guarantee, despite high node churn rate, that almost all nodes reach a stable agreement. Our main contributions are randomized distributed algorithms that guarantee stable almost-everywhere agreement with high probability even under high adversarial churn in a polylogarithmic number of rounds. In particular, we present the following results: 1. An O(log n)-round (n is the stable network size) randomized algorithm that achieves almost-everywhere agreement with high probability under up to linear churn per round (i.e., εn, for some small constant ε > 0), assuming that the churn is controlled by an oblivious adversary (that has complete knowledge and control of what nodes join and leave and at what time and has unlimited computational power, but is oblivious to the random choices made by the algorithm). 2. An O(log m log3 n)-round randomized algorithm that achieves almost-everywhere agreement with high probability under up to ε√n churn per round (for some small ε > 0), where m is the size of the input value domain, that works even under an adaptive adversary (that also knows the past random choices made by the algorithm). Our algorithms are the first-known, fully-distributed, agreement algorithms that work under highly dynamic settings (i.e., high churn rates per step). Furthermore, they are localized (i.e., do not require any global topological knowledge), simple, and easy to implement. These algorithms can serve as building blocks for implementing other non-trivial distributed computing tasks in dynamic P2P networks. John Augustine 0001, Gopal Pandurangan, Peter Robinson 0002, Eli Upfal |
SODA | 3 |
| 2011 | Easy Impossibility Proofs for k-Set Agreement in Message Passing Systems
Martin Biely, Peter Robinson 0002, Ulrich Schmid 0001 |
OPODIS | 2 |
| 2011 | Easy impossibility proofs for k-set agreement in message passing systemsabstractNo abstract available. Martin Biely, Peter Robinson 0002, Ulrich Schmid 0001 |
PODC | 2 |
| 2011 | The Asynchronous Bounded-Cycle modelabstractThis paper shows how synchrony conditions can be added to the purely asynchronous model in a way that avoids any reference to message delays and computing step times, as well as system-wide constraints on execution patterns and network topology. Our Asynchronous Bounded-Cycle (ABC) model just bounds the ratio of the number of forward- and backward-oriented messages in certain ("relevant") cycles in the space-time diagram of an asynchronous execution. We show that clock synchronization and lock-step rounds can be implemented and proved correct in the ABC model, even in the presence of Byzantine failures. Furthermore, we prove that any algorithm working correctly in the partially synchronous Θ-Model also works correctly in the ABC model. In our proof, we first apply a novel method for assigning certain message delays to asynchronous executions, which is based on a variant of Farkas' theorem of linear inequalities and a non-standard cycle space of graphs. Using methods from point-set topology, we then prove that the existence of this delay assignment implies model indistinguishability for time-free safety and liveness properties. We also introduce several weaker variants of the ABC model, and relate our model to the existing partially synchronous system models, in particular, the classic models of Dwork, Lynch and Stockmayer and the query-response model by Mostefaoui, Mourgaya, and Raynal. Finally, we discuss some aspects of the ABC model's applicability in real systems, in particular, in the context of VLSI Systems-on-Chip. Peter Robinson 0002, Ulrich Schmid 0001 |
Theor. Comput. Sci. | 1 |
| 2009 | Weak Synchrony Models and Failure Detectors for Message Passing (k-)Set Agreement
Martin Biely, Peter Robinson 0002, Ulrich Schmid 0001 |
OPODIS | 2 |
| 2009 | Brief Announcement: Weak Synchrony Models and Failure Detectors for Message Passing (k-)Set Agreement
Martin Biely, Peter Robinson 0002, Ulrich Schmid 0001 |
DISC | 2 |
| 2008 | The asynchronous bounded-cycle modelabstractIn this paper, we introduce the Asynchronous Bounded-Cycle (ABC) model, which considerably relaxes the Theta-Model proposed by Le Lann and Schmid. The ABC model just bounds the ratio of the number of forward and backward messages in certain cycles in the space-time diagram of an asynchronous execution. It hence avoids any reference to end-to-end delays, allows individual messages to have arbitrary delays, and does not involve global synchrony conditions. We show that clock synchronization and lock-step rounds can easily be implemented and proved correct in the ABC model, even in the presence of Byzantine failures. Moreover, we show that any correct Theta-algorithm also works correctly in the ABC model. Our proof is based on a novel technique for assigning message delays to asynchronous executions, which is of independent interest. Peter Robinson 0002, Ulrich Schmid 0001 |
PODC | 1 |
| 2008 | The Asynchronous Bounded-Cycle Model
Peter Robinson 0002, Ulrich Schmid 0001 |
SSS | 1 |