Load balancing in parallel database systems using multi-reordering
Load balancing in a parallel database system is performed using multi-reordering, in which a sequence of multiple processors (two, three, or more) that have small average load (for example, the smallest of any such sequence) is selected to participate in load balancing. In the case of three adjacent low-load processors, the load of all three is evenly distributed over two of them, and the remaining processor becomes a free processor. The free processor is moved adjacent to a high-load processor, the load of which is then shared with the free processor. Data is moved in a “transaction” so that the view of data is always consistent. The database continues to service requests from a former location while data undergoes a transfer.
1. A method of load balancing in a parallel database system comprising a plurality of database nodes, the method comprising:
identifying a database node having an excessive load;
identifying a group of database nodes having light loads relative to the database node having an excessive load;
performing a load adjustment on a relation distributed among the group of database nodes so that at least one database node of the group has zero or low load for that relation; and
performing a load adjustment between the database node having an excessive load and the at least one database node of the group, wherein both the database node having an excessive load and the at least one database node of the group have loads closer to an average load over said group of database nodes for said relation.
2. The method of claim 1 , further comprising using a system skew measure to determine whether or not to make a load adjustment, wherein load balance thrashing is reduced.
3. The method of claim 1 , further comprising:
causing multiple database nodes of the group to have zero or low load for that relation; and
using the multiple database nodes having zero or low load for that relation in a load adjustment operation, wherein load balance thrashing is reduced.
4. The method of claim 1 , wherein performing load adjustment occurs while handling load balancing operations on other relations on an overlapping set of database nodes.
5. The method of claim 1 , wherein performing load adjustment occurs while handling load balancing operations on a non-overlapping set of database nodes.
6. The method of claim 1 ,
wherein performing load adjustment on a relation includes moving tuples of the relation from the at least one database node of the group; and
wherein data movement for performing load adjustment occurs in parallel.
7. The method of claim 1 , further comprising selecting the at least one database node of the group that is logically adjacent to the database node having an excessive load.
8. The method of claim 1 , wherein database nodes of the group of database nodes are logically adjacent to one another.
9. The method of claim 1 , further comprising:
identifying a plurality of candidate groups of database nodes having light loads;
evaluating the plurality of candidate groups of database nodes; and
selecting one of the candidate groups of database nodes as said group of database nodes to participate in load balancing.
10. The method of claim 1 , further comprising:
computing a load imbalance measure based on a difference between actual load and average load;
using the load imbalance measure to determine a database node having an excessive load; and
using the load imbalance measure to determine the group of database nodes having light loads.
11. The method of claim 1 , further comprising:
copying tuples from at least the database node having an excessive load to the at least one database node; and
deleting from the database node having an excessive load the copied tuples.
12. The method of claim 11 , further comprising continuing to service requests for nodes participating in rebalancing during the rebalancing.
13. The method of claim 11 , further comprising continuing to service requests for a relation participating in rebalancing during the rebalancing.
14. A non-transitory computer-readable medium comprising instructions for load balancing a parallel database system comprising a plurality of database nodes, the instructions comprising instructions for:
identifying a database node having an excessive load;
identifying a group of database nodes having light loads relative to the database node having an excessive load;
performing a load adjustment on a relation among the group of database nodes so that at least one database node of the group has zero or low load for that relation; and
performing a load adjustment between the database node having an excessive load and the at least one database node of the group so that both the database node having an excessive load and the at least one database node of the group have loads closer to an average load over said group of database nodes for said relation.
15. The non-transitory computer-readable medium of claim 14 , further comprising instructions for,
during said load adjustment, allowing load balancing operations on other relations on an overlapping set of database nodes.
16. The non-transitory computer-readable medium of claim 14 , further comprising instructions for, during said load adjustment, allowing balancing operations on a non-overlapping set of database nodes.
17. The non-transitory computer-readable medium of claim 14 , further comprising instructions for moving data for performing load adjustment in parallel.
18. The non-transitory computer-readable medium of claim 14 , further comprising instructions for causing the at least one database node of the group to be logically adjacent to the database node having an excessive load.
19. The non-transitory computer-readable medium of claim 14 , wherein database nodes of the group of database nodes are logically adjacent to one another.
20. The non-transitory computer-readable medium of claim 14 , further comprising instructions for:
identifying a plurality of candidate groups of database nodes having light loads;
evaluating the plurality of candidate groups of database nodes; and
selecting one of the candidate groups of database nodes as said group of database nodes to participate in load balancing.
21. The non-transitory computer-readable medium of claim 14 , further comprising instructions for:
computing a load imbalance measure based on a difference between actual load and average load;
using the load imbalance measure to determine a database node having an excessive load; and
using the load imbalance measure to determine the group of database nodes having light loads.
22. The non-transitory computer-readable medium of claim 14 , further comprising instructions for:
copying tuples to be transferred from at least the database node having an excessive load to the at least one database node; and
deleting from the database node having an excessive load the tuples to be transferred.
23. The non-transitory computer-readable medium of claim 22 , further comprising instructions for continuing to service requests for nodes participating in rebalancing during the rebalancing.
24. The non-transitory computer-readable medium of claim 22 , further comprising instructions for continuing to service requests for a relation participating in rebalancing during the rebalancing.
25. A method of load balancing in a parallel database system comprising a plurality of database nodes, the method comprising:
monitoring loads of database nodes
by comparing a load of each database node with an average load of the database nodes for a relation; and
when the load of a database node is excessive for that relation, scheduling a load balancing operation taking into account any database nodes that are expected to be idle during a scheduled time of the load balancing operation.
26. The method of claim 25 , wherein database nodes expected to be idle are selected preferentially to participate in the load balancing operation.
27. The method of claim 25 , further comprising:
identifying a group of database nodes having light loads relative to the database node having an excessive load;
performing a load adjustment among the group of database nodes so that at least one database node of the group has zero or low load; and
performing a load adjustment between the database node having an excessive load and the at least one database node of the group, wherein both the database node having an excessive load and the at least one database node of the group have loads closer to an average load over said group of database nodes.
28. A parallel database system comprising:
a plurality of database nodes, each node of the plurality of database nodes being coupled to a communication medium, and including a processor and a memory coupled to the processor; and
one or more load balance controllers residing in the memory of at least one node and configured to:
monitor loads of database nodes by comparing the load of each database node with an average load of the database nodes for a relation; and
when the load of one of the database node is excessive for that relation, schedule a load balancing operation taking into account the database nodes that are expected to be idle during a scheduled time of the load balancing operation.
29. The database system of claim 28 , wherein the one or more load balance controllers are further configured to preferentially select database nodes expected to be idle to participate in the load balancing operation.
30. The database system of claim 28 , wherein the one or more load balance controllers are further configured to:
identify a group of database nodes having light loads relative to the database node having an excessive load;
cause a load adjustment among the group of database nodes so that at least one database node of the group has zero or low load; and
cause a load adjustment between the database node having an excessive load and the at least one database node of the group so that both the database node having an excessive load and the at least one database node of the group have loads closer to an average load over the group of database nodes for the relation.
31. A parallel database system comprising:
a plurality of database nodes each node coupled to a communication medium, and including a processor and a memory coupled to the processor; and
one or more load balance controllers residing in the memory of at least one node and configured to:
identify a group of database nodes having light loads for a relation relative to a database node having an excessive load for that relation;
cause a load adjustment among the group of database nodes so that at least one database node of the group has zero or low load for said relation; and
cause a load adjustment between the database node having an excessive load and the at least one database node of the group so that both the database node having an excessive load and the at least one database node of the group have loads closer to an average load over said group of database nodes for said relation.
32. The database system of claim 31 , wherein the one or more load balance controllers are configured to:
monitor loads of database nodes
by comparing a load of each database node with an average load of the database nodes; and
when load of a database node is excessive, schedule a load balancing operation taking into account the database nodes that are expected to be idle during a scheduled time of the load balancing operation.
33. A method of load balancing in a parallel database system comprising a plurality of database nodes, the method comprising:
computing a system skew based on an estimated number of tuples stored on each node for every relation in the database;
determining that the system skew exceeds a system skew threshold and a load of a node for a relation is greater than the load of all other nodes for the relation, according to current partition boundaries of the relation;
calculating a number of free nodes needed to reduce the system skew;
selecting a sequence of nodes that supplies at least the calculated number of free nodes;
computing new partition boundaries for the relation based on the selected sequence of nodes;
starting a transaction that transfers tuples of the relation according to the new partition boundaries; and
committing the transaction.
34. A method of balancing a load in a parallel database system having a plurality of database nodes, the method comprising:
detecting a system skew among partition nodes of a relation, the partition nodes being associated with partition boundaries for the relation;
determining that the system skew exceeds a system skew threshold; and
performing a rebalancing database transaction for the relation by altering the partition boundaries boundary to form new partition boundaries and moving tuples of the relation between nodes according to the new partition boundaries to reduce the system skew among the nodes for the relation.
35. A non-transitory computer-readable medium comprising instructions for load balancing a parallel database system comprising a plurality of database nodes, the instructions comprising instructions for:
computing a system skew based on an estimated number of tuples stored on each node for every relation in the database;
determining that the system skew exceeds a system skew threshold and a load of a node for a relation is greater than the load of all other nodes for the relation, according to current partition boundaries of the relation;
calculating a number of free nodes needed to reduce the system skew;
selecting a sequence of nodes that supplies at least the calculated number of free nodes;
computing new partition boundaries for the relation based on the selected sequence of nodes;
starting a transaction that transfers tuples of the relation according to the new partition boundaries; and
committing the transaction.
36. A non-transitory computer-readable medium comprising instructions for load balancing a parallel database system comprising a plurality of database nodes, the instructions comprising instructions for:
detecting a system skew among partition nodes of a relation, the partition nodes being associated with partition boundaries for the relation;
determining that the system skew exceeds a system skew threshold; and
performing a rebalancing database transaction for the relation by altering the partition boundaries boundary to form new partition boundaries and moving tuples of the relation between nodes according to the new partition boundaries to reduce the system skew among the nodes for the relation.