Network Topology and Communication-Computation Tradeoffs in Decentralized Optimization

Angelia Nedić, Alex Olshevsky, Michael G. Rabbat

I Introduction

Because each node only has access to local information, the nodes must communicate over a network to find a minimizer of f(x)f(x). Multi-agent consensus optimization algorithms are iterative, where each iteration typically involves some local computation followed by communication over the network.

Centralized gradient descent for minimizing the function f(x)f(x) starts with an initial value x0x^{0} and recursively updates it for k=1,2,…,k=1,2,\dots, by setting

where α1,α2,…,\alpha_{1},\alpha_{2},\dots, is a sequence of positive scalar step-sizes. When ff is convex, it has a unique minimum, and it is well-known that, for appropriate choices of the step-sizes αk\alpha_{k}, the sequence of values f(xk)f(x^{k}) converges to this minimum .

and where the gradient ∇fi(x)\nabla f_{i}(x) can only be evaluated at agent ii. There are a variety of distributed architectures one may consider in this setting, and we discuss three here: 1) the master-worker architecture, 2) the fully-connected architecture, and 3) a general, connected architecture. They are depicted in Fig. 1 and described next.

When f(x)f(x) decomposes as in (2), then the gradient also decomposes, and

that is, the gradient of the overall objective is the average of the gradients of the local objectives.

In a master-worker architecture, one node acts as the master (sometimes also called the fusion center), maintaining the authoritative copy of the optimization variable xkx^{k}. At each iteration, it sends xkx^{k} to every agent, and agent ii returns ∇fi(xk)\nabla f_{i}(x^{k}) to the master. The master averages the gradients it receives from the agents, and once it has received a gradient from every agent it can perform the gradient descent update (1), before proceeding to the next iteration.

The master-worker architecture is useful in that it is relatively simple to implement. However in many applications it may be unattractive or impractical for a variety of reasons. First, as the number of nodes nn grows large, the master node may become a communication bottleneck if it has limited communication resources (e.g., if its bandwidth does not grow linearly with the size of the network), and at the same time, scaling the bandwidth of the master with the size of the network may be expensive or impractical. Also, the master node may become a robustness bottleneck, in the sense that if the master node fails then the entire network fails. In addition, in many scenarios it may not be practical to have a single master node that communicates with all agents. For example, if agents are low-power devices communicating via wireless radios, then two devices may only be able to communicate if they are nearby and it may not be practical to have all nodes within the required proximity of the master.

I-A2 The fully-connected architecture

Thus, using the average of the gradients received from its neighbors, node ii can update

and xi1x_{i}^{1} is exactly equivalent to having taken one step of centralized gradient descent. Furthermore, the values xi1x_{i}^{1} will be identical at all nodes, and so we can repeat this process recursively to essentially implement centralized gradient descent exactly in a distributed manner.

For the fully-connected architecture just describedWe refer to it as fully-connected because every node communicates with every other node at each iteration., each node acts like a master in the master-worker architecture, and so the fully-connected architecture suffers from the same issues as the master-worker architecture. Moreover, the communication overhead of having all nodes communicate at every iteration is even worse than the master-worker architecture (it grows quadratically in the number of nodes nn, whereas the communication overhead was linear in nn for the master-worker architecture). Nevertheless, the fully-connected architecture provides a conceptual transition from the master-worker architecture to general connected (but not fully-connected) architectures.

I-A3 General multi-agent architectures

where ∣Ni∣|N_{i}| is the size of node ii’s neighborhood.

This approach given in (4) is prototypical of most multi-agent optimization algorithms, in that the update equation can be implemented in the following steps, which are executed in parallel at every node, i=1,…,ni=1,\dots,n:

Node ii locally computes ∇fi(xik)\nabla f_{i}(x_{i}^{k}).

Node ii transmits its gradient ∇fi(xik)\nabla f_{i}(x_{i}^{k}) and receives gradients ∇fj(xjk)\nabla f_{j}(x_{j}^{k}) from its neighbors j∈Nij\in N_{i}.

Node ii uses this new information to compute the new value xik+1x_{i}^{k+1}, e.g., via equation (4).

Different multi-agent optimization algorithms may differ in terms of what information gets exchanged in the second step, and in the precise way they compute the update in the last step, not necessarily using (4), as well as in the assumptions they make about the local objective functions fif_{i} or the communication topology, captured by the neighborhoods NiN_{i}. For example: the communication topology may be static or it may vary from iteration to iteration; communications may be undirected (node ii receives messages from node jj if and only if jj also receives messages from ii) or directed.

The general multi-agent approach to implementing gradient descent, given in (4), also raises a set of issues which did not come up in the other architectures. Since the master-worker and fully-connected architectures exactly implement gradient descent, the well-established convergence theory for gradient descent directly applies to those architectures. However, when nodes update using the rule (4), they no longer exactly implement centralized gradient descent, because they use a search direction

which is the average of a subset, rather than all, of the gradients at other nodes. Thus, after the first iteration, the local values xj1x_{j}^{1} at different nodes are no longer equivalent. Subsequently, at the next iteration, the local gradients being averaged at node ii will have been evaluated at different values xj1x_{j}^{1}. Therefore, there are a few ways in which the values produced by multi-agent gradient descent and deviates from centralized gradient descent. One may hope that, under the right conditions, the values at different nodes will not be too different from each other and that the local search directions will be sufficiently similar to the gradient search direction that the nodes still converge to (and agree on!) a minimizer of f(x)f(x).

Indeed, we will see that we can identify a variety of conditions under which multi-agent optimization algorithms are guaranteed to converge, and we can precisely quantify how the convergence rate differs from that of centralized gradient descent. Most often, this difference depends directly on the communication topology. In many applications of interest, either it is not possible or one does not allow each node to communicate with every other node. The connectivity of the network (i.e., which pairs of nodes may communicate directly with each other) can be represented as a graph with nn vertices and with an edge from jj to ii if node jj receives messages from node ii. We will see that the communication network topology plays a key role in the convergence theory of multi-agent optimization methods in that it may limit the flow of information between distant nodes and thereby hinder convergence.

During the past decade, multi-agent consensus optimization has been the subject of intense interest, motivated by a variety of applications which we discuss in next.

I-B Motivating Applications

The general multi-agent optimization problem described above was originally introduced and studied in the 1980’s in the context of parallel and distributed numerical methods . The surge of interest in multi-agent convex optimization during the past decade has been fueled by a variety of applications where a network of autonomous agents must coordinate to achieve a common objective. We describe three such examples next; for a survey describing additional applications of multi-agent methods for coordination, see .

Consider a wireless sensor network with nn nodes where node ii has a measurement ζi\zeta_{i} which is modeled as a random variable with density p(ζi∣x)p(\zeta_{i}|x) depending on unknown parameters xx. For example, the network may be deployed to monitor a remote or difficult to reach location, and the estimate of xx may be used for scientific observation (e.g., bird migration patterns) or for detecting events (e.g., avalanches) . In many applications of sensor networks, uncertainty is primarily due to thermal measurement noise introduced at the sensor itself, and so it is reasonable to assume that the observations {ζi}i=1n\{\zeta_{i}\}_{i=1}^{n} are conditionally independent given the model parameters xx. In this case, the maximum likelihood estimate of xx can obtained by solving

which can be addressed by using multi-agent consensus optimization methods with fi(x)=−log⁡p(ζi∣x)f_{i}(x)=-\log p(\zeta_{i}|x).

In this example, the data are already being gathered in a decentralized manner at different sensors. When the data dimension is large (e.g., for image or video sensors), it can be more efficient to perform decentralized estimation and simply transmit the estimate of xx to the end-user, rather than transmitting the raw data and then performing centralized estimation . Similarly, even if the data dimension is not large, if the number nn of nodes in the network is large, it may still be more efficient to perform decentralized estimation rather than sending raw data to a fusion center, since the fusion center will become a bottleneck.

When the nodes communicate over a wireless network, whether or not a given pair of nodes can directly communicate is typically a function of their physical proximity as well as other factors (e.g., fading, shadowing) affecting the wireless channel, which may possibly result in time-varying and directed network connectivity.

I-B2 Big data and machine learning

Many methods for supervised learning (e.g., classification or regression) can also be formulated as fitting a model to data. This task may generally be expressed as finding model parameters xx by solving

where the loss function lj(x)l_{j}(x) measures how well the model with parameters xx describes the jjth training instance, and the training data set contains mm instances in total. For many popular machine learning models—including linear regression, logistic regression, ridge regression, the LASSO, support vector machines, and their variants—the corresponding loss function is convex .

When mm is large, it may not be possible to store the training data on a single server, or it may be desirable, for other reasons, to partition the training data across multiple nodes (e.g., to speedup training by exploiting parallel computing resources, or because the data is being gathered and/or stored at geographically distant locations). In this case, the training task (5) can be solved using multi-agent optimization with local objectives of the form

where Ji\mathcal{J}_{i} is the set of indices of training instances at node ii.

