🔍Uber Redesigns M3DB Shard Placement with Fixed-Size Subclusters
Uber's M3DB now limits shard impact with subclusters
TL;DR
Uber has revamped M3DB's shard placement to use fixed-size subclusters, reducing the impact of node failures and maintenance. Each subcluster owns a distinct portion of shard space, improving stability and scalability. Worth watching for anyone managing distributed time series databases.
Uber has overhauled M3DB's shard placement strategy, introducing fixed-size subclusters to minimize the impact of node failures and maintenance. Previously, a node failure could affect up to (n-1)/n of the cluster, but now each subcluster owns a distinct, non-overlapping portion of shard space. This change significantly improves stability and scalability, especially for large clusters. For instance, in a 12-node cluster with a replication factor of three, each subcluster owns half of the shards, ensuring more balanced shard distribution and reduced single-point-of-failure risks. Developers managing distributed time series databases should pay attention to this update, as it could streamline their cluster management and reduce operational overhead.

Key Points
Uber's new M3DB shard placement model partitions nodes into fixed-size subclusters, each owning distinct shard portions.
In a 12-node cluster with a replication factor of three, each subcluster owns half of the shards, ensuring balanced distribution.
Uber uses a greedy algorithm to evaluate shard movement, minimizing load imbalance during scaling operations.
Each subcluster must have an instance weight equal to the replication factor, and scaling occurs in multiples of subcluster size.
M3DB's existing instance-level placement operations remain, preserving compatibility with existing tooling.
Why It Matters
If you're managing a distributed time series database like M3DB, Uber's new shard placement model could significantly reduce the impact of node failures and maintenance. For instance, in a 12-node cluster with a replication factor of three, each subcluster owns half of the shards, ensuring more balanced shard distribution and reduced single-point-of-failure risks. This change improves stability and scalability, making it easier to manage large clusters without triggering large-scale migrations.
Comments
Be the first to comment
Enjoyed this article?
Get it daily. 7am. Free. Reads in 5 minutes.
Join 3,483 builders reading daily.