Uber redesigns M3DB sharding with subclusters to limit failure impact
Uber has redesigned shard placement in its distributed time series database M3DB by partitioning nodes into fixed-size subclusters, each owning a distinct portion of the shard space. The change limits the blast radius of node failures, maintenance and scaling, which previously could affect up to 66.67% of a cluster.
- Fixed-size subclusters own non-overlapping portions of the shard space
- Under the old model a node failure could affect up to 66.67% of the cluster
- Scaling uses a greedy algorithm with O(S log S) sorting and O(S × N) simulation
- Constraints: equal instance weights and subcluster size a multiple of the replication factor
Read next
Software