In this setting, where the nodes are typically servers communicating over a wired network, it may be feasible for every node to send and receive messages from all other nodes. However, it is often still preferable, for a variety of reasons, to run multi-agent algorithms over a network with sparser connectivity. Communicating a message takes time, and reducing the number of edges in the communication graph at any iteration corresponds to reducing the number of messages to be transmitted. This results in iterations that take less time and also that consume less bandwidth.

I-B3 Multi-robot systems

Similar to the previous example, multi-agent methods have attracted attention in applications requiring the coordination of multiple robots because they naturally lead to decentralized solutions. One well-studied problem arising in such systems is that of rendezvous—collectively deciding on a meeting time and location. When the robots have different battery levels or are otherwise heterogeneous, it may be desirable to design a rendezvous time and place, and corresponding control trajectories, which minimize the energy to be expended collectively by all robots. This can be formulated as a multi-agent optimization problem where the local objective fi(x)f_{i}(x) at agent ii quantifies the energy to be expended by agent ii and xx encodes the time and place for rendezvous .

When robots communicate over a wireless network, the network connectivity will be dependent on the proximity of nodes as well as other factors affecting channel conditions, similar to in the first example. Moreover, as the robots move, the network connectivity is likely to change. It may be desirable to ensure that a certain minimal level of network connectivity is maintained while the multi-robot system performs its task, and such requirements can be enforced by introducing constraints in the optimization formulation .

I-C Outline of the rest of the paper

The purpose of this article is to provide an overview of the main advances in this field, highlighting the state-of-the-art methods and their analyses, and pointing out open questions. During the past decade, a vast literature has amassed on multi-agent optimization and related methods, and we do not attempt to provide an exhaustive review (which, in any case, would not be feasible in the space of one article). Rather, in addition to describing the main advances and results leading to the current state-of-the-art, we also seek to provide an accessible survey of theoretical techniques arising in the analysis of multi-agent optimization methods.

As we have already seen, decentralized averaging algorithms—where each node initially holds a number or vector, and the aim is to calculate the average at every node—form a fundamental building block of multi-agent optimization methods. Section II reviews decentralized averaging algorithms and their convergence theory in the setting of undirected graphs, where node ii receives message from node jj if and only if jj also receives messages from ii,. The main results of this section provide conditions under which decentralized averaging algorithms converge asymptotically to the exact average, and they quantify how close the values at each node are to the exact average after a finite number kk of iterations. We initially consider the general scenario where the communication topology is time-varying, finding that a sufficient condition for convergence is that the topology be sufficiently well-connected over periodic windows of time. Then, for the particular case where the communication topology is static, we present stronger results illustrating how the rate of convergence depends intimately on properties of the communication topology.

Section IV then describes how decentralized averaging and multi-agent optimization methods can be extended to run over networks with directed connectivity (i.e., where node ii may receive messages from jj although jj does not receive messages from ii). The key technique we study, which enables this extension, is the so-called “push-sum” approach. We provide a novel, concise analysis of the push-sum method for decentralized averaging, and then we describe how it can be used to obtain decentralized optimization methods.

Section V discusses a variety of ways that the basic approaches described in Sections III and IV can be extended. For example, in both Sections III and IV we limit our attention to methods for unconstrained optimization problems running in a synchronous manner. Sections V discusses how to handle constrained optimization problem and how multi-agent optimization methods can be implemented in an asynchronous manner. It also describes other extensions, such as handling stochastic gradient information or online optimization (where the objective function varies over time), and discusses connections to other methods for distributed optimization.

We conclude in Section VI and highlight some open problems.

I-D Notation

Given a sequence of stochastic matrices A0,A1,A2…A^{0},A^{1},A^{2}\ldots, for k>lk>l, we denote by Ak:lA^{k:l} the product of matrices AkA^{k} to AlA^{l} inclusive, i.e.,

We say that a matrix sequence is BB-strongly-connected if the graph with vertex set {1,…,n}\{1,\ldots,n\} and edge set

is strongly connected for each l=0,1,2,…l=0,1,2,\ldots. Intuitively, we partition the iterations k=1,2,…k=1,2,\ldots into consecutive blocks of length BB, and the sequence is BB-strongly-connected when the graph obtained by unioning the edges within each block is always strongly connected. When the graph sequence is undirected, we will simply say BB-connected.

The out-neighbors of node ii at iteration kk refers to the set of nodes that can receive messages from it,

and similarly, the in-neighbors of ii at iteration kk are the nodes from which ii receives messages,

We assume that ii is always an neighbor of itself (i.e., the diagonal entries of AkA^{k} are non-zero), which means that we always have i∈Niout,ki\in N^{\text{out},k}_{i} and i∈Niin,ki\in N^{\text{in},k}_{i}. When the graph is not time-varying, we simply refer to the out-neighbors NioutN^{\text{out}}_{i} and in-neighbors NiinN^{\text{in}}_{i}. When the graph is undirected, the sets of in-neighbors and out-neighbors are identical, so we will simply refer to the neighbors NiN_{i} of node ii, or NikN_{i}^{k} in the time-varying setting. The out-degree of node ii at iteration kk is defined as the cardinality of Niout,kN^{\text{out},k}_{i} and is denoted by diout,k=∣Niout,k∣d^{\text{out},k}_{i}=|N^{\text{out},k}_{i}|. Similarly, diin,kd^{\text{in},k}_{i}, dioutd^{\text{out}}_{i}, diind^{\text{in}}_{i}, dikd_{i}^{k}, and did_{i} denote the cardinalities, respectively, of the sets Niin,kN^{\text{in},k}_{i}, NioutN^{\text{out}}_{i}, NiinN^{\text{in}}_{i}, NikN_{i}^{k}, and NiN_{i}.

II Decentralized Averaging over Undirected Graphs

This section reviews methods for decentralized averaging that will form a key building block in our subsequent discussion of methods for multi-agent optimization.

We begin by examining the linear consensus process defined as

For example, consider a collection of nodes interconnected in a directed graph and suppose node ii holds the ii’th coordinate of the vector xkx^{k}. Consider the following update rule: at step kk, node ii broadcasts the value xikx_{i}^{k} to its out-neighbors, receives values xjkx_{j}^{k} from its in-neighbors, and sets xik+1x_{i}^{k+1} to be the average of the messages it has received, so that

This is sometimes called the equal neighbor iteration, and by stacking up the variables xikx_{i}^{k} into the vector xkx^{k} it can be written in the form of Eq. (6) with an appropriate choice for the matrix AkA^{k}.

Intuitively, we may think of the equal-neighbor updates in terms of an opinion dynamics process over a network wherein node ii repeatedly revises it’s opinion vector xikx_{i}^{k} by averaging the opinions of it’s neighbors. As we will see later, under some relatively mild conditions this process converges to a state where all opinions are identical, explaining why Eq. (6) is usually referred to as the “consensus iteration.”

Over undirected graphs, an alternative popular choice of update rule is to set

where ϵ>0\epsilon>0 is sufficiently small. Unfortunately, finding an appropriate choice of ϵ\epsilon to guarantee convergence of this iteration can be bothersome (especially when the graphs are time-varying), and it generally requires knowing an upper bound on the degrees of nodes in the network.

Another possibility (when the underlying graphs are undirected) is the so-called Metropolis update

The Metropolis update requires node ii to broadcast both xikx_{i}^{k} and its degree dikd_{i}^{k} to its neighbors. Observe that the Metropolis update of Eq. (7) can be written in the form of Eq. (6) where the matrices AkA^{k} are doubly stochastic.

A variation on this is the so-called lazy Metropolis update,

with the key difference being the factor of 22 in the denominator. It is standard convention within the probability literature that such updates are called “lazy,” since they move half as much per iteration. As we will see later, the lazy Metropolis iteration possesses a number of attractive convergence properties.

As we have alluded to above, under certain technical conditions, the iteration of Eq. (6) results in consensus, meaning that all of the xikx_{i}^{k} (for i=1,…,ni=1,\dots,n) approach the same value as k→∞k\rightarrow\infty. We describe one such condition next. The key properties needed to ensure asymptotic consensus are that the matrices AkA^{k} should exhibit sufficient connectivity and aperiodicity (in the long term). In the following, we use the shorthand GkG^{k} for GAkG_{A^{k}}, the graph corresponding to the matrix AkA^{k}. The starting point of our analysis is the following assumption.

The sequence of directed graphs G0,G1,G2,…G^{0},G^{1},G^{2},\ldots is BB-strongly-connected. Moreover, each graph GkG^{k} has a self-loop at every node.

As the next theorem shows, a variation on this assumption suffices to ensure that the update of Eq. (6) converges to consensus.

