Systems and methods for balancing storage resources in a distributed database
Abstract
Embodiments are provided for balancing storage resources in a distributed database. According to certain aspects, various hardware components may facilitate a three-stage technique including a node balancer technique, a shard balancer technique, and a replica balancer technique. The node balancer technique may create a set of pods from a set of nodes residing across a set of storage racks. The shard balancer technique may redistribute, among the set of pods, a portion of a set of shards assigned to respective pods of the set of pods. The replica balancer technique may, for each pod, distribute the set of replicas so that the replicas reside across the storage racks of that pod.
Claims
exact text as granted — not AI-modified1 .- 22 . (canceled)
23 . A computer-implemented method of balancing, within a pod, replicas of a set of shards, the pod having assigned a set of nodes residing across a set of storage racks, the set of shards assigned to the set of nodes, the method comprising:
determining, by a computer processor, that the replicas of the set of shards are not balanced among the set of nodes; responsive to the determining, identifying a candidate replica of a given shard of the set of shards for rebalancing the replicas of the set of shards, wherein when there is more than one replica of a given shard within the pod, identifying the candidate replica of the given shard comprises identifying, as the candidate replica, one of the more than one replica of the given shard; identifying a destination node from the pod that does not contain a replica of the given shard; and reassigning the candidate replica to the destination node.
24 . The computer-implemented method of claim 23 , wherein when there is not more than one replica of the given shard within the pod, identifying the candidate replica of the given shard comprises identifying the candidate replica of the given shard that is assigned to a failed node of the set of nodes.
25 . The computer-implemented method of claim 24 , wherein when there is not a failed node of the set of nodes, identifying the candidate replica of the given shard comprises identifying the candidate replica of the given shard that is assigned to a node of the set of nodes having the least amount of free space.
26 . The computer-implemented method of claim 25 , wherein identifying the candidate replica of the given shard further comprises identifying the candidate replica of the given shard that is assigned to the node of the set of nodes (i) having the least amount of free space and (ii) having a highest value for a maximum node shard overlap factor (SOF).
27 . The computer-implemented method of claim 26 , wherein identifying the candidate replica of the given shard further comprises identifying the candidate replica of the given shard that is assigned to the node of the set of nodes whose storage rack has a highest value for a maximum rack shard overlap factor (SOF).
28 . The computer-implemented method of claim 23 , wherein identifying the destination node from the pod that does not contain a replica of the given shard comprises:
identifying a subset of the set of storage racks that do not contain a replica of the given shard; determining, from the subset of the set of storage racks, at least one node having the greatest amount of free space; selecting, from the at least one node having the greatest amount of free space, a subset of nodes that would result in a lowest value for a maximum node shard overlap factor (SOF) between any pair of nodes within the pod; selecting, from the subset of nodes, an additional subset of nodes that would result in a lowest value for a maximum rack SOF between any pair of storage racks associated with the pod; designating, in a deterministic order from the additional subset of nodes, the destination node.
29 . The computer-implemented method of claim 28 , wherein selecting the subset of nodes that would result in the lowest value for the maximum node shard overlap factor (SOF) between any pair of nodes within the pod comprises identifying, from the at least one node having the greatest amount of free space, a node currently having the lowest value for the maximum node SOF within the pod.
30 . The computer-implemented method of claim 28 , further comprising repeating the identifying the candidate replica, the identifying the subset of the set of storage racks, the determining the at least one node, the selecting the subset of the nodes, the selecting the additional subset of nodes, the designating the destination node, and the reassigning, until each replica of each of the set of shards is assigned to a node of a different storage rack.
31 . The computer-implemented method of claim 28 , further comprising repeating, for additional replicas of an additional set of shards within at least one additional pod, the identifying the candidate replica, the identifying the subset of the set of storage racks, the determining the at least one node, the selecting the subset of the nodes, the selecting the additional subset of nodes, the designating the destination node, and the reassigning.
32 . The computer-implemented method of claim 23 , further comprising distributing at least a portion of the set of shards to the pod from an additional pod according to an available capacity of each of the pod and the additional pod.
33 . The computer-implemented method of claim 23 , wherein at least two of the set of nodes assigned to the pod reside on different storage racks from the set of storage racks.
34 . A system for balancing, within a pod, replicas of a set of shards, the system comprising:
a computer processor; a set of storage racks; a set of nodes (i) residing across the set of storage racks, (ii) assigned to the pod, and (iii) having the set of shards assigned thereto; and a balancer module executed by the computer processor and configured to:
determine that the replicas of the set of shards are not balanced among the set of nodes;
responsive to the determination, identify a candidate replica of a given shard of the set of shards for rebalancing the replicas of the set of shards, wherein when there is more than one replica of a given shard within the pod, identifying the candidate replica of the given shard comprises identifying, as the candidate replica, one of the more than one replica of the given shard;
identify a destination node from the pod that does not contain a replica of the given shard; and
reassign the candidate replica to the destination node.
35 . The system of claim 34 , wherein when there is not more than one replica of the given shard within the pod, the balancer module is configured to identify the candidate replica of the given shard that is assigned to a failed node of the set of nodes.
36 . The system of claim 35 , wherein when there is not a failed node of the set of nodes, the balancer module is configured to identify the candidate replica of the given shard that is assigned to a node of the set of nodes having the least amount of free space.
37 . The system of claim 36 , wherein to identify the candidate replica of the given shard, the balancer module is configured to further identify the candidate replica of the given shard that is assigned to the node of the set of nodes (i) having the least amount of free space, and (ii) having a highest value for a maximum node shard overlap factor (SOF).
38 . The system of claim 37 , wherein the balancer module further identifies the candidate replica of the given shard that is assigned to the node of the set of nodes whose storage rack has a highest value for a maximum rack shard overlap factor (SOF).
39 . The system of claim 34 , wherein to identify the destination node from the pod that does not contain a replica of the given shard, the balancer module is configured to:
identify a subset of the set of storage racks that do not contain a replica of the given shard; determine, from the subset of the set of storage racks, at least one node having the greatest amount of free space; select, from the at least one node having the greatest amount of free space, a subset of nodes that would result in a lowest value for a maximum node shard overlap factor (SOF) between any pair of nodes within the pod; select, from the subset of nodes, an additional subset of nodes that would result in a lowest value for a maximum rack SOF between any pair of storage racks associated with the pod; designate, in a deterministic order from the additional subset of nodes, the destination node.
40 . The system of claim 39 , wherein to select the subset of nodes that would result in the lowest value for the maximum node shard overlap factor (SOF) between any pair of nodes within the pod, the balancer module is further configured to identify, from the at least one node having the greatest amount of free space, a node currently having the lowest value for the maximum node SOF within the pod.
41 . The system of claim 39 , wherein the balancer module is further configured to repeat the identifying the candidate replica, the identifying the subset of the set of storage racks, the determining the at least one node, the selecting the subset of the nodes, the selecting the additional subset of nodes, the designating the destination node, and the reassigning, until each replica of each of the set of shards is assigned to a node of a different storage rack.
42 . The system of claim 39 , wherein the balancer module is further configured to repeat, for additional replicas of an additional set of shards within at least one additional pod, the identifying the candidate replica, the identifying the subset of the set of storage racks, the determining the at least one node, the selecting the subset of the nodes, the selecting the additional subset of nodes, the designating the destination node, and the reassigning.
43 . The system of claim 34 , wherein the balancer module is further configured to distribute at least a portion of the set of shards to the pod from an additional pod according to an available capacity of each of the pod and the additional pod.
44 . The system of claim 34 , wherein at least two of the set of nodes assigned to the pod reside on different storage racks from the set of storage racks.Join the waitlist — get patent alerts
Track US2019334991A1 — get alerts on status changes and closely related new filings.
We store only your email — no account needed. See our privacy policy.