ByRDiE: Byzantine-resilient distributed coordinate descent for decentralized learning
Zhixiong Yang, Waheed U. Bajwa
I Introduction
One of the fundamental goals in machine learning is to learn a model that minimizes the statistical risk. This is typically accomplished through stochastic optimization techniques, with the underlying principle referred to as empirical risk minimization (ERM) . The ERM principle involves the use of a training dataset and tools from optimization theory. Traditionally, training data have been assumed available at a centralized location. Many recent applications of machine learning, however, involve the use of a dataset that is either distributed across different locations (e.g., the Internet of Things) or that cannot be processed at a single machine due to its size (e.g., social network data). The ERM framework in this setting of distributed training data is often referred to as decentralized or distributed learning .
While there exist excellent works that solve the problem of distributed learning, all these works make a simplified assumption that all nodes in the network function as expected. Unfortunately, this assumption does not always hold true in practice; examples include cyber attacks, malfunctioning equipments and undetected failures . When a node arbitrarily deviates from its intended behavior, it is termed to have undergone Byzantine failure . While Byzantine failures are hard to detect in general, they can easily jeopardize the operation of the whole network .
In particular, with just a simple strategy, one can show that a single Byzantine node in the network can lead to failures of most state-of-the-art distributed learning algorithms . The main contribution of this paper is to introduce and analyze an algorithm that completes the distributed learning task in the presence of Byzantine failures in the network.
To achieve the goal of distributed learning, one usually sets up a distributed optimization problem by defining and minimizing a (regularized) loss function on the training data of each node. The resulting problem can then be solved by distributed optimization algorithms. Several types of distributed algorithms have been introduced in the past . The most common of these are gradient-based methods , which usually have low local computational complexity. Another type includes augmented Lagrangian-based methods , which iteratively update the primal and dual variables. These methods often require the ability to locally solve an optimization subproblem at each node. A third type includes second-order distributed methods , which tend to have high computational and/or communications costs. While any of these distributed optimization algorithms can be used for distributed learning, they all make the assumption that there are no failures in the network.
Byzantine-resilient algorithms for a variety of problems have been studied extensively over the years . Byzantine-resilient algorithms for scalar- and vector-averaging distributed consensus were studied in . The algorithms proposed in extend some of these works from scalar consensus to scalar-valued distributed optimization, but they cannot be used in a straightforward manner for vector-valued distributed optimization problems. In the context of distributed learning, introduces a method to implement distributed support vector machines (SVMs) under Byzantine failures, but it neither provides theoretical guarantees nor generalizes to other learning problems. A number of recent works have also investigated Byzantine-resilient distributed learning in networks that are equipped with a central processor (often referred to as a paramater server) . Theoretical guarantees developed in such works make extensive use of the fact that the parameter server is connected to all network nodes; as such, these guarantees cannot be generalized to Byzantine-resilient learning in fully distributed networks.
I-B Our Contributions
We have already noted several limitations of works like within the context of Byzantine-resilient learning in fully distributed settings. Additionally, one of the limitations of existing Byzantine-resilient algorithms such as is that, when required to work with vector-valued data, they make a strong assumption on the network topology. Specifically, the smallest size of neighborhood of each node in the vector setting depends linearly on the dimensionality of the problem. This is impractical for most learning problems since the dimensionality of the training samples is usually much larger than the size of the neighborhood of each node. Within the specific context of Byzantine-resilient distributed machine learning, the main limitation of fully distributed algorithms proposed in is that they only yield the minimizer of a convex combination of local empirical risk functions for scalar-valued problems. Since this (scalar) minimizer is usually different from the minimizer of the exact average of local loss functions, these works cannot guarantee by themselves alone that the outputs of any vector-valued algorithms that leverage similar ideas will converge to either the minimum empirical risk or the minimum statistical risk.
In contrast to prior works, our work has two main contributions. First, we propose a Byzantine-resilient algorithm that scales well with the dimensionality of distributed learning problems. The proposed algorithm is an inexact distributed variant of coordinate descent that first breaks vector-valued distributed learning problems into a (possibly infinite) sequence of scalar-valued distributed subproblems and then solves these subproblems in a Byzantine-resilient manner by leveraging ideas from works such as . The inexactness here stems from the fact that—even when using an exact line search procedure—the subproblems cannot be solved exactly in the presence of Byzantine failures (cf. Sec. III). This inexactness in the solution of each subproblem, whether or not an exact line search procedure is used, is one of the reasons the analytical techniques of prior works such as do not lead to theoretical guarantees for the proposed coordinate descent algorithm. In contrast, under the assumption of independent and identically distributed training samples across the network, we provide theoretical guarantees that the output of the proposed algorithm will lead to the minimum statistical risk with high probability. Our theoretical analysis, which forms the second main contribution of this work, also highlights the fact that the output of our algorithm can statistically converge to the minimizer of the statistical risk faster than using only local information by a factor that will be shown in the sequel.
I-C Notation and Organization
The rest of this paper is organized as follows. Section II gives the problem formulation. Section III discusses the proposed algorithm along with theoretical guarantees for consensus as well as statistical and algorithmic convergence. Numerical results corresponding to distributed convex and nonconvex learning on two different datasets are provided in Section IV, while Section V concludes the paper.
II Problem Formulation
Given a connected network in which nodes have access to local training data, our goal is to learn a machine learning model from the distributed data, even in the presence of Byzantine failures. We begin with a model of our basic problem, which will be followed with a model for Byzantine failures and the final problem statement.
Note that while the main problem is being formulated here under the supervised setting, our proposed framework and the final results are equally applicable under both unsupervised and semi-supervised settings.
Note that Assumption 1 implies the risk function itself is also Lipschitz continuous :
Centralized machine learning focuses on learning the “best” mapping from to in terms of the following stochastic optimization problem (referred to as statistical risk minimization):
While the set has no algorithmic significance, it is later shown that the iterates of the proposed algorithm stay within it. Because of this reason, makes an appearance in our theoretical guarantees concerning statistical convergence.
In particular, the minimizer of the empirical risk can be shown to converge to with high probability . In the case of i.i.d. training samples, and for a fixed probability of failure and ignoring some terms, the rate of this statistical convergence is known to be under mild assumptions on the problem .
In many distributed learning problems, training data cannot be made available at a single location. This then requires learning in a distributed fashion, which can be done by employing distributed optimization techniques. The main idea of distributed optimization-based learning is to minimize the average of local empirical risks, i.e.,
To achieve this goal, we need nodes to cooperate with each other by communicating over network edges. Specifically, define the neighborhood of node as . We say that node is a neighbor of node if . Distributed learning algorithms proceed iteratively. In each iteration of the algorithm, node is expected to accomplish two tasks:
Update a local variable according to some (deterministic or stochastic) rule , and
Broadcast the updated local variable to other nodes, where node can receive the broadcasted information from node only if .
While distributed learning via message passing is well understood , existing protocols require all nodes to operate as intended. In contrast, the main assumption in this paper is that some of the nodes in the network can undergo Byzantine failures, formally defined as follows.
A node is said to be Byzantine if during any iteration, it either updates its local variable using an update function or it broadcasts some value other than the intended update to its neighbors.
In this paper, we assume there are Byzantine nodes in the network. Knowing the exact value of , however, is not necessary. One can, for example, set to be an upper bound on the number of Byzantine nodes. Let denote the set of nonfaulty nodes. Without loss of generality, we assume the nonfaulty nodes are labeled from 1 to . We now provide some definitions and an assumption that are common in the literature; see, e.g., .
A subgraph of is called a reduced graph if it is generated by () removing all Byzantine nodes along with their incoming and outgoing edges, and () removing additionally up to incoming edges from each nonfaulty node.
A “source component” of a graph is a collection of nodes such that each node in the source component has a directed path to every other node in the graph.
All reduced graphs generated from contain a source component of cardinality at least .
The purpose of Assumption 3 is to ensure there is enough redundancy in the graph to tolerate Byzantine failures. Note that the total number of different reduced graphs one can generate from is finite as long as is finite. So, in theory, Assumption 3 can be certified for any graph. But efficient certification of this assumption remains an open problem. In the case of Erdős–Rényi graphs used in our experiments, however, we have observed that Assumption 3 is typically satisfied whenever the ratio of the average incoming degree of the graph and the number of Byzantine nodes is high enough.
Assumption 3 has been stated to guarantee resilience against the worst-case attack scenario in which all Byzantine nodes concentrate in the neighborhood of any one of the network nodes. While such worst-case analysis is customary in the literature on Byzantine fault tolerance , it does impose stringent constraints on the network topology. Such topology constraints, however, seem to be unavoidable for the types of “local screening” algorithms being considered in the fully distributed setting of this paper.
II-B Problem Statement
Our goal is to develop a Byzantine fault-tolerant algorithm for distributed learning under Assumptions 1–3. Specifically, under the assumption of at most Byzantine nodes in the network, we need to accomplish the following:
Achieve consensus among nonfaulty nodes in the network, i.e., as the number of algorithmic iterations ; and
Guarantee as the number of training samples at nonfaulty nodes .
In particular, the latter objective requires understanding both the statistical convergence and the algorithmic convergence of the distributed empirical risk minimization problem in the presence of Byzantine failures.
III Byzantine-resilient Distributed Coordinate Descent for Decentralized Learning
In distributed learning, one would ideally like to solve the empirical risk minimization (ERM) problem
at each node and show that as . However, we know from prior work that this is infeasible when . Nonetheless, we establish in the following that distributed strategies based on coordinate descent algorithms can still be used to solve a variant of (4) at nonfaulty nodes and guarantee that the solutions converge to the minimizer of the statistical risk for the case of i.i.d. training data. We refer to our general approach as Byzantine-Resilient Distributed coordinate dEscent (ByRDiE), which is based on key insights gleaned from two separate lines of prior work. First, it is known that certain types of scalar-valued distributed optimization problems can be inexactly solved in the presence of Byzantine failures . Second, coordinate descent algorithms break down vector-valued optimization problems into a sequence of scalar-valued problems . The algorithmic aspects of ByRDiE leverage these insights, while the theoretical aspects leverage tools from Byzantine-resilient distributed consensus , optimization theory , stochastic convex optimization , and statistical learning theory .
ByRDiE involves splitting the ERM problem (4) into one-dimensional subproblems using coordinate descent and then solving each scalar-valued subproblem using the Byzantine-resilient approach described in . The exact implementation is detailed in Algorithm 1. The algorithm can be broken into an outer loop (Step 3) and an inner loop (Step 5). The outer loop is the coordinate descent loop, which breaks the vector-valued optimization problem in each iteration into scalar-valued subproblems. The inner loop solves a scalar-valued optimization problem in each iteration and ensures resilience to Byzantine failures. We assume the total number of iterations for coordinate descent are specified during initialization. We use to denote the -th element of at the -th iteration of the coordinate descent loop and the -th iteration of the -th subproblem (coordinate). Without loss of generality, we initialize .
We now fix some and , and focus on the implementation of the inner loop (Step 5). Every node has some at the start of the inner loop (). During each iteration of this loop, all (nonfaulty) nodes engage in the following: broadcast, screening, and update. In the broadcast step (Step 7), all nodes broadcast ’s and each node receives . During this step, a node can receive values from both nonfaulty and Byzantine neighbors. The main idea of the screening step (Step 8) is to reject values at node that are either “too large” or “too small” so that the values being used for update by node in each iteration will be upper and lower bounded by a set of values generated by nonfaulty nodes. To this end, we partition into 3 subsets , and , which are defined as following:
The step is called screening because node only uses ’s from to update its local variable. Note that there might still be ’s received from Byzantine nodes in . We will see later, however, that this does not effect the workings of the overall algorithm.
The final step of the inner loop in ByRDiE is the update step (Step 9). Using to denote the -th element of , we can write this update step as follows:
where are square-summable (but not summable), diminishing stepsizes: , , and . Notice that is updated after the -th subproblem of coordinate descent in iteration finishes and it stays fixed until the start of the -th subproblem in the -th iteration of coordinate descent. An iteration of the coordinate descent loop is considered complete once all subproblems within the loop are solved. The local variable at each node at the end of this iteration is then denoted by (Step 13). We also express the output of the whole algorithm as . Finally, note that while Algorithm 1 cycles through the coordinates of the optimization variables in each iteration in the natural order, one can use any permutation of in place of this order.
We conclude this discussion by noting that the parameter in ByRDiE, which can take any value between and , trades off consensus among the nonfaulty nodes and the convergence rate; see Sec. III-E for further discussion on this tradeoff and Sec. IV for numerical experiments that highlight this tradeoff. Our theoretical analysis of ByRDiE focuses on the two extreme cases of and . In both cases, we establish in the following that the output of ByRDiE at nonfaulty nodes converges in probability to the minimum of the statistical risk. Convergence guarantees for a finite-valued can be obtained from straightforward modifications of the analytical techniques used in the following.
III-B Theoretical Guarantees: Consensus
We now turn our attention to theoretical guarantees for ByRDiE. We first show that it leads to consensus among nonfaulty nodes in the network, i.e., all nonfaulty nodes agree on the same variable, in the limit of large and/or . Then we show in the next sections that the output of ByRDiE converges to the statistical optimum in the limit of large for the two extreme choices of : and .
where is a matrix that is fully specified in the following.
Let and denote the nonfaulty nodes and the Byzantine nodes, respectively, in the neighborhood of , i.e., and . Notice that one of two cases can happen during each iteration at node :
For case (10(a)), since and , we must have and . Then for each , and satisfying . Therefore, we have for each that such that the following holds:
We can now rewrite the update at node as follows:
It can be seen from (III-B) that the screening rule in ByRDiE effectively enables nonfaulty nodes to replace data received from Byzantine nodes with convex combinations of data received from nonfaulty nodes in their neighborhoods. This enables us to express the updates at nonfaulty nodes in the form (9), with the entries of given by
Notice further that, since (cf. (III-B)), case (10(b)) corresponds to a special case of case (10(a)) in which we keep only the first, second, and last rows of (13). It is worth noting here that our forthcoming proof does not require knowledge of ; in particular, since the choices of and are generally not unique, itself is also generally not unique. The main thing that matters here is that will always be a row stochastic matrix; we refer the reader to for further properties of .
In order to complete our claim that nonfaulty nodes achieve consensus under ByRDiE, even in the presence of Byzantine failures in the network, fix an arbitrary and consider . It then follows from (9) that
, and . Notice that is also row stochastic since it is a product of row-stochastic matrices. We can then express (III-B) as
Next, we need two key properties of from .
Suppose Assumption 3 holds. Then for any , there exists a stochastic vector such that
In words, Lemma 1 states that the product of the row-stochastic matrices converges to a steady-state matrix whose rows are identical and stochastic. We can also characterize the rate of this convergence; to this end, let denote the total number of reduced graphs that can be generated from , and define , , and
Suppose Assumption 3 holds. We then have that ,
Lemma 2 describes the rate at which the rows of converge to . We now leverage this result and show that the nonfaulty nodes achieve consensus under ByRDiE in the limit of large , which translates into and/or . To this end, under the assumption of , we have from (15) the following expression:
Next, suppose the nonfaulty nodes stop computing local gradients at time step and use for . Then, defining , we obtain:
Let Assumptions 1 and 3 hold and fix . Then,
as and/or .
Let Assumptions 1 and 3 hold. Then, fixing , the iterates of ByRDiE satisfy:
where denotes the consensus vector in iteration .
Theorem 2, which is a straightforward consequence of the proof of Theorem 1 (cf. (A)), guarantees a sublinear rate for consensus; indeed, choosing the stepsize evolution to be gives us .
III-C Theoretical Guarantees: Convergence for T→∞→𝑇T\to\infty
We now move to the second (and perhaps the most important) claim of this paper. This involves showing that the output of ByRDiE converges in probability to the minimizer (and minimum) of the statistical risk (cf. (2)) for two extreme cases: Case I: and Case II: . We start our discussion with the case of , in which case an auxiliary lemma simply follows from [13, Theorem 2] (also, see ).
Let Assumptions 1 and 3 hold, and let the -th subproblem of the coordinate descent loop in iteration of ByRDiE be initialized with some . Then, as ,
for some such that .
Lemma 3 shows that each subproblem of the coordinate descent loop in ByRDiE under Case I converges to the minimizer of some convex combination of local empirical risk functions of the nonfaulty nodes with respect to each coordinate. In addition, Lemma 3 guarantees that consensus is achieved among the nonfaulty nodes at the end of each coordinate descent loop under Case I. Note that while this fact is already known from Theorems 1 and 2, Lemma 3 helps characterize the consensus point. In summary, when nonfaulty nodes begin a coordinate descent subproblem with identical local estimates and , they are guaranteed to begin the next subproblem with identical local estimates.
for some such that . Note that is strictly convex and Lipschitz continuous. Now for fixed and , define
In words, if one were to solve the statistical risk minimization problem (2) using (centralized) coordinate descent then will be the -th component of the output of coordinate descent after update of each coordinate in every iteration . In contrast, is the -th component of the outputs of ByRDiE after update of each coordinate in every iteration (cf. Lemma 3). While there exist works that relate the empirical risk minimizers to the statistical risk minimizers (see, e.g., ), such works are not directly applicable here because of the fact that in this paper changes from one pair of indices to the next. Nonetheless, we can provide the following uniform statistical convergence result for ByRDiE under Case I that relates the empirical minimizers to the statistical minimizers .
where , , and .
The proof of this theorem is provided in Appendix B. In words, ignoring minor technicalities that are resolved in Theorem 4 in the following, Theorem 3 states that the coordinate-wise outputs of ByRDiE for all under Case I achieve, with high probability, almost the same statistical risk as that obtained using the corresponding coordinate-wise statistical risk minimizers . We now leverage this result to prove that the iterates of ByRDiE at individual nodes achieve statistical risk that converges to the minimum statistical risk achieved by the statistical risk minimizer (vector) (cf. (2)).
Let Assumptions 1–3 hold. Then, and we have
where , , and for , and are as defined in Theorem 3.
The proof of this theorem is given in Appendix C. We now make a couple of remarks concerning Theorem 4. First, note that the uniqueness of the minimum of strictly convex functions coupled with the statement of Theorem 4 guarantee that with high probability.
Second, Theorem 4 helps crystallize the advantages of distributed learning over local learning, in which nodes individually solve the empirical risk minimization problem using only their local data samples. Prior work on stochastic convex optimization (see, e.g., [47, Theorem 5 and (11)]) tells us that, with high probability and ignoring the terms, the gap between the statistical risk achieved by the empirical risk minimizer and the statistical risk minimizer scales as in the centralized setting. This learning rate translates into for local learning and for the idealized centralized learning. In contrast, Theorem 4 can be interpreted as resulting in the following distributed learning rate (with high probability):We are once again ignoring the terms in our discussion; it can be checked, however, that the terms resulting from Theorem 4 match the ones in prior works such as on centralized learning.
where denotes the effective number of training samples available during distributed learning.
In order to understand the significance of (30), notice that
In particular, if and only if there are no Byzantine failures in the network, resulting in the coordinate descent-based distributed learning of , which matches the centralized learning rate. (This, to the best of our knowledge, is the first result on the explicit learning rate of coordinate descent-based distributed learning.) In the presence of Byzantine nodes, however, the maximum number of trustworthy samples in the network is , and (30) and (31) tell us that the learning rate of ByRDiE in this scenario will be somewhere between the idealized learning rate of and the local learning rate of .
where the parameter is given by
while , and the constants , and are as defined in Theorem 4.
III-D Theoretical Guarantees: Convergence for T=1𝑇1T=1
Case I for ByRDiE, in which , is akin to doing exact line search during minimization of each coordinate, which is one of the classic ways of implementing coordinate descent. Another well-adopted way of performing coordinate descent is to take only one step in the direction of descent in a dimension and then switch to another dimension . Within the context of ByRDiE, this is equivalent to setting (Case II); our goal here is to provide convergence guarantees for ByRDiE in this case. Our analysis in this section uses the compact notation . In the interest of space, and since the probabilistic analysis of this section is similar to that of the previous section, our probability results are stated asymptotically here, rather than in terms of precise bounds.
where denotes the standard basis vector (i.e., it is zero in every dimension except and ) and the relationship between , and is as defined earlier. The sequence effectively helps capture update of the optimization variable after each coordinate-wise update of ByRDiE. In particular, we have the following result concerning the sequence .
Let Assumptions 1, 2, and 3 hold for ByRDiE and choose . We then have that in probability.
The proof of this lemma is provided in Appendix E. We are now ready to state the convergence result for ByRDiE under Case II (i.e, ).
Let Assumptions 1, 2, and 3 hold for ByRDiE and choose . Then, , in probability.
We have from Theorem 1 that for all . The definition of along with Lemma 4 also implies that in probability. This completes the proof of the theorem. ∎
III-E How to Choose the Parameter T𝑇T in ByRDiE?
The parameter in ByRDiE trades off consensus among the nonfaulty nodes and the convergence rate as a function of the number of communication iterations , defined as
In particular, given a fixed , each iteration of ByRDiE involves (scalar-valued) communication exchanges among the neighboring nodes. In the previous section, we have provided theoretical guarantees for the two extreme cases of and . In the limit of large , our results establish that both extremes result in consensus and convergence to the statistical risk minimizer. In practice, however, different choices of result in different behaviors as a function of , as discussed in the following and as illustrated in our numerical experiments in the next section.
When is large, the two time-scale nature of ByRDiE ensures the disagreement between nonfaulty nodes does not become too large in the initial stages of the algorithm; in particular, the larger the number of iterations in the inner loop, the smaller the disagreement among the nonfaulty nodes at the beginning of the algorithm. Nonetheless, this comes at the expense of slower convergence to the desired minimizer as a function of the number of communication iterations .
On the other hand, while choosing also guarantees consensus among nonfaulty nodes, it only does so asymptotically (cf. Theorem 1). Stated differently, ByRDiE cannot guarantee in this case that the disagreement between nonfaulty nodes will be small in the initial stages of the algorithm (cf. Theorem 2). This tradeoff between small consensus error and slower convergence (as a function of communication iterations ) should be considered by a practitioner when deciding the value of . We conclude by noting that the different nature of the two extreme cases requires markedly different proof techniques, which should be of independent interest to researchers.
Our discussion so far has focused on the use of a static parameter within ByRDiE. Nonetheless, it is plausible that one could achieve somewhat better tradeoffs between consensus and convergence through the use of an adaptive parameter in lieu of that starts with a large value and gradually decreases as increases. Careful analysis and investigation of such an adaptive two-time scale variant of ByRDiE, however, is beyond the scope of this paper.
IV Numerical Results
In this section, we validate our theoretical results and make various observations about the performance of ByRDiE using two sets of numerical experiments. The first set of experiments involves learning of a binary classification model from the infinite MNIST datasethttps://leon.bottou.org/projects/infimnist that is distributed across a network of nodes. This set of experiments fully satisfies all the assumptions in the theoretical analysis of ByRDiE. The second set of experiments involves training of a small-scale neural network for classification of the Iris dataset distributed across a network. The learning problem in this case corresponds to a nonconvex one, which means this set of experiments does not satisfy the main assumptions of our theorems. Nonetheless, we show in the following that ByRDiE continues to perform well in such distributed nonconvex learning problems.
We consider a distributed linear binary classification problem involving MNIST handwritten digits dataset. The (infinite) MNIST dataset comprises images of handwritten digits from ‘0’ to ‘9’. Since our goal is demonstration of the usefulness of ByRDiE in the presence of Byzantine failures, we focus only on distributed training of a linear support vector machine (SVM) for classification between digits ‘5’ and ‘8’, which tend to be the two most inseparable digits. In addition to highlighting the robustness of ByRDiE against Byzantine failures in this problem setting, we evaluate its performance for different choices of the parameters , , and .
In terms of the experimental setup, we generate Erdős–Rényi networks () of nodes, of which are randomly chosen to be Byzantine nodes. All nonfaulty nodes are allocated samples—equally divided between the two classes—from the dataset, while a Byzantine node broadcasts random data uniformly distributed between and to its neighbors in each iteration. When running ByRDiE algorithm, each node updates one dimension times before proceeding to the next dimension. All tests are performed on the same test set with 1000 samples of digits ‘5’ and ‘8’ each.
We first report results that confirm the idea that ByRDiE can take advantage of cooperation among different nodes to achieve better performance even when there are Byzantine failures in the network. This involves varying the local sample size and comparing the classification accuracy on the test data. We generate a network of nodes, randomly pick nodes within the network to be Byzantine nodes ( failures), vary from to , and average the final set of results over 10 independent (over network, Byzantine nodes, and data allocation) Monte Carlo trials of ByRDiE. The performance of ByRDiE is compared with two approaches: () coordinate descent-based training using only local data (local CD); and () distributed gradient descent-based training involving network data (DGD). To achieve the best convergence rate for ByRDiE, is chosen to be 1 in these experiments. The final set of results are shown in Fig 1, in which the average classification accuracy is plotted against the number of algorithmic iterations, corresponding to the number of (scalar-valued) communication iterations for ByRDiE, the number of per-dimension updates for local CD, and the number of (vector-valued) communication iterations for DGD. It can be seen that the performance of local CD is not good enough due to the small local sample size. On the other hand, when trying to improve performance by cooperating among different nodes, DGD fails for lack of robustness against Byzantine failures. In contrast, the higher accuracy of ByRDiE shows that ByRDiE can take advantage of the larger distributed dataset while being Byzantine resilient.
Next, we investigate the impact of different values of in ByRDiE on the tradeoff between consensus and convergence rate. This involves generating a network of nodes that includes randomly placed Byzantine nodes within the network ( failures), randomly allocating training samples to each nonfaulty node, and averaging the final set of results over independent trials. The corresponding results are reported in Fig. 2 for four different values of as a function of the number of communication iterations in terms of () average classification accuracy (Fig. 2(a)) and () average pairwise distances between local classifiers (Fig. 2(b)). It can be seen from these figures that leads to the fastest convergence in the initial stages of the algorithm, but this fast convergence comes at the expense of the largest differences among local classifiers, especially in the beginning of the algorithm. In contrast, while results in the slowest convergence, it ensures closeness of the local classifiers at all stages of the algorithm.
Finally, we investigate the impact of different values of (the actual number of Byzantine nodes) on the performance of ByRDiE. In order to amplify this impact, we focus on a smaller network of nodes with a small number of samples per node and . Assuming resilience against the worst-case Byzantine scenario (cf. Remark 3), ByRDiE requires the neighborhood of each node to have at least nodes. Under this constraint, we find that in this setting, which translates into or more of the nodes as being Byzantine, often renders ByRDiE unusable.As noted in Remark 3, one could overcome this limit by opting instead for resilience against the average-case scenario and replacing with . We therefore report our results in terms of both consensus and convergence behavior by varying from to . The final results, averaged over independent trials, are shown in Fig. 3. It can be seen from these figures that both the classification accuracy (Fig. 3(a)) and the consensus performance (Fig. 3(b)) of ByRDiE suffer as increases from to . Such behavior, however, is in line with our theoretical developments. First, as increases, the post-screening graph becomes sparser, which slows down information diffusion and consensus. Second, as increases, fewer samples are incorporated into distributed learning, which limits the final classification accuracy of distributed SVM.
IV-B Distributed Neural Network Using Iris Dataset
While the theoretical guarantees for ByRDiE have been developed for convex learning problems, we now demonstrate the usefulness of ByRDiE in distributed nonconvex learning problems. Our setup in this regard involves distributed training of a small-scale neural network using the three-class Iris dataset. This dataset consists of 50 samples each from three species of irises, with each sample being a four-dimensional feature vector. While one of the species in this data is known to be linearly separable from the other two, the remaining two species cannot be linearly separated. We employ a fully connected three-layer neural network for classification of this dataset, with a four-neuron input layer, a three-neuron hidden layer utilizing the rectified linear unit (ReLU) activation function and a three-neuron output layer using the softmax function for classification. The distributed setup corresponds to a random network of nodes with one Byzantine node (), in which each nonfaulty node is allocated samples that are equally divided between the three classes.
Since distributed training of this neural network requires solving a distributed nonconvex problem, our theoretical guarantees for ByRDiE do not hold in this setting. Still, we use ByRDiE with to train the neural network for a total of 200 independent trials with independent random initializations. Simultaneously, we use DGD for distributed training and also train the neural network in a centralized setting using CD for comparison purposes. The final results are reported in Table I, which lists the average number of outer iterations needed by each training algorithm to achieve classification accuracy. It can be seen from these results that while DGD completely breaks down in the presence of a single Byzantine node, ByRDiE continues to be resilient to Byzantine failures in distributed nonconvex learning problems and comes close to matching the performance of centralized CD in this case.
V Conclusion
In this paper, we have introduced a coordinate descent-based distributed algorithm, termed ByRDiE, that can carry out distributed learning tasks in the presence of Byzantine failures in the network. The proposed algorithm pursues the minimizer of the statistical risk in a fully distributed environment and theoretical results guarantee that ByRDiE achieves this objective under mild assumptions on the risk function and the network topology. In addition, numerical results presented in the paper validate resilience of ByRDiE against Byzantine failures in the network for distributed convex and nonconvex learning tasks. Future works in this direction include strengthening the results in terms of almost sure convergence, deriving explicit rates of convergence, and extensions of theoretical analysis to constant step size and nonconvex risk functions.
Appendix A Proof of Theorem 1
Fix any dimension and recall from (III-B) that
Since and , we get
where we have used the fact that . Thus,
where the last summation follows from the observation that . Further, Assumption 1 implies there exists a coordinate-wise bound such that . Using this and Lemma 2, we obtain
Note that the last inequality in the above expression exploits the fact that . In the case of an arbitrary initialization of ByRDiE, however, we could have bounded the first term in the second inequality in (A) as , which still goes to zero as .
We now expand the term in (A) asWe focus here only on the case of an even for the sake of brevity; the expansion for an odd follows in a similar fashion.
It can be seen from the above expression that . Since , it follows that . This fact in concert with (A) and the diminishing nature of give
Next, take , , and note that and in this case. Further, when and/or . It therefore follows from (41) that as and/or . Finally, since in our analysis was arbitrary, the same holds for all .∎
Appendix B Proof of Theorem 3
Next, notice from (II), (22), and (23) that
Further, since the -dimensional vector is an arbitrary element of the standard simplex, defined as
the probability bound in (44) also holds for any , i.e.,
We now define the set . Our next goal is to leverage (46) and derive a probability bound similar to (44) that uniformly holds for all . To this end, let
Therefore, picking any , and defining and , we have from (B) and (51) that
We now fix , and define and . We then obtain the following from (B)–(57):
for any . Conditioned on this event, we have
Therefore, given any , we have from (B)–(B) that
Appendix C Proof of Theorem 4
We now fix an arbitrary and note that
where () follows from the definition of and () follows from Assumption 1. Plugging in (C), we obtain
Using (67), (C), and the Cauchy–Schwarz inequality yields
Setting in (C) and removing conditioning on (62) using Theorem 3 give us the desired bound.
Appendix D Proof of Theorem 5
with probability as long as (see, e.g., (65) and the discussion around it). Conditioning on the probability event described by (70), the definition of and (70) then give us the following recursion in :
Plugging the inequality in (72) completes the proof.∎
Appendix E Proof of Lemma 4
because of convexity of , the Cauchy–Schwarz inequality, and our assumptions. Since this is the desired result, we need now focus on the claim of strict monotonicity of for this lemma. To prove this claim, note from Assumption 1 that
where (a) follows from the fact that is only nonzero in dimension .
where E(q):=\eta(q)\sum_{i=1}^{|J^{\prime}|}[\pi(r+1)]_{i}\big{(}[\nabla\widehat{f}(Q(q),S_{i})]_{k}-[\nabla\widehat{f}(w_{i}^{r},S_{i})]_{k}\big{)}e_{k}. Plugging this into (E) results in
The right-hand side of (E) strictly lower bounded by implies strict monotonicity of . Simple algebraic manipulations show that this is equivalent to the condition
converges to uniformly for all (equivalently, all ) for any . We therefore have with probability (as )
Next, we consider two cases: () , and () . When , we have in probability . This fact along with (E), the realization that for large enough since , and some tedious but straightforward algebraic manipulations show that the following condition is sufficient for (E) to hold in probability:
Using similar arguments, one can also show that the case also results in (E) as a sufficient condition for (E) to hold in probability. We now note that and in (E) can be upper bounded by some constants and by virtue of Assumption 1 and the definition of . This results in the following sufficient condition for (E):
The right-hand side of (E) can be made arbitrarily small (and, in particular, equal to ) through appropriate choice of and large enough ; indeed, we have from our assumptions, Theorem 1, and the definitions of and that both and converge to as .
Taking summation on both sides of (83) from to , and noting that and , gives us . This contradicts the fact that is lower bounded, thereby validating our claim.