[Consensus Convergence over Time-Varying Graphs] Suppose the sequence of stochastic matrices A0,A1,A2,…A^{0},A^{1},A^{2},\ldots has the property that there exists an α>0\alpha>0 such that the sequence of graphs G[A0]α,G[A1]α,G[A2]α,…G_{[A^{0}]_{\alpha}},G_{[A^{1}]_{\alpha}},G_{[A^{2}]_{\alpha}},\ldots satisfies Assumption 1. Then x(t)x(t) converges to a limit in span{1}{\rm span}\{{\bf 1}\} and the convergence is geometricA sequence of vectors z(t)z(t) converges to the limit zz geometrically if ∥z(t)−z∥2≤cαt\|z(t)-z\|_{2}\leq c\alpha^{t} for some c≥0c\geq 0, and 0<α<10<\alpha<1.. Moreover, if all the matrices AkA^{k} are doubly stochastic then for all i=1,…,ni=1,\ldots,n,

On an intuitive level, the theorem works by ensuring two things. First, there needs to be an assumption of repeated connectivity in the system over the long-term, and this is what the strong-connectivity condition does. Furthermore, thresholding the weights at some strictly positive α\alpha rules out counterexamples where the weights decay to zero with time. Secondly, the assumption that every node has a self-loop rules out a class of counterexamples where the underlying graph is bipartite and the underlying opinions oscillateIndeed, observe that if we do not require that each node has a self-loop, the dynamics x^{k+1}=\left(\begin{array}[]{cc}0&1\\ 1&0\end{array}\right)x^{k}, started at x0≠0x^{0}\neq 0, would be a counterexample to Theorem 1.. Once these potential counterexamples are ruled out, Theorem 1 guarantees convergence.

We now turn to the proof of this theorem, which while being reasonably short, still builds on a sequence of preliminary lemmas and definitions which we present first. Given a sequence of directed graphs G0,G1,G2,…,G^{0},G^{1},G^{2},\ldots,, we say that node bb is reachable from node aa in time period k:lk:l if there exists a sequence of directed edges ek,ek−1,…,el+1,ele^{k},e^{k-1},\ldots,e^{l+1},e^{l}, such that: (i) eje^{j} is present in GjG^{j} for all j=l,…,kj=l,\dots,k, (ii) the origin of ele^{l} is aa, (iii) the destination of eke^{k} is bb. Note that this is the same as stating that [Wk:l]ba>0[W^{k:l}]_{ba}>0 if the matrices WkW^{k} are nonnegative with [Wk]ij>0[W^{k}]_{ij}>0 if and only if (j,i)(j,i) belongs to GkG^{k}. We use Nk:l(a)N^{k:l}(a) to denote the set of nodes reachable from node aa in time period k:lk:l.

The first lemma discusses the implications of Assumption 1 for products of the matrices AkA^{k}.

Suppose A0,A1,A2,…A^{0},A^{1},A^{2},\ldots is a sequence of nonnegative matrices with the property that there exists α>0\alpha>0 such that the sequence of graphs G[A0]α,G[A1]α,G[A2]α,…G_{[A^{0}]_{\alpha}},G_{[A^{1}]_{\alpha}},G_{[A^{2}]_{\alpha}},\ldots satisfies Assumption 1 . Then for any integer ll, A(l+n)B−1:lBA^{(l+n)B-1:lB} is a strictly positive matrix. In fact, every entry of A(l+n)B−1:lBA^{(l+n)B-1:lB} is at least αnB\alpha^{nB}.

The proof, given next, is a mathematical formalization of the observation that sufficiently long paths exist between any two nodes.

Consider the set of nodes reachable from node ii in time period kstartk_{\rm start} to kfinishk_{\rm finish} in the graph sequence G[A0]α,G[A1]α,G[A2]α,…G_{[A^{0}]_{\alpha}},G_{[A^{1}]_{\alpha}},G_{[A^{2}]_{\alpha}},\ldots, and denote this set by Nkfinish:kstart(i)N^{k_{\rm finish}:k_{\rm start}}(i). Since each of these graphs has a self-loop at every node by Assumption 1, the reachable set can only be enlarged, i.e.,

A further immediate consequence of Assumption 1 is that if NmB−1:lB(i)≠{1,…,n}N^{mB-1:lB}(i)\neq\{1,\ldots,n\}, then N(m+1)B−1:lB(i)N^{(m+1)B-1:lB}(i) is strictly larger than NmB−1:lB(i)N^{mB-1:lB}(i), because during times (m+1)B−1:mB(m+1)B-1:mB there is an edge in some [Gi]α[G^{i}]_{\alpha} leading from the set of nodes already reachable from ii to those not already reachable from ii. Putting together these two properties, we obtain that from lBlB to (l+n)B−1(l+n)B-1 every node is reachable, i.e.,

But since every non-zero entry of [Ak]α[A^{k}]_{\alpha} is at least α\alpha by construction, this implies that A(l+n)B−1:lB≥αnBA^{(l+n)B-1:lB}\geq\alpha^{nB}, and the lemma is proved. ∎

Lemma 2 tells us that, over sufficiently long horizons, the products of the matrices AkA^{k} have entries bounded away from zero. The next lemma discusses what multiplication by such a matrix does to the spread of the values in a vector.

Suppose WW is a stochastic matrix, every entry of which is at least β>0\beta>0. If v=Wuv=Wu then

Without loss of generality, let us assume that the largest entry of uu is u1u_{1} and the smallest entry of uu is unu_{n}. Then, for l∈{1,…,n}l\in\{1,\ldots,n\},

so that for any a,b∈{1,…,n}a,b\in\{1,\ldots,n\}, we have

With these two lemmas in place, we are ready to prove Theorem 1. Our strategy is to apply Lemma 3 repeatedly to show that the spread of the underlying vectors keeps getting smaller.

Since we have assumed that G[A0]α,G[A1]α,G[A2]α,…G_{[A^{0}]_{\alpha}},G_{[A^{1}]_{\alpha}},G_{[A^{2}]_{\alpha}},\ldots satisfy Assumption 1, by Lemma 2, we have that

for all l=0,1,2,…l=0,1,2,\ldots, and i,j=1,…,ni,j=1,\ldots,n. Applying Lemma 3 gives that

Applying this recursively, we obtain that ∣xak−xbk∣→0|x_{a}^{k}-x_{b}^{k}|\rightarrow 0 for all a,b∈{1,…,n}a,b\in\{1,\ldots,n\}.

To obtain further that every xikx_{i}^{k} converges, it suffices to observe that xikx_{i}^{k} lies in the convex hull of the vectors xtx^{t} for t≤kt\leq k. Finally, since each AkA^{k} is doubly stochastic,

where 1{\bf 1} denotes a vector with all entries equal to one, and thus all xikx_{i}^{k} must converge to the initial average. ∎

A potential shortcoming of the proof of Theorem 1 is that the convergence time bounds it leads to tend to scale poorly in terms of the number of nodes nn. We can overcome this shortcoming as illustrated in the following propositions. These results apply to a much narrower class of scenarios, but they tend to provide more effective bounds when they are applicable.

The first step is to introduce a precise notion of convergence time. Let T(n,ϵ,{A0,A2,…,})T\left(n,\epsilon,\{A^{0},A^{2},\ldots,\}\right) denote the first time kk when

In other words, the convergence time is defined as the time until the deviation from the mean shrinks by a factor of ϵ\epsilon. The convergence time is a function of the desired accuracy ϵ\epsilon and of the underlying sequence of matrices. In particular, we emphasize the dependence on the number of nodes, nn. When the sequence of matrices is clear from context, we will simply write T(n,ϵ)T(n,\epsilon).

where each AkA^{k} is a doubly stochastic matrix. Then

where σ2(Al)\sigma_{2}(A^{l}) denotes the second-largest singular value of the matrix AlA^{l}.

We skip the proof, which follows quickly from the definition of singular value.

We adopt the slightly non-standard notation

so that the previous proposition can be conveniently restated as

Recalling that log⁡(1/λ)≤1/(1−λ)\log(1/\lambda)\leq 1/(1-\lambda), a consequence of this equation is that

so the number λ\lambda provides an upper bound on the convergence rate of decentralized averaging.

In general, there is no guarantee that λ<1\lambda<1, and the equations we have derived may be vacuous. Fortunately, it turns out that for the lazy Metropolis matrices on connected graphs, it is true that λ<1\lambda<1, and furthermore, for many families of undirected graphs it is possible to give order-accurate estimates on λ\lambda, which translate into estimates of convergence time. This is captured in the following proposition. Note that all of these bounds should be interpreted as scaling laws, explaining how the convergence time increases as the network size nn increases, when the graphs all come from the same family.

If each AkA^{k} is the lazy Metropolis matrix on the …

…path graph, then T(n,ϵ)=O(n2log⁡(1/ϵ))T(n,\epsilon)=O\left(n^{2}\log(1/\epsilon)\right).

…22-dimensional grid, then T(n,ϵ)=O(nlog⁡nlog⁡(1/ϵ))T(n,\epsilon)=O\left(n\log n\log(1/\epsilon)\right).

…22-dimensional torus, then T(n,ϵ)=O(nlog⁡(1/ϵ))T(n,\epsilon)=O\left(n\log(1/\epsilon)\right).

