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

Isometric editorial illustration of stacked modular storage containers in an organized grid, illustrating distributed database partitioning.
Uber has reorganized its M3DB time series database architecture into isolated subclusters to prevent large-scale node failure impact. AI Illustration. Upload story photo >

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

  1. 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?