Post

Uber Redesigns M3DB Sharding with Subclusters to Limit Failure Impact

Uber’s new M3DB placement model bounds the blast radius of node failures and maintenance by partitioning nodes into fixed-size subclusters instead of allowing shard dependencies to span most of a cluster.

Each subcluster owns a non-overlapping shard range while preserving replica isolation across racks or availability zones. Scaling uses a greedy O(S log S) sort plus O(S × N) simulation to choose shard moves that keep the donor balanced, avoiding a second rebalance pass. The trade-offs are real: equal instance weights, scale steps tied to subcluster size and replication factor, and temporary cross-subcluster sharing during expansion.