…kk-dimensional torus, then T(n,ϵ)=O(n2/klog⁡(1/ϵ)).T(n,\epsilon)=O\left(n^{2/k}\log(1/\epsilon)\right).

…star graph, then T(n,ϵ)=O(n2log⁡(1/ϵ))T(n,\epsilon)=O\left(n^{2}\log(1/\epsilon)\right).

…two-star graphsA two-star graph is composed of two star graphs with a link connecting their centers., then T(n,ϵ)=O(n2log⁡(1/ϵ))T(n,\epsilon)=O\left(n^{2}\log(1/\epsilon)\right).

…complete graph, then T(n,ϵ)=O(1)T(n,\epsilon)=O(1).

…expander graph, then T(n,ϵ)=O(log⁡(1/ϵ))T(n,\epsilon)=O(\log(1/\epsilon)).

…Erdős-Rényi random graphAn Erdős-Rényi random graph with nn nodes and parameter pp has a symmetric adjacency matrix AA whose (n2)\binom{n}{2} distinct off-diagonal entries are independent Bernoulli random variables taking the value 1 with probability pp. In this article we focus on the case where p=(1+ε)log⁡(n)/np=(1+\varepsilon)\log(n)/n, where ε>0\varepsilon>0, for which it is known that the random graph is connected with high probability . then with high probabilityA statement is said to hold “with high probability” if the probability of it holding approaches 11 as n→∞n\rightarrow\infty. In this context, nn is the number of nodes of the underlying graph. T(n,ϵ)=O(log⁡(1/ϵ))T(n,\epsilon)=O(\log(1/\epsilon)).

…geometric random graphA geometric random graph is one where nn nodes are placed uniformly and independently in the unit square 2^{2} and two nodes are connected with an edge if their distance is at most rnr_{n}. In this article we focus on the case where rn2=(1+ε)log⁡(n)/nr_{n}^{2}=(1+\varepsilon)\log(n)/n for some ε>0\varepsilon>0, for which it is known that the random graph is connected with high probability ., then with high probability T(n,ϵ)=O(nlog⁡nlog⁡(1/ϵ))T(n,\epsilon)=O(n\log n\log(1/\epsilon)).

…any connected undirected graph, then T(n,ϵ)=O(n2log⁡(1/ϵ)).T(n,\epsilon)=O\left(n^{2}\log(1/\epsilon)\right).

The spectral gap 1/(1−λ)1/(1-\lambda) can be bounded as O(H)O({\cal H}) where H{\cal H} is the largest hitting time of the Markov chain whose probability transition matrix is the lazy Metropolis matrix (Lemma 2.13 of ). We thus only need to bound hitting times on the graphs in question, and these are now standard exercises. For example, the fact that the hitting time on the path graph is quadratic in the number of nodes is essentially the main finding of the standard “gambler’s ruin” exercise (see e.g., Proposition 2.1 of ). The result for the 22-d grid follows from putting together Theorem 2.1 and Theorem 6.1 of . For corresponding results on 22-d and kk-dimensional tori, please see Theorem 5.5 of ; note that we are treating kk as a fixed number and looking at the scaling as a function of the number of nodes nn. Hitting times on star, two-star, and complete graphs are elementary exercises. The result for an expander graph is a consequence of Cheeger’s inequality; see Theorem 6.2.1 in . For Erdős-Rényi graphs the result follows because such graphs are expanders with high probability; see the discussion on page 170 of . For geometric random graphs a bound can be obtained by partitioning the unit square into appropriately-sized regions, thereby reducing to the case of a 22-d grid; see Theorem 1.1 of . Finally the bound for connected graphs is from Lemma 2.2 of . ∎

Fig. 2 depicts examples of some of the graphs discussed in Proposition 5. Clearly the network structure affects the time it takes information to diffuse across the graph. For graphs such as the path or 2-d torus, the dependence on nn is intuitively related to the long time it takes information to spread across the network. For other graphs, such as stars, the dependence is due to the central node (i.e., the “hub”) becoming a bottleneck. For such graphs this dependence is strongly related to the fact that we have focused on the Metropolis scheme for designing the entries of the matrices AkA^{k}. Because the hub has a much higher degree than the other nodes, the resulting Metropolis updates lead to very small changes and hence slower convergence (i.e., AkA^{k} is diagonally dominant); see Eq. (7). In general, for undirected graphs in which neighboring nodes may have very different degrees, it is known that faster rates can be achieved by using linear iterations of the form Eq. (6), where AkA^{k} is optimized for the particular graph topology . However, unlike using the Metropolis weights—which can be implemented locally by having neighboring nodes exchange their degrees—determining the optimal matrices AkA^{k} involves solving a separate network-wide optimization problem; see for a decentralized approximation algorithm.

