Uber Redesigned M3DB Shard Placement for Resilience
The database update limits failure impact by partitioning nodes into fixed-size subclusters.
Updated on Sept. 21, 2026 in Data Centers

Live Poll
Do you believe complex infrastructure updates improve system reliability for the average user?
Uber has transitioned its M3DB distributed time series database to a subcluster-based architecture. This update replaces a previous model where node failures could impact up to 66.67% of a cluster.
Why it matters
The change limits the blast radius of node failures and maintenance tasks, which had become difficult to manage as cluster sizes grew. By introducing isolated subclusters, the system prevents complex dependency chains during scaling.
The new architecture organizes nodes into subclusters that own non-overlapping shard spaces, requiring equal instance weights and sizes that are multiples of the replication factor. This replaces a more permissive model that historically allowed failures to cascade across two-thirds of a cluster.
The players
Uber
A global technology company that manages large-scale distributed systems, including M3DB for time series metrics.
The details
Uber implemented this by partitioning nodes into groups that independently manage distinct portions of the shard space. To maintain stability, the system uses a greedy algorithm—a method that makes locally optimal choices at each step—to ensure that shards are moved between nodes so that the remaining load stays balanced. When scaling, destination nodes stream data directly from existing peers to ensure continuity before they officially take ownership of the shards.
Timeline
September 21, 2026: The M3DB redesign was published.
The Tech Race
This transition marks a departure from standard monolithic cluster management in M3DB, reflecting a broader industry push toward strict isolation in distributed databases. It follows a pattern of maturing infrastructure by prioritizing fault containment over the architectural simplicity of smaller clusters.
This change improves operational reliability for the backend systems supporting Uber’s services by isolating the fallout of hardware failures. It provides engineers with a more predictable scaling path, though it requires existing clusters to adhere to equal instance weights.
The takeaway
The move to fixed-size subclusters establishes a new standard for fault tolerance within Uber's time series infrastructure. Monitor future documentation for updates on how this approach scales when clusters exceed current node count constraints.
Further reading
Learn more about the latest infrastructure trends in Data Centers.
Source note: This article includes information reported by InfoQ.
Live Poll
Do you believe complex infrastructure updates improve system reliability for the average user?









