Source-linked AI summary
Distributing the Kalman Filter for Large-Scale Systems
Usman A. Khan, Jose M. F. Moura
TL;DR
Large-scale Kalman filtering is difficult to distribute while preserving coherence with centralized error processes. The paper uses overlapping low-order subsystem filters, covariance assimilation, and distributed fusion to achieve convergence toward the global Kalman filter as information-matrix bands increase.
Problem
Independently evolving local error covariances may lose coherence with the centralized error process, complicating distributed estimation for large-scale systems.
Method
The paper spatially decomposes the system into coupled low-order models and combines local Kalman filters with covariance assimilation and distributed observation fusion.
Results
The distributed implementation converges to the global Kalman filter as the number of bands in the information matrices increases.
Takeaways & Limitations
Communication, computation, and storage are local and distributed, using low-dimensional quantities rather than full state vectors and matrices.
Takeaways & Limitations
The paper notes that existing large-scale models can be infeasible, nonoptimal, and unable to guarantee the desired properties.
Abstract
from arXiv · showhide
This paper derives a \emph{distributed} Kalman filter to estimate a sparsely connected, large-scale, $n-$dimensional, dynamical system monitored by a network of $N$ sensors. Local Kalman filters are implemented on the ($n_l-$dimensional, where $n_l\ll n$) sub-systems that are obtained after spatially decomposing the large-scale system. The resulting sub-systems overlap, which along with an assimilation procedure on the local Kalman filters, preserve an $L$th order Gauss-Markovian structure of the centralized error processes. The information loss due to the $L$th order Gauss-Markovian approximation is controllable as it can be characterized by a divergence that decreases as $L\uparrow$. The order of the approximation, $L$, leads to a lower bound on the dimension of the sub-systems, hence, providing a criterion for sub-system selection. The assimilation procedure is carried out on the local error covariances with a distributed iterate collapse inversion (DICI) algorithm that we introduce. The DICI algorithm computes the (approximated) centralized Riccati and Lyapunov equations iteratively with only local communication and low-order computation. We fuse the observations that are common among the local Kalman filters using bipartite fusion graphs and consensus averaging algorithms. The proposed algorithm achieves full distribution of the Kalman filter that is coherent with the centralized Kalman filter with an $L$th order Gaussian-Markovian structure on the centralized error processes. Nowhere storage, communication, or computation of $n-$dimensional vectors and matrices is needed; only $n_l \ll n$ dimensional vectors and matrices are communicated or used in the computation at the sensors.
I. INTRODUCTION
The introduction motivates a distributed Kalman filter for sparse, coupled large-scale systems, avoiding centralized computation and communication of n-dimensional quantities. It combines spatially overlapping sub-systems, observation fusion, and DICI covariance assimilation to provide a robust and scalable implementation.
- Contribution: The proposed distributed filter uses low-order local Kalman filters and never stores, communicates, or computes n-dimensional quantities.Local quantities have dimensions nl ≪ n, while prior replicated filters require local communication and inversion of n×n matrices, generally an O(n3) operation.
- Spatial Decomposition: Spatial decomposition creates overlapping, coupled sub-systems whose local filters preserve centralized error-covariance structure through covariance assimilation and coupled-state inputs.Overlap preserves the structure of centralized approximated error covariances, while coupled state variables retain system coupling.
- Observation Fusion: Shared observations are fused across overlapping sub-systems using bipartite fusion graphs and local average consensus algorithms.Fusion interactions are constrained to a small neighborhood, potentially using multi-hop communication.
- DICI Algorithm: DICI cooperatively assimilates local error covariances to preserve an Lth-order Gauss-Markovian structure while iteratively approximating centralized information matrices.DICI complexity is independent of n, and its error process is bounded above by that of the distributed Jacobi algorithm.
II. BACKGROUND · A. Global Model
The paper frames large-scale data-assimilation systems as discretized, sparse and localized dynamical models, then specifies their centralized information-filter formulation and sensor observations. It establishes Gaussian noise, observability, coupling, and time-invariance assumptions for the global model.
- II. BACKGROUND: The background concerns time-varying random fields governed by partial differential equations and extends to arbitrary dynamical systems in a particular structural class.
- II. BACKGROUND: The centralized version of the Information filter is presented to fix notation before developing the distributed framework.
- A. Global Model: Because nonlinear models are infeasible for data assimilation, linearized approximations are used for systems such as ocean, heat, and wave models.
- A. Global Model: The motivating discrete models exhibit sparse and localized structure, enabling the spatial distribution developed later.
- A. Global Model: Spatial discretization of elliptical operators on an M × J uniform mesh produces a linear system whose local equations involve neighboring grid values.
- A. Global Model: The system is assumed coupled, irreducible, and globally observable, while N sensors provide localized observations that are stacked into a global observation model.
B. Centralized Information Filter · C. Centralized L−Banded Information filters
The centralized information filter (CIF) performs filtering and prediction using state estimates and inverse error covariances, while sensor observations can be centralized or fused. The centralized L-banded information filter (CLBIF) replaces global information matrices with L-banded approximations, linking computational savings to an Lth-order Gauss-Markovian error-process approximation but remaining centralized.
- B. Centralized Information Filter: The CIF represents filtered and predicted state estimates together with their corresponding error covariances and information-matrix inverses.The notation distinguishes filtered quantities at k|k from prediction quantities at k|k−1.
- B. Centralized Information Filter: The CIF can collect all distributed sensor observations centrally or use observation fusion to construct the global observation variables.The observation-fusion formulation rewrites the global observation variables from the local sensor observations.
- B. Centralized Information Filter: The CIF consists of separate filter and prediction steps applied to the information-form state estimation relations.The filter and prediction steps are specified separately in the centralized formulation.
- C. Centralized L−Banded Information filters: O(n3) computation of global quantities motivates replacing the optimal information matrices with L-banded matrices.The cited example is inversion of a global information matrix.
- C. Centralized L−Banded Information filters: The CLBIF applies L-banded approximations in both filter and prediction steps, with the approximation measured by a divergence between the original and banded information matrices.The divergence is defined using the Frobenius norm and eigenvalues of the approximating matrix.
- C. Centralized L−Banded Information filters: The L-banded information approximation is equivalent to modeling the Gaussian error processes as Lth-order Gauss-Markovian processes.The paper assumes an optimal L-banded approximation in the Kullback-Leibler or maximum-entropy sense.
- C. Centralized L−Banded Information filters: O(n2) complexity reductions are possible for the CLBIF, but the resulting method remains centralized and operates on the n-dimensional state.The paper therefore proceeds by distributing the global model to distribute the CLBIF.
III. SPATIAL DECOMPOSITION OF COMPLEX LARGE-SCALE SYSTEMS … 2) Reduced Models from the Cut-point Sets:
The paper distributes a sparse, localized global model into overlapping reduced-order subsystems at the sensors using system digraphs and cut-point sets. Local dimensions determine an L-banded centralized information approximation, while coupled states and shared observations are exchanged or fused to preserve global dynamics.
- III. SPATIAL DECOMPOSITION OF COMPLEX LARGE-SCALE SYSTEMS: Spatial decomposition exploits the model matrix’s sparse, localized structure to implement local information filters on sensor-specific subsystems instead of the global model.The approach is applicable to general model matrices but is practical when F is sparse and localized.
- A. Reduced Models at Each Sensor: Each sensor’s reduced subsystem is formed by distributing the global dynamical and observation models across local state variables and sensor observations.The illustration uses a five-dimensional system monitored by N = 3 sensors, with the global observation vector stacking local observations.
- 1) Graphical Representation using System Digraphs:: System digraphs represent states and noise inputs as vertices and encode nonzero entries of F and G in the interconnection matrix E.The digraph visualizes dynamical interdependence and supports identifying the structure used for subsystem construction.
- 2) Reduced Models from the Cut-point Sets:: At sensor l, the cut-point set V^(l) contains the locally observed or required states and defines the local state vector and Kalman-filter dimension n_l.If a sensor observes a linear combination of states, the cut-point set includes all states in that combination; unobserved states can be incorporated by extending the sets.
- 2) Reduced Models from the Cut-point Sets:: Choosing L for the desired performance gives a lower bound on the dimensions of the sensor subsystems.The bound follows because additional states may need to be included in the cut-point sets to achieve the selected L.
- 2) Reduced Models from the Cut-point Sets:: Directed edges entering a cut-point set become local internal inputs, so sensors communicate estimates of coupled states to complete local models and preserve global dynamics.In the example, sensor 2 communicates its estimate of x4,k to sensor 1 because x4,k is required there as an internal input but is not in sensor 1’s local state vector.
B. Local Information Filters
Local Information Filters distribute global-state estimation by computing reduced-order local objects at each sensor and exchanging information with neighbors. Their local filtering and prediction steps combine observation, matrix, and estimate fusion without requiring centralized estimation knowledge throughout the network.
- Local Information Filters: Local Information Filters (LIFs) at each sensor distribute estimation of the global state vector using sensor-based reduced models.Each LIF computes local matrices and vectors of dimension n_l.
- Information Fusion: Local matrices and vectors are fused when required by exchanging information among neighboring sensors, with some update procedures performed iteratively.The distributed procedure uses local communication rather than centralized estimation knowledge.
- Distributed Representation: The union of local state vectors represents, in distributed form, the knowledge available centrally in the CLBIF.The passage states that centralized estimation knowledge need not exist at any single location in most applications.
- Filter Structure: Each LIF comprises initial conditions, a local filter step, and a local prediction step.The filter step includes observation fusion and distributed matrix inversion, while prediction includes estimate fusion.
1) Notation: · IV. OVERLAPPING REDUCED MODELS · A. Observation Fusion
The paper defines overlapping local covariance and information matrices for reduced sensor-based models, then fuses observations of shared states through bipartite fusion graphs and local consensus. This fusion supports arbitrarily connected communication networks and, under stationarity, allows observation-matrix fusion to be performed offline.
- 1) Notation:: Local error covariance and information matrices are overlapping diagonal submatrices of corresponding global matrices because reduced sensor-based models share state variables.The information matrices correspond to global L-banded information matrices.
- IV. OVERLAPPING REDUCED MODELS: Reduced models across sensors may share states, so observations associated with shared states must be fused.The shared-state observations are independent across sensors.
- A. Observation Fusion: Because reduced observation variables correspond to different local state vectors, they cannot be added directly as full-state observation variables can.The paper introduces a bipartite fusion graph to coordinate this reduced observation fusion.
- A. Observation Fusion: The bipartite fusion graph connects sensor s_i to state x_j exactly when sensor s_i observes x_j through a nonzero column of its local observation matrix H_i.Its vertex sets are the sensors and state variables.
- A. Observation Fusion: Fusion for each state x_j is performed over the sensors attached to that state, using the induced communication subgraph G_j.States observed by multiple sensors require fusion because they have multiple observations.
- A. Observation Fusion: Iterative weighted averaging computes fused observations over arbitrarily connected networks using only local, potentially multi-hop communication.The same procedure can fuse reduced observation matrices over the corresponding state-variable subgraphs.
- A. Observation Fusion: For stationary observation models with time-independent H and R, reduced observation-matrix fusion is performed once offline; otherwise, it is repeated at each time k.The analysis assumes communication is fast enough for consensus to converge; convergence is geometric and can be accelerated by optimizing weights with SDP.
- A. Observation Fusion: Consensus on observations leads to consensus on estimates, so separate fusion of estimates for shared states is not required.The paper defers the detailed estimate-fusion discussion to the local filtering and prediction steps.
V. DISTRIBUTED MATRIX INVERSION WITH LOCAL COMMUNICATION · A. Centralized Jacobi Overrelaxation (JOR) Algorithm · 1) Convergence:
The paper develops DICI-OR by exploiting an L-banded global information matrix to avoid impractical centralized inversion. Centralized JOR converges under spectral-norm conditions, while its distributed implementation targets local error covariances using only local communication.
- V. DISTRIBUTED MATRIX INVERSION WITH LOCAL COMMUNICATION: The distributed implementation uses local information matrices corresponding to overlapping diagonal submatrices of the global estimation information matrix.Local error covariance matrices are obtained from local filter information matrices and support local prediction.
- V. DISTRIBUTED MATRIX INVERSION WITH LOCAL COMMUNICATION: DICI-OR is introduced to invert the global estimation information matrix while avoiding n×n communication and O(n^3) computation.The method exploits the assumed L-banded structure of the global matrix.
- A. Centralized Jacobi Overrelaxation (JOR) Algorithm: The centralized JOR algorithm iteratively solves matrix equations for an unknown matrix S, with a relaxation parameter γ > 0.Setting T = I_n×n yields S = Z^-1.
- A. Centralized Jacobi Overrelaxation (JOR) Algorithm: Unlike vector-based Jacobi or Gauss-Seidel schemes, DICI uses a nonlinear collapse operator that exploits the inverse structure of symmetric positive definite L-banded matrices.This makes DICI complexity independent of n.
- 1) Convergence:: The JOR error process decays to zero when ||P_γ||_2 < 1, and JOR converges for symmetric positive definite Z with sufficiently small γ > 0.The information matrix Z is symmetric positive definite because it is the inverse of an error covariance matrix, so JOR always converges.
- 1) Convergence:: With γ = 1, JOR becomes the centralized Jacobi matrix algorithm, which converges for diagonally dominant matrices.The error norm can be bounded using the maximum eigenvalue magnitude of P_γ.
- 1) Convergence:: Centralized JOR requires complete n×n matrices, global communication, and nth-order computation at every iteration, motivating its distributed implementation.The distributed method focuses on local error covariances lying on the L-band.
2) Distributed JOR algorithm: · B. Distributed Iterate Collapse Inversion Overrelaxation (DICI-OR) Algorithm · 1) Convergence of the DICI-OR algorithm:
The paper replaces distributed JOR, whose per-sensor computation scales linearly with n, with DICI-OR, an iterate-and-collapse inversion method whose computation is independent of n. Its convergence is formulated through a contraction composition and numerically supported by α values strictly below 1.
- 2) Distributed JOR algorithm:: Distributed JOR requires per-sensor computation that scales linearly with the system dimension n despite using local communication.A single iteration sweeps entire rows in S because off L-band elements recursively require further off L-band elements.
- B. Distributed Iterate Collapse Inversion Overrelaxation (DICI-OR) Algorithm: DICI-OR addresses this limitation with two steps: an iterate step and a collapse step.The collapse step computes non L-band elements from L-band elements, preventing further iteration on non L-band elements.
- B. Distributed Iterate Collapse Inversion Overrelaxation (DICI-OR) Algorithm: The DICI-OR formulation combines local information and communication from neighboring sensors to implement the distributed inversion procedure.Its initial conditions can be computed directly at each sensor from the local information matrix without communication.
- B. Distributed Iterate Collapse Inversion Overrelaxation (DICI-OR) Algorithm: DICI extends to L > 1 using replacement collapse formulae, and its computation requirements are independent of n, providing a scalable matrix-inversion implementation.DICI is obtained from DICI-OR by setting γ = 1.
- 1) Convergence of the DICI-OR algorithm:: The collapse operator converts a symmetric positive definite matrix into one whose inverse is L-banded, while DICI-OR composes it with the linear iterate operator Pγ.The composition is Υ = ζ ◦ Pγ over the set of symmetric positive definite matrices.
- 1) Convergence of the DICI-OR algorithm:: Convergence requires showing that the composition map Υ is a contraction under the spectral norm || · ||2.The iterate operator Pγ is proved contractive, while contraction of the collapse operator ζ is assessed numerically.
- 1) Convergence of the DICI-OR algorithm:: 1.17 × 10^6 simulations produced α values ranging from 0.1938 to 0.9955, with the maximum remaining strictly below 1.The simulations used n = 100 and random L between 1 and 50; the authors therefore numerically verified α ∈ (0, 1).
2) Error Bound for the DICI-OR algorithm: … B. Local Filter Step
The DICI-OR error process is bounded by the JOR error process and therefore converges for symmetric positive definite matrices. The local information filters use fused observations, asymptotically recover the global filter step, and convert local information matrices into covariance and state estimates.
- 2) Error Bound for the DICI-OR algorithm:: The DICI-OR error-process spectral norm is bounded above by the corresponding JOR error-process spectral norm.
- 2) Error Bound for the DICI-OR algorithm:: Because JOR converges for symmetric positive definite matrices, the DICI-OR algorithm also converges.
- 2) Error Bound for the DICI-OR algorithm:: 4490 Monte Carlo trials numerically verify the error bound by comparing DICI-OR and JOR error processes.The simulations use relaxation parameter γ = 0.1 and examine differences in their spectral norms.
- A. Initial Conditions: The local predictor is initialized using a condition derived from the local information representation.Because local information matrices and local error covariances are not inverses, the prediction information matrix uses the L−banded inversion theorem, possibly requiring local communication.
- B. Local Filter Step: The local filter step fuses observations through iterative weighted averaging over sensor communication sub-graphs.Asymptotic convergence requires conditions such as connectivity of the relevant sensor communication sub-graph Gj.
- B. Local Filter Step: Fused local information variables asymptotically converge to the global information variables, making the local filter step converge to the CLBIF global filter step.
- B. Local Filter Step: After filtering, DICI converts local information matrices into local error covariance matrices and then supports conversion to Kalman-filter-domain state estimates.The conversion from information-domain estimates uses a matrix-vector product.
VII. LOCAL INFORMATION FILTERS: LOCAL PREDICTION STEP
This section distributes the global Kalman prediction step into local information filters and predictors using localized, sparse model structure. Local prediction covariances and information matrices rely on neighboring sensors’ communicated internal inputs and selected covariance submatrices, avoiding long-distance communication and full-system matrices.
- Section overview: The global prediction step is distributed across local information filters using results from the DICI algorithm for L-banded matrices.The section formulates both local prediction information matrices and local predictors from the global prediction equations.
- A. Computing the local prediction information matrix, Z(l): Coupled reduced-model dynamics require sensors to communicate estimated states as internal inputs, so each local prediction error depends on neighboring estimation errors.These communicated states enter the local filters, and the dependence is reflected in each local prediction error covariance.
- A. Computing the local prediction information matrix, Z(l): Because F is sparse and localized, each sensor computes its local prediction covariance from selected covariance submatrices available at neighboring sensors, without long-distance communication.Only certain submatrices of the global error covariance are required, and the local information matrix uses additional neighboring sensors with nlth order computation.
- B. Computing the local predictor, bz(l): Under the L = 1-banded tridiagonal assumption, the local predictor uses a limited covariance term, while sparse localized dynamics restrict computation to a small subset of estimated states.The required state estimates may come from sensors modeling those states in reduced models, potentially requiring multi-hop communication.
C. Estimate Fusion … B. Simulations
The distributed local filters fuse shared-state observations through consensus and use DICI to compute local covariance and estimation quantities. Simulations show that increasing consensus or DICI iterations and approximation order L improves agreement with centralized solutions, while larger L increases communication cost.
- C. Estimate Fusion: As consensus iterations increase, local filters reach consensus on estimates of shared states when local predictors also reach consensus.The observation-consensus iterations produce consensus on shared-state estimates under consensus on the local predictors.
- A. Summary of the LIFs: The local information-filter pipeline fuses observations, computes local information matrices and estimators, and applies DICI to obtain covariance and prediction quantities.DICI is used for local error covariance, estimator, and prediction error covariance computations.
- B. Simulations: The simulations use a 100-dimensional system monitored by 10 sensors, with L = 20- and L = 36-banded model matrices and randomly selected observations.The model matrices satisfy ||F||2 = 1, and the global observation matrix is distributed by sensor rows.
- B. Simulations: The experiments average tr(S_k|k) over 1000 Monte Carlo trials for L in [n, 1, 2, 5, 10, 15, 20].DICI and consensus stopping uses a last-10-iterations deviation below 10^-5.
- B. Simulations: The LIFs asymptotically track the CLBIF, and as L increases their performance becomes virtually indistinguishable from the CIF.This behavior is consistent with the approximation’s Kullback-Leibler optimality, while larger L may require communication over a larger neighborhood.
- B. Simulations: The decoupled LIFs at t = 1 are unstable because error covariances are not fused and individual sub-systems are not observable.Increasing Monte Carlo trials reduces variation, and the filters eventually follow the Riccati solution.
- VIII. RESULTS: The results section reports that the distributed filters’ computational advantage over alternatives is discussed after the simulation evidence.The supplied passage introduces, but does not quantify, this computational comparison.
C. Complexity … IX. CONCLUSIONS
The paper contrasts centralized and replicated-model Kalman-filter implementations with a distributed local-information-filter design whose computation and communication remain neighborhood-constrained. Its conclusion combines reduced-order model decomposition, distributed observation fusion, and DICI inversion, with simulations showing convergence toward the global Kalman filter as information-matrix bands increase.
- C. Complexity: Matrix multiplication and inversion are modeled as O(n^3), while multiplying an n × n matrix by an n × 1 vector is O(n^2).These operation costs underpin the subsequent complexity comparisons.
- 1) Centralized Information Filter, CIF:: O(n^3k) is the centralized information filter’s computation complexity, with inordinate back-and-forth communication between sensors and the fusion center.Each sensor sends local observations to a centralized location, where the global observation vector is assembled and processed.
- 2) Information Filters with Replicated Models at each Sensor and Observation Fusion:: The LIF computation is much smaller than the referenced centralized solutions, while communication may be multi-hop but remains neighborhood-constrained by the structure of F.With t_ϖ = max(t_o, t_J1, t_J2), the stated LIF complexity is truncated in the supplied passage but is described as lower-cost.
- IX. CONCLUSIONS: The conclusion distributes communication, computation, and storage across sensors without any sensor processing n-dimensional vectors or matrices.The design distributes the global model into coupled low-order sensor-based models and uses local communication throughout.
- IX. CONCLUSIONS: The method fuses common observations by distributed averaging and inverts full matrices through DICI using only order-n_l matrices and vectors, preserving local-filter coupling.It avoids replicated n-dimensional filters and decoupled reduced models, which are described as infeasible or non-optimal and unable to guarantee centralized error-covariance structure.
- IX. CONCLUSIONS: Simulations show convergence to the global Kalman filter as the number of bands in the information matrices increases.This result concerns the distributed implementation with local Kalman filters at each sensor.
APPENDIX I L−BANDED INVERSION THEOREM
The appendix presents an L-banded inversion theorem for recovering Z = S^-1 from L-band submatrices of S. Each principal submatrix of Z can be computed using only three neighboring L-band submatrices, without constructing the full matrix S.
- Notation: The method partitions matrix addition and subtraction into constituent submatrices and defines principal submatrices by contiguous row and column ranges.S_i^j denotes the principal submatrix spanning rows and columns i through j.
- L-banded inversion theorem: When Z = S^-1 is L-banded, the inversion is expressed using L-banded submatrices of S.The appendix states that the inverse formula is given in (86).
- Local computation: Three neighboring submatrices in the L-band of S suffice to compute a principal submatrix of Z, so the entire S is unnecessary.The algorithm obtains Z from submatrices in the L-band of S.