On the other hand, the algorithm is evidently fast on certain graphs. For the complete graph (where every node is directly connected to every other node, this is not surprising—since A=(1/n)11TA=(1/n){\bf 1}{\bf 1}^{T}, the average is computed exactly at every other node after a single iteration. Expander graphs can be seen as sparse approximations of the complete graph (sparse here is in the sense of having many fewer edges) which approximately preserve the spectrum, and hence the hitting time . In applications where one has the option to design the network, expanders are particularly of practical interest since they allow fast rates of convergence—hence, few iterations—while also having relatively few links—so each iteration requires few transmissions and is thus fast to implement .

II-B Worst-case scaling of decentralized averaging

One might wonder about the worst-case complexity of average consensus: how long does it take to get close to the average on any graph? Initial bounds were exponential in the number of nodes . However, Proposition 5 tells us that this is at most O(n2)O(n^{2}) using a Metropolis matrix. A recent result shows that if all the nodes know an upper bound UU on the total number of nodes which is reasonably accurate, this convergence time can be brought down by an order of magnitude. This is a consequence of the following theorem.

[LinearNote that we do not adhere to the common convention of using “linear convergence” as a synonym for “geometric convergece”; rather, “linear time” convergence in this paper refers to a convergence time which scales as O(n)O(n) in terms of the number of nodes nn. Time Convergence for Consensus] Suppose each node in an undirected connected graph GG implements the update

where w‾=(1/n)∑i=1nwi0\overline{w}=(1/n)\sum_{i=1}^{n}w_{i}^{0} is the initial average.

Thus if every node knows the upper bound UU, the above theorem tells us that the number of iterations until every element of the vector wkw^{k} is at most ϵ\epsilon away from the initial average w‾\overline{w} is O(Uln⁡(1/ϵ))O(U\ln(1/\epsilon)). In the event that UU is within a constant factor of nn, (e.g., n≤U≤10nn\leq U\leq 10n) this turns out to be linear in the number of nodes nn. One situation in which this is possible is if every node precisely knows the number of nodes in the network, in which case they can simply set U=nU=n. However, this scheme is also useful in a number of settings where the exact number of nodes in the system is not known (e.g., if nodes fail) as long as approximate knowledge of the total number of nodes is available.

Intuitively, Eq. (12) takes a lazy Metropolis update and accelerates it by adding an extrapolation term. Strategies of this form are known as over-relaxation in the numerical analysis literature and as Nesterov acceleration in the optimization literature . On a non-technical level, the extrapolation speeds up convergence by reducing the inherent oscillations in the underlying sequence. A key feature, however, is that the degree of extrapolation must be carefully chosen, which is where knowledge of the bound UU is required. At present, open questions are whether any improvement on the quadratic convergence time of Proposition 5 is possible without such an additional assumption, and whether a linear convergence time scaling can be obtained for time-varying graphs.

III Decentralized optimization over undirected graphs

We now shift our attention from decentralized averaging back to the problem of optimization. We begin by describing the (centralized) subgradient method, which is one of the most basic algorithms used in convex optimization.

The subgradient may be viewed as a generalization of the notion of the gradient to non-differentiable (but convex) functions. Indeed, if the function hh is continuously differentiable, then g=∇h(x)g=\nabla h(x) is the only subgradient at xx. In general, there are multiple subgradients at points xx where the function hh is not differentiable. See Figure 3 for a graphical illustration.

The subgradient methodThe earliest work on subgradient methods appears in . for minimizing the function hh is defined as the iterate process

where gkg^{k} is a subgradient of the function hh at the point uku^{k}. The quantity αk\alpha^{k} is a nonnegative step-size.

If the nonnegative step-size sequence αk\alpha^{k} is “summable but not square summable,” i.e.,

Then, the iterate sequence {uk}\{u^{k}\} converges to some minimizer u∗∈U∗u^{*}\in U^{*}.

If the subgradient method is run for TT steps with the (constant) choice of stepsize αk=1/T\alpha^{k}=1/\sqrt{T} for k=0,…,T−1k=0,\ldots,T-1, then

where h∗h^{*} is the minimal value of the function, i.e., h∗=h(u∗)h^{*}=h(u^{*}) for any u∗∈U∗u^{*}\in U^{*}.

(1) A proof can be found in Lemma 7 of . (2) The distance ∥uk−u∗∥22\|u^{k}-u^{*}\|^{2}_{2}, for an arbitrary u∗∈U∗u^{*}\in U^{*} is used to measure the progress of the basic subgradient method. From the definition of the method, for the constant stepsize it follows that

Then, using the subgradient defining inequality in Eq. 13 and the assumption that the subgradient norms are bounded by LL, we obtain

By summing these inequalities over k=0,1,…Tk=0,1,\ldots T, re-arranging the terms, and dividing by 2αT2\alpha T, one can see that

The result follows by using the convexity of h(⋅)h(\cdot) which yields

and by letting α=1T.\alpha=\frac{1}{\sqrt{T}}. ∎

For the diminishing step in part (1), since the iterates uku^{k} converge to some minimizer u∗u^{*}, so does any weighted average of the iterates (with positive weights). In particular, it follows that

Furthermore, it is a fact that any convex function whose domain is the entire space of the decision variables is continuous at every point. Thus, by continuity of h(⋅)h(\cdot), it follows that

In the case of a fixed stepsize, part (2) provides an error bound in terms of the function values. On a conceptual level, the main takeaway is that the subgradient method produces O(1/T)O\left(1/\sqrt{T}\right) convergence to the optimal function value in terms of the number of iterations TT.

Similar to gradients, for the convex functions defined over the entire space, the subgradients are “linear” in the sense that a subgradient of the sum of two convex functions can be obtained as a sum of two subgradients (one for each function). Formally, if g1g_{1} is a subgradient of a function f1f_{1} at xx and g2g_{2} is a subgradient of f2f_{2} at xx, then g1+g2g_{1}+g_{2} is a subgradient of f1+f2f_{1}+f_{2} at xx. (This follows directly from the subgradient definition in Eq. (13).)

III-B Decentralizing the subgradient method

in a decentralized way. If all the functions f1(x),…,fn(x)f_{1}(x),\ldots,f_{n}(x) were available at a single location, we could directly apply the subgradient method to their average f(x)f(x):

where g‾ik\overline{g}_{i}^{k} is a subgradient of the function fi(⋅)f_{i}(\cdot) at uku^{k}. Unfortunately, this is not a decentralized method under our assumptions, since only node ii knows the function fi(⋅)f_{i}(\cdot), and thus only node ii can compute a subgradient of fi(⋅)f_{i}(\cdot).

A decentralized subgradient method solves this problem by interpolating between the subgradient method and an average consensus scheme from Section II. In this scheme, node ii maintains the variable xikx_{i}^{k} which is updated as

Intuitively, the decentralized subgradient method of Eq. (16) pulls the value xikx_{i}^{k} at each node in two directions: on the one hand towards the minimizer (via the subgradient term) and on the other hand towards neighboring nodes (via the averaging term). Eq. (16) can be thought of as reconciling these pulls; note that the strength of the consensus pull does not change, but the strength of the subgradient pull is controlled by the stepsize, and this stepsize αk\alpha^{k} will be later chosen to decay to zero, so that in the limit the consensus term will prevail. However, if the rate at which the stepsize decays to zero is slow enough, then under appropriate conditions consensus will be achieved not on some arbitrary point, but rather on a global minimizer of f(⋅)f(\cdot).

[Convergence and Convergence Time for the Decentralized Subgradient Method] Let X∗{\cal X}^{*} denote the set of minimizers of the function ff. We assume that: (i) each fif_{i} is convex; (ii) X∗{\cal X}^{*} is nonempty; (iii) each function fif_{i} has the property that its subgradients at any point are bounded by the constant LL; (iv) the matrices AkA^{k} are doubly stochastic and there exists some α>0\alpha>0 such that the graph sequence [GA0]α,[GA1]α,[GA2]α,…[G_{A^{0}}]_{\alpha},[G_{A^{1}}]_{\alpha},[G_{A^{2}}]_{\alpha},\ldots satisfies Assumption 1; and (v) the initial values xi0x_{i}^{0} are the same across all nodesThis assumption is not necessary for the results stated here, but we use it to simplify the exposition. When this assumption is violated, the bound in part (ii) has an additional term depending on the spread of the initial values. This term decays exponentially on the order of λk\lambda^{k}. (e.g., xi0=0x_{i}^{0}=0). Then:

If the positive step-size sequence αk\alpha^{k} is “summable but not square summable,” i.e.,

thenIn fact we can show a stronger result that, as k→∞k\to\infty, the iterate sequences {xik}\{x_{i}^{k}\} converge to a common minimizer x∗∈X∗x^{*}\in X^{*}, for all ii. However, the proof is more involved; see . for any x∗∈X∗x^{*}\in{\cal X}^{*}, we have that for all i=1,…,ni=1,\ldots,n,

If we run the protocol for TT steps with (constant) step-size αk=1/T\alpha^{k}=1/\sqrt{T}, and with the notation yk=(1/n)∑i=1nxiky^{k}=(1/n)\sum_{i=1}^{n}x_{i}^{k}, then we have that for all i=1,…,ni=1,\ldots,n,

We remark that the quantity (∑l=0T−1yl)/T(\sum_{l=0}^{T-1}y^{l})/T on which the suboptimality bound is proved can be computed via an average consensus protocol after the protocol is finished if node ii keeps track of (∑l=0T−1xil)/T(\sum_{l=0}^{T-1}x_{i}^{l})/T.

Comparing part (2) of Theorems 7 and 8, and ignoring the similar terms involving the initial conditions, we see that the convergence bound gets multiplied by 1/(1−λ)1/(1-\lambda). This term may be thought of as measuring the “price of decentralization” resulting from having knowledge of the objective function decentralized throughout the network rather than available entirely at one place.

We can use Proposition 5 to translate this into concrete convergence times on various families of graphs, as the next result shows. For ϵ>0\epsilon>0, let us define the ϵ\epsilon-convergence time to be the first time when

Naturally, the convergence time will depend on ϵ\epsilon and on the underlying sequence of matrices/graphs.

Suppose all the hypotheses of Theorem 8 are satisfied, and suppose further that the weights aijka_{ij}^{k} are the lazy Metropolis weights defined in Eq. (8). Then the convergence time can be upper bounded as

…path graphs, then Pn=O(n2)P_{n}=O\left(n^{2}\right);

…22-dimensional grid, then Pn=O(nlog⁡n)P_{n}=O\left(n\log n\right);

…22-dimensional torus, then Pn=O(n)P_{n}=O\left(n\right);

…kk-dimensional torus, then Pn=O(n2/k)P_{n}=O\left(n^{2/k}\right);

…star graphs, then Pn=O(n2)P_{n}=O\left(n^{2}\right);

…two-star graphs, then Pn=O(n2)P_{n}=O\left(n^{2}\right);

…Erdős-Rényi random graphs, then Pn=O(1)P_{n}=O(1);

…geometric random graphs, then Pn=O(nlog⁡n)P_{n}=O(n\log n);

…any connected undirected graph, then Pn=O(n2)P_{n}=O\left(n^{2}\right).

These bounds follow immediately by putting together the upper bounds on 1/(1−λ)1/(1-\lambda) discussed in Proposition 5 with Eq. (18).

We remark that it is possible to decrease the scaling from O(L4Pn2)O(L^{4}P_{n}^{2}) to O(L2Pn)O(L^{2}P_{n}) in the above corollary if the constant LL, the type of the underlying graph (e.g., star graph, path graph), and the number of nodes nn is known to all nodes. Indeed, this can be achieved by setting the stepsize αk=β/T\alpha^{k}=\beta/\sqrt{T} and using a hand-optimized β\beta (which will depend on LL, nn, as well as the kind of underlying graphs). We omit the details but this is very similar to the optimization done in .

We now turn to the proof of Theorem 8. We will need two preliminary lemmas covering some background in optimization. The first lemma discusses how the bound on the norms of the subgradients translate into Lipschitz continuity of the underlying function.

On the one hand, we have by definition of subgradient

so that, by the Cauchy-Schwarz inequality,

Together Eq. (19) and Eq. (20) imply the lemma. ∎

Our overall proof strategy is to view the decentralized subgradient method as a kind of perturbed consensus process. To that end, the next lemma extends our previous analysis of the consensus process to deal with perturbations.

If sup⁡k∥Δk∥2≤L′\sup_{k}\|\Delta^{k}\|_{2}\leq L^{\prime} then

If Δk→0\Delta^{k}\rightarrow 0, then xk−1Txkn1→0x^{k}-\frac{{\bf 1}^{T}x^{k}}{n}{\bf 1}\rightarrow 0.

For convenience, let us introduce the notation

Now using the fact that the vectors e0e^{0} and Δi−mi1\Delta^{i}-m^{i}{\bf 1} have mean zero, by Proposition 4 we have

This equation immediately implies the first claim of the lemma.

Combining this with Eq. (22), we have that ∥ek∥2→0\|e^{k}\|_{2}\rightarrow 0 and this proves the second claim of the lemma. ∎

With these lemmas in place, we now turn to the proof of Theorem 8. Our approach will be to view the decentralized subgradient method as a perturbation of a subgradient-like process followed by averaging of the entries of the vector xkx^{k}. Provided that the step-size αk\alpha^{k} decays to zero at the appropriate rate, we will argue that (i) the vector xkx^{k} is not too far from its average and (ii) this average makes continual progress towards a minimizer of the function f(⋅)f(\cdot).

Since the matrices AkA^{k} are doubly stochastic, 1TAk=1T{\bf 1}^{T}A^{k}={\bf 1}^{T} so that

Now for any x∗∈X∗x^{*}\in{\cal X}^{*}, we have

where the first inequality uses a rearrangement of the definition of the subgradient and the last inequality uses Lemma 10. Plugging this into Eq. (23), we obtain

We now turn to the first claim of the theorem statement. The first term on the right-hand side goes to zero because its numerator is bounded while its denominator is unbounded (due to the assumption that the step-size is summable but not square summable). For the second term, we view −αkgk-\alpha^{k}g^{k} as the perturbation Δk\Delta^{k} in Lemma 6 to obtain that xl−yl1→0x^{l}-y^{l}{\bf 1}\rightarrow 0, and in particular xil−yil→0x_{i}^{l}-y_{i}^{l}\rightarrow 0 for each ii. It follows that the Cesàro sum (which is exactly the second term on the right-hand side) must go to zero as well. We have thus shown that

Putting this together with Lemma 11, which implies that xil−yl→0x_{i}^{l}-y^{l}\rightarrow 0 for all ii, we complete the proof of the first claim.

which completes the proof of the second claim. ∎

III-C Improved scaling with the number of nodes

The results of Corollary 9 improve upon those reported in . A natural question is whether it is possible to further improve the scalings even further. In particular, one might wonder how the worst-case convergence time of decentralized optimization scales with the number of nodes in the network. In general, this question is open. Partial progress was made in , where, under the assumption that all nodes know an order-accurate bound on the total number of nodes in the network, it was shown that we can use the linear time convergence of average consensus described in Theorem 6 to obtain a corresponding convergence time for decentralized optimization when the underlying graph is fixed and undirected. Specifically, consider the following update rule

where gikg_{i}^{{k}} is a subgradient of the function fi(⋅)f_{i}(\cdot) at the point yiky_{i}^{{k}}. As in Section II, here the number UU is an upper bound on the number of nodes known to each individual node, and it is assumed that UU is within a constant factor of the true number of nodes, i.e., n≤U≤cnn\leq U\leq cn for some constant cc (not depending on n,Un,U or any other problem parameters).

By relying on Theorem 6, it is shown in that the corresponding time until this scheme (followed by a round of averaging) is ϵ\epsilon close to consensus on a minimizer of (1/n)∑i=1nfi(⋅)(1/n)\sum_{i=1}^{n}f_{i}(\cdot) is O(nlog⁡n+n/ϵ2)O(n\log n+n/\epsilon^{2}). It is an open question at present whether a similar convergence time can be achieved over time-varying graphs or without knowledge of the upper bound UU.

III-D The strongly convex case

The error decrease of 1/T1/\sqrt{T} with the number of iterations TT is, in general, the best possible rate for dimension-independent convex optimization . Under the stronger assumption that the underlying functions fi(⋅)f_{i}(\cdot) are strongly convex with Lipschitz-continuous gradients, gradient descent will converge geometrically. Until recently, however, there were no corresponding decentralized protocols with a geometric rate.

Significant progress on this issue was first made in , where, over fixed undirected graphs, the following scheme was proposed:

It is not immediately obvious how to extend the EXTRA update to handle time-varying directed graphs; the original proof in only covered static, undirected graphs. Progress on this question was made in which, in addition to providing a geometrically convergent method in the time-varying and directed cases, also provides a new intuitive interpretation of EXTRA. Indeed, observes that the scheme

is a special case of the EXTRA update of Eq. (26). Here, the initialization x0x^{0} can be arbitrary, while y0=∇f(x0)y^{0}=\nabla f(x^{0}). The matrices WkW^{k} are doubly stochastic. Moreover, Eq. (27) has a natural interpretation. In particular, the second line of Eq. (27) is a tracking recursion: yky^{k} tracks the time-varying gradient average 1T∇f(xk)/n{\bf 1}^{T}\nabla f(x^{k})/n. Indeed, observe that, by the double stochasticity of WkW^{k}, we have that

In other words, the vector yky^{k} has the same average as the average gradient. Moreover, it can be seen that if xk→x^x^{k}\rightarrow\widehat{x}, then yk→∇f(x^)y^{k}\rightarrow\nabla f(\widehat{x}); this is due to the “consensus effect” of repeated multiplications by WkW^{k}. Such recursions for tracking were studied in .

While the second line of Eq. (27) tracks the average gradient, the first line of Eq. (27) performs a decentralized gradient step as if yky^{k} was the exact gradient direction. The method can be naturally analyzed using methods for approximate gradient descent. It was shown in that this method converges to the global optimizers geometrically under the same assumptions as EXTRA, even when the graphs are time-varying; further, the complexity of reaching an ϵ\epsilon neighborhood of the optimal solution is polynomial in nn.

We conclude by remarking that there is quite a bit of related work in the literature. Indeed, the idea to use a two-layered scheme as in Eq. (27) originates from . Furthermore, improved analysis of convergence rates over an undirected graph is available in .

IV Averaging and Optimization Over Directed graphs

We have seen in Section II that over time-varying undirected graphs, the lazy Metropolis update results in consensus on the initial average. In this section, we ask whether this is possible over a sequence of directed graphs.

By way of motivation, we remark that many applications of decentralized optimization involve directed graphs. For example, in wireless networks the communication radius of a node is a function of its broadcasting power; if nodes do not all transmit at the same power level, communications will naturally be directed. Any decentralized optimization protocol meant to work in wireless networks must be prepared to deal with unidirectional communications.

Unfortunately, it turns out that there is no direct analogue of the lazy Metropolis method for average consensus over directed graphs. In fact, if we consider deterministic protocols where, at each step, nodes broadcast information to their neighbors and then update their states based on the messages they have received, then it can be proven that no such protocol can result in average consensus; see . The main obstacle is that the consensus iterations we have considered up to now (e.g., in Section II-A) relied on doubly stochastic matrices in their updates, which cannot be done over graphs that are time-varying and directed. We thus need to make an additional assumption to solve the average consensus problem over directed graphs.

A standard assumption in the field is that every node always knows its out-degree. In other words, whenever a node broadcasts a message it knows how many other nodes are within listening range. In practice, this can be accomplished in practice via a two-level scheme, wherein nodes broadcast hello-messages at an identical and high power level, while the remainder of the messages are transmitted at lower power levels. The initial exchange of hello-messages provides estimates of distance to neighboring nodes, allowing each node to see how many listeners it has as a function of its transmission power. Alternatively, the out-degrees can be estimated in a decentralized manner using linear iterations , assuming that the underlying communication topology is strongly connected.

Under this assumption, it turns out that average consensus is indeed possible and may be accomplished via the following iteration,

initialized at an arbitrary x0x^{0} and y0=1y^{0}={\bf 1}. This is known as the Push-Sum iteration; it was introduced in , where its correctness was shown for a fully-connected graph (allowing only pairwise communications), while it was extended to arbitary strongly connected graphs in . In the push-sum was applied to address distributed energy resources over a static directed graph, with a more recent extensions including imperfect communications such as those with delays in and with packet drops .

On an intuitive level, the update of the variables xikx_{i}^{k} does not lead to consensus because of the lack of doubly stochasticity. Instead, at each time kk, each xikx_{i}^{k} is some linear combination of xjkx_{j}^{k} where jj runs over a large enough neighborhood of ii. The main idea of Push-Sum is that an identical iteration started at the all-ones vector (i.e., the update for yiky_{i}^{k}) allows the algorithm to estimate the weights of that linear combination. Once these weights are known, average consensus can be achieved via rescaling. Indeed, we will show later how a decentralized algorithm can use both xikx_{i}^{k} and yiky_{i}^{k} to achieve average consensus.

The name Push-Sum derives from the nature of the decentralized implementation of Eq. (28). Observe that Eq. (28) may be implemented with one-directional communication. Specifically, every node ii transmits (or broadcasts) the values xik/dioutx_{i}^{k}/d^{\text{out}}_{i} and yik/diouty_{i}^{k}/d^{\text{out}}_{i} to its out-neighbors. After these transmissions, each node has the information it needs to perform the update (28), which involves summing the pushed values. In contrast, the algorithms for undirected graphs described in the previous sections required that each node ii send a message to all of its neighbors and receive a message from each neighbor. Protocols of this sort are known as “push-pull” in the decentralized computing literature because the transmission of a message from node ii to node jj (the “push”) implies that ii also expects to receive a message from jj (the “pull”).

Our next theorem, which is the main result of this subsection, tells us that Push-Sum works. For this result, we define matrices AkA^{k} as follows:

[Convergence of Push-Sum] Suppose the sequence of graphs GA0,GA1,GA2,…G_{A^{0}},G_{A^{1}},G_{A^{2}},\ldots satisfies Assumption 1. Then for each i=1,…,ni=1,\ldots,n,

It is somewhat remarkable that the convergence to the average happens for the ratios xik/yikx_{i}^{k}/y_{i}^{k}. Adopting the notation a./ba./b for the element-wise ratio of two vectors aa and bb, the above theorem may be restated as

We now turn to the proof of the theorem. Using the matrices AkA^{k} as defined in Eq. (29), the iterations in Eq. (28) may be written as

Observe that AkA^{k} is column stochastic by design, i.e.,

As a consequence of this, the sums of xkx^{k} and yky^{k} are preserved, i.e.,

For our proof, we will need to use the fact that the vector yky^{k} remains strictly positive and bounded away from zero in each entry. This is shown in the following lemma.

For all i,ki,k, yik≥1/(n2nB)y_{i}^{k}\geq 1/(n^{2nB}).

By Assumption 1, every node has a self-loop, so we have that

and consequently the lemma is true for k=1,…,nB−1k=1,\ldots,nB-1. Let lnBlnB be the largest multiple of nBnB which is at most kk. If k>nB−1k>nB-1 then

By Lemma 2, the matrix AnlB−1:0A^{nlB-1:0} is the transpose product of stochastic matrices satisfying Assumption 1, and consequently each of its entries is at least αnB\alpha^{nB} (where α=1/n\alpha=1/n) by Lemma 2. Thus

Applying Eq. (31) to the last k−nlBk-nlB steps now proves the lemma. ∎

With this lemma in mind, we can give a proof of Theorem 12 that is essentially a quick reduction to the result already obtained in Theorem 1.

Let us introduce the notation zik=xik/yikz_{i}^{k}=x_{i}^{k}/y_{i}^{k}. Then xik=zikyikx_{i}^{k}=z_{i}^{k}y_{i}^{k}, and therefore, we can rewrite the Push-Sum update as

where the last step used the fact that yik≠0y_{i}^{k}\neq 0, which follows from Lemma 13. Therefore, defining

We have thus written the Push-Sum update as an “ordinary” consensus update after a change of coordinates. However, to apply Theorem 1 about the convergence of the basic consensus process, we need to lower bound the entries of PkP^{k}, which we proceed to do next.

Indeed, as a consequence of the definition of PkP^{k}, if we choose α\alpha to be some fixed number such that

always holds, then the sequence of graphs G[P0]α,G[P1]α,G[P2]α,…G_{[P^{0}]_{\alpha}},G_{[P^{1}]_{\alpha}},G_{[P^{2}]_{\alpha}},\ldots will satisfy Assumption 1. To find an α\alpha that satisfies this condition, we make use of the fact that 1/(n2nB)≤yik≤n1/(n^{2nB})\leq y_{i}^{k}\leq n, which is a consequence of Eq. (30) and Lemma 13. It follows that the choice α=(1/n)⋅(1/n)⋅1/(n2nB)\alpha=(1/n)\cdot(1/n)\cdot 1/(n^{2nB}) suffices. Thus, we can apply Theorem 1 and obtain that zkz^{k} converges to a multiple of the all-ones vector.

It remains to show that the final limit point is the initial average. Let z∞z_{\infty} be the ultimate limit of each zikz_{i}^{k}. Then for each k=0,1,2,…k=0,1,2,\ldots,

where the last equality used that each yiky_{i}^{k} is a positive number upper bounded by nn. Finally, appealing to the first relation in Eq. (30) we complete the proof of the theorem. ∎

IV-B Push-Sum based subgradient method

Suppose now that every agent ii has a (scalar) convex objective function fi(⋅)f_{i}(\cdot), and the system objective is to minimize f(x)=(1/n)∑i=1nfi(x)f(x)=(1/n)\sum_{i=1}^{n}f_{i}(x). We next describe a decentralized subgradient method for determining a minimizer of ff using the Push-Sum algorithm. Every node ii maintains scalar variables xik,yik,wikx_{i}^{k},y_{i}^{k},w_{i}^{k}, and updates them according to the following rules: for all k≥0k\geq 0 and all i=1,…,ni=1,\ldots,n,

In the next theorem, we establish the convergence properties of the subgradient method of Eq. (32).

[Convergence of the Push-Sum Subgradient Method] Let X∗{\cal X}^{*} be the set of minimizers of the function ff. Assume that: (i) each fif_{i} is convex; (ii) X∗{\cal X}^{*} is nonempty; (iii) each function fif_{i} has the property that its subgradients at any point are bounded by a constant LL; and (iv) the graph sequence GA0,GA1,GA2,…G_{A^{0}},G_{A^{1}},G_{A^{2}},\ldots satisfies Assumption 1.

If the stepsizes α1,α2,...\alpha^{1},\alpha^{2},... are positive, non-increasing, and satisfy the conditions

then the decentralized subgradient method of Eq. (32) converges asymptotically:

where S0=0S^{0}=0 and Sk=∑s=0k−1αs+1S^{k}=\sum_{s=0}^{k-1}\alpha^{s+1} for k≥1k\geq 1, then for all k≥1k\geq 1, i=1,…,ni=1,\ldots,n, and any x∗∈X∗x^{*}\in X^{*},

The first work to have employed Push-Sum decentralized averaging within a decentralized optimization methods is , and it was further investigated in . This work focused on static graphs, and it has been proposed as an alternative to the algorithm based on synchronous decentralized averaging over undirected graphs in order to avoid deadlocks and synchronization issues, among others. This work also described a decentralized method based on Push-Sum for multi-agent optimization problems with constraints by using Nesterov’s dual-averaging approach. This Push-Sum consensus-based algorithm has been extended to the subgradient-push algorithm in that can deal with convex optimization problems over time-varying directed graphs. More recently, the paper has extended the Push-Sum algorithm to a larger class of decentralized algorithms that are applicable to nonconvex objectives, convex constraint sets, and time-varying graphs.

References combine EXTRA with the Push-Sum approach to produce the DEXTRA (Directed Extra-Push) algorithm for optimization over a directed graph. It has been shown that DEXTRA converges at a geometric (R-linear) rate for a strongly convex objective function, but it requires a careful stepsize selection. It has been noted in that the feasible region of stepsizes which guarantees this convergence rate can be empty in some cases.

V Extensions and Other Work on Decentralized Optimization

We discuss here some extension as well as other algorithms for minimizing the average sum f(⋅)=1n∑i=1nfi(⋅)f(\cdot)=\frac{1}{n}\sum_{i=1}^{n}f_{i}(\cdot) in a decentralized manner.

where ΠX[⋅]\Pi_{X}[\cdot] is the Euclidean projection on the set XX. The subgradient gikg_{i}^{k} of the function fif_{i} can be evaluated at the past iterate xikx_{i}^{k} or at the point ∑j∈Nikaijkxjk\sum_{j\in N_{i}^{k}}a_{ij}^{k}x_{j}^{k}. Since the projection mapping ΠX[⋅]\Pi_{X}[\cdot] is non-expansive (i.e., ∥ΠX[x]−ΠX[y]∥2≤∥x−y∥2\|\Pi_{X}[x]-\Pi_{X}[y]\|_{2}\leq\|x-y\|_{2} for all x,yx,y), the convergence properties of the algorithm with projections remain the same as that of the algorithm without projections.

A more complicated case arises when the constraint set is given as the intersection of per-node constraint sets, i.e., X=∩i=1nXiX=\cap_{i=1}^{n}X_{i}, where each XiX_{i} is a convex closed set and only known to node ii. In this case, the node ii update in Eq. (39) is modified by replacing ΠX[⋅]\Pi_{X}[\cdot] with ΠXi[⋅]\Pi_{X_{i}}[\cdot], thus resulting in the following updates

The projections on the individual agents’ constraint sets XiX_{i} (instead of the true constraint set X=∩i=1nXiX=\cap_{i=1}^{n}X_{i}) introduce additional “perturbations”, which can be controlled with the step-size αk\alpha^{k}, provided that the sets XiX_{i} exhibit some form of regularity. Regularity is a condition requiring that the sum of the distances of a point from the individual sets XiX_{i} is lower bounded by the distance of the point to the intersection of the sets, i.e., ∑i=1n∥x−ΠXi∥22≥c∥x−ΠX∥22\sum_{i=1}^{n}\|x-\Pi_{X_{i}}\|^{2}_{2}\geq c\|x-\Pi_{X}\|^{2}_{2} for all xx. As a result of these additional perturbations coming from the sets XiX_{i}, the convergence analysis of the method is much more involved. This algorithm, including random set-selections, has been studied in for synchronous updates over time-varying graphs and for (random) asynchronous updates over a static graph. A variant of this algorithm (using the Laplacian formulation of the consensus problem) for decentralized optimization with decentralized constraints in noisy networks has been studied in .

V-A2 Effect of noise

We will discuss two possibilities, the case of (stochastic) noisy (sub)gradients and the case of noisy links. In the former case, the decentralized subgradient method proceeds by using stochastic (sub)gradients instead of subgradients. In particular, it assumes the following form:

When the links are noisy, the agent ii may receive xjk+ξijkx_{j}^{k}+\xi^{k}_{ij} instead of the actual quantity xjkx_{j}^{k} that was sent by its neighbor jj, where ξijk\xi^{k}_{ij} is a random link noise. The decentralized algorithm has the following form in this case:

Assuming that the noise process {ξijk}\{\xi_{ij}^{k}\} has zero mean and bounded variance, one can show that all the iterate sequences {xik}\{x_{i}^{k}\} converge to the same minimizer of f(⋅)f(\cdot) almost surely (see for example ). Decentralized inference algorithms for general estimation problems (including nonlinear least squares) in stochastically time-varying networks under noisy gradient computation have been considered in .

V-A3 Random graphs

In the literature of the decentralized methods for multi-agent optimization, the graph sequence G1,G2,…G_{1},G_{2},\ldots is typically assumed to be externally given. The objective is to develop decentralized methods, given a graph sequence that constrains the agent communications. Under such a point of view, the algorithmic design does not address the question of designing the graph sequence, hence, does not optimize the network connectivity structure.

where the neighbor sets NikN_{i}^{k} are random. To ease the representation, the method is re-written as

A special case of such a random iid graph sequence corresponds to the case when the agents use a random gossip or a random broadcast to communicate over a network. These random protocols have traditionally been used in network communication literature as protocols designed for asynchronous information exchange. They have also been used in design of decentralized multi-agent optimization methods, as discussed in the next subsection.

V-A4 Asynchronous vs synchronous computations

All the algorithms we discussed so far have been synchronous in the sense that all nodes update at the same time and also use the same stepsize αk\alpha^{k} at iteration kk. To accommodate the asynchronous updates and, also, allow that agents use different stepsizes, one may resort to a random gossip or broadcast communications, where a random link is activated for communication (gossip) or a random node is activated to broadcast its information to the neighbors. In this case, the decentralized method assumes the form as given in Eq. (40) where the matrix AkA^{k} takes a particular form. Specifically, for the random gossip scheme, the underlying undirected graph GG is static and, at any time kk, only one edge is activated at random, say the edge connecting agents iki_{k} and jkj_{k}. In this case, the matrix AkA^{k} has the following form

where γ∈(0,1)\gamma\in(0,1), and eie_{i} denotes the unit-norm vector with ii entry equal to 1 and all other entries equal to 0. Each matrix AkA^{k} is doubly stochastic, implying that the expected matrix Aˉ\bar{A} is also doubly stochastic.

In the case of a random broadcast, at every iteration kk, each node can be activated with probability 1/n1/n, and the activated node broadcasts its value xikx_{i}^{k} to all of the neighbors j∈Nij\in N_{i}. Given that a node iki_{k} was activated at time kk, the updates follow the rule given in Eq. (40), where

Consensus algorithms implemented in a network using a gossip-based or a broadcast-based communications have been studied in , while a different consensus algorithm (the push-sum method) has been considered in . A nonlinear gossip method is investigated in , while the survey paper provides a detailed account of gossip algorithms for decentralized averaging and their applications to signal processing in sensor networks.

V-B Additional work on decentralized optimization

A decentralized algorithm preserving an optimality condition at every iteration has been proposed in . Decentralized convex optimization algorithms for weight-balanced directed graphs have been investigated in continuous-time .

A different type of a decentralized algorithm for convex optimization has been proposed in , where each agent keeps an estimate for all agents’ decisions. This algorithm solves a problem where the agents have to minimize a global cost function f(x1,…,xm)f(x_{1},\ldots,x_{m}) while each agent ii can control only its variable xix_{i}. The algorithm of has been recently extended to the online optimization setting in . Decentralized algorithms based on the augmented Lagrangian approach with gossip-type communications have been studied in , and accelerated versions of decentralized gradient methods have been proposed and studied in . A consensus-based algorithm for solving problems with a separable constraint structure and the use of primal-dual decentralized methods have been studied in , , while a decentralized primal-dual approach with perturbations have been explored in . Work in provides algorithms for centralized and decentralized convex optimization from the control perspective, while considers an event-triggered decentralized optimization for sensor networks. In , a decentralized simplex algorithm has been developed for linear programming problems, while a Newton-Raphson consensus-based method has been proposed in for decentralized convex problems.

Although our discussion has mainly focused on studying asymptotic rates of convergence of iterative decentralized optimization methods, in it was shown that a related approach based on consensus can solve general constrained abstract optimization problems in a finite number of iterations.

All of the work mentioned above relies on the use of state-independent weights, i.e., the weights that do not depend on the agents’ iterates. A consensus-based algorithm employing state-dependent weights has been proposed and analyzed in .

Another popular decentralized approach for consensus optimization over a static network is the alternating direction method of multipliers (ADMM). This method is based on an equivalent formulation of the consensus constraints. Unlike consensus-based (sub)-gradient method, which operates in the space of the primal-variables, the ADMM solves a corresponding Lagrangian dual problem (obtained by relaxing the equality constraints that are associated with consensus requirement). Just as any dual method, the ADMM is applicable to problems where the structure of the objective functions fif_{i} is simple enough so that the ADMM updates can be executed efficiently. The algorithm has the potential solve the problem with a geometric convergence rate, which requires global knowledge of some parameters including eigenvalues of a weight matrix associated with the graph. A recent survey on the ADMM and its various applications is given in . The first work to address the development of decentralized ADMM over a network is , and it has been investigated in , while its linear rate has been shown in . For an explicit analysis of the relationship to network topology, see . In the ADMM with linearization has been proposed for special composite optimization problems over graphs.

The work in utilizes an adapt-then-combine (ATC) strategy of dynamic weighted-average consensus approach to develop a distribute algorithm, termed Aug-DGM algorithm. This algorithm can be used over static directed or undirected graphs (but requires doubly stochastic matrix). The most interesting aspect of the Aug-DGM algorithm is that it can produce convergent iterates even when different agents use different (constant) stepsizes.

Simultaneously and independently, the idea of tracking the gradient averages through the use of consensus has been proposed in for convex unconstrained problems and in for non-convex problems with convex constraints. The work in develops a large class of decentralized algorithms, referred to as NEXT, which utilizes various “function-surrogate modules” thus providing a great flexibility in its use and rendering a new class of algorithms that subsumes many of the existing decentralized algorithms. The work in and in have also been proposed independently, with the former preceding the latter. The algorithm framework of is applicable to nonconvex problems with convex constraint sets over time-varying graphs, but requires the use of doubly stochastic matrices. This assumption was recently removed in by using column-stochastic matrices, which are more general than the degree-based column-stochastic matrices of the push-sum method. Simultaneously and independently, the papers and have appeared to treat nonconvex problems over graphs. The work in proposes and analyzes a decentralized gradient method based on the push-sum consensus in deterministic and stochastic setting for unconstrained problems.

VI Conclusion and Open Problems

We have discussed decentralized optimization methods for minimizing the average of the nodes’ objectives over graphs. We have considered undirected and directed time varying graphs, and computational models for solving consensus problem in such graphs. Then, we have discussed decentralized optimization algorithms that combine optimization techniques with decentralized averaging algorithms. We have also discussed extensions of the consensus-based approaches and other decentralized optimization algorithms.

In terms of algorithm scalability with the number nn of nodes, at present, it is an open question whether any improvement on the quadratic convergence time of Proposition 5 is possible without an additional assumption about the knowledge of nn (recall that the assumption of knowing a reasonable upper bound on nn was made in Theorem 8 to show a linear scaling with nn in the convergence time). Also, it is not known whether a linear convergence-time scaling can be obtained for time-varying graphs.

Another question for future research is the implementation of decentralized algorithms with lower communication requirements. In particular, even broadcast based communications can be expensive, in terms of the power needed to broadcast in some sensor networks. A question is how to implement decentralized algorithms with fewer communications, and what trade-offs are involved in such implementations. Some initial investigations along these lines were presented in in the context of stochastic optimization, where progressively more time is spent calculating gradients between each round of communication as the number of iterations progresses. There remains much further work to be done along these lines.

Finally, we remark that although there are well-understood lower bounds on the number of iterations required to achieve an ϵ\epsilon-optimal solution in the context of centralized convex optimization , much less is understood about the fundamental limits of decentralized optimization. Although bounds on the number of iterations for centralized algorithms carry over directly to synchronous decentralized algorithms, since any decentralized algorithm can always be emulated on a centralized processor, these results do not provide insight into how much communication is fundamentally required to reach consensus on an ϵ\epsilon-optimal solution. In communication-constrained settings (e.g., where network links have very low bandwidth), it remains an open question as to how many iterations may be required, and a related line of questioning would be to understand when there may be tradeoffs between communication and computation (e.g., to reach an ϵ\epsilon-optimal solution there may be algorithms which require significant computation and lower communication, or vice versa).

Acknowledgements

M.R. thanks Mido Assran for a careful reading and suggestions that improved this paper.

References