IP Library Granted Patent US 11,150,958
Granted Patent B2
US 11,150,958 · App. 16/548,773 · Granted Oct 19, 2021

Quorum based transactionally consistent membership management in distributed storage

Inventors: Santeri Olavi Voutilainen (Seattle, WA); Gopala Krishna Reddy Kakivaya (Sammamish, WA); Ajay Kalhan (Redmond, WA); Lu Xun (Kirkland, WA)
Assignee: Microsoft Technology Licensing, LLC
G06F9/5061G06F16/27
View Patent ↗
Loading inventors, assignments & file history…
Monitor This Case
Get email alerts when status or documents change.
Order Certified Copies
Most orders are placed with the USPTO same day — all within 24 business hours.
Order via The Patent Place →
Pre-filled with this patent's details
Quick Facts
Patent No.
US 11,150,958
App. No.
16/548,773
Granted
Oct 19, 2021
Kind
B2
Abstract

Systems and methods that restore a failed reconfiguration of nodes in a distributed environment. By analyzing reports from read/write quorums of nodes associated with a configuration, automatic recovery for data partitions can be facilitated. Moreover, a configuration manager component tracks current configurations for replication units and determines whether a reconfiguration is required (e.g., due to node failures, node recovery, replica additions/deletions, replica moves, or replica role changes, and the like.) Reconfigurations of data activated as being replicated from an old configuration to being replicated on a new configuration may be performed in a transactionally consistent manner based on dynamic quorums associated with the new configuration and the old configuration.

Claims (97)

1. A computing device, comprising:

a memory and a processor that are respectively configured to store and execute instructions, including instructions for causing the computing device to perform operations for reconfiguring at least a portion of a distributed system including a first set of nodes in a first configuration to a second set of nodes in a second configuration, the at least the portion of the distributed system storing a plurality of transactions, the operations including:

determining an identifier associated with the second configuration;

updating transactions stored on the second configuration with transactions stored on the first configuration;

transitioning from the first configuration to the second configuration;

wherein the proceeding operations are performed in response to acceptances from at least a quorum number of nodes on which the operations are proposed to be performed; and

receiving a response to the reconfiguration proposal from each node in at least a subset of the first configuration and from each node in at least a subset of the second configuration, wherein each response received from a node accepting the reconfiguration proposal contains an indication of the acceptance of the reconfiguration proposal, and wherein each response received from a node rejecting the reconfiguration proposal includes an indication of the rejection of the reconfiguration proposal.

2. The computing device of claim 1 , wherein:

determining the identifier associated with the second configuration comprises:

selecting the identifier associated with the second configuration;

sending to each of a plurality of nodes in the first configuration and to each of a plurality of nodes in the second configuration a reconfiguration proposal comprising the selected identifier associated with the second configuration;

receiving a response to the reconfiguration proposal from each node in at least a subset of the first configuration and from each node in at least a subset of the second configuration, wherein each response received from a node contains an indication of an acceptance or a rejection of the reconfiguration proposal;

determining whether the received responses to the reconfiguration proposal indicate acceptance of the reconfiguration proposal by at least the quorum number of nodes; and

in response to a determination that the received responses to the reconfiguration proposal indicate acceptance of the reconfiguration proposal by at least the quorum number of nodes, determining that the selected identifier is to be associated with the second configuration.

3. The computing device of claim 1 , wherein the operations further comprise:

in response to a determination that received responses to a reconfiguration proposal indicate rejection of the reconfiguration proposal, restarting the reconfiguring.

4. The computing device of claim 1 , wherein:

updating transactions stored on the second configuration with transactions stored on the first configuration comprises:

instructing each node in the second configuration to update a replica of the at least a portion of the distributed system stored on the node, wherein the updated replica contains each transaction that has been stored as being committed on at least a subset of the first configuration, the size of the subset of the first configuration representing a first quorum value associated with the first configuration; and

determining whether a number of nodes in the second configuration that have updated the replica is at least a second quorum value associated with the second configuration.

5. The computing device of claim 1 , wherein:

transitioning from the first configuration to the second configuration comprises:

sending a deactivation message to each of a plurality of nodes of the first configuration;

receiving a response to the deactivation message from each node in at least a subset of the plurality of nodes of the first configuration, wherein the responses to the deactivation message contain an indication of an acceptance or a rejection of the deactivation;

determining whether the received responses to the deactivation message indicate acceptance of the deactivation by at least the quorum number of nodes; and

in response to a determination that the received responses to the deactivation message indicate acceptance of the deactivation by at least the quorum number of nodes, deactivating the first configuration.

6. The computing device of claim 5 , wherein the operations further comprise:

in response to a determination that a request to deactivate the first configuration was rejected, restarting the reconfiguring.

7. The computing device of claim 1 , wherein:

transitioning from the first configuration to the second configuration comprises:

sending an activation message to each of a plurality of nodes of the second configuration;

receiving a response to the activation message from each node in at least a subset of the plurality of nodes of the second configuration, wherein the responses to the activation messages contain an indication of an acceptance or a rejection of the activation;

determining whether the received responses to the activation message indicate acceptance of the activation by at least the quorum number of nodes; and

in response to a determination that the received responses to the activation message indicate acceptance of the activation by at least the quorum number of nodes, activating the second configuration.

8. A computer-readable storage medium, comprising at least one of a memory, disk, disc, or device, that is encoded with computer-executable instructions that, in response to execution, cause operations for replicating at least a portion of a distributed system including a first set of nodes in a first configuration to a second set of nodes in a second configuration to be performed, the operations including:

determining an identifier associated with the second configuration;

updating transactions stored on the second configuration with transactions stored on the first configuration;

transitioning from the first configuration to the second configuration;

wherein the operations are performed in response to indications of acceptance of the operations from at least a quorum number of nodes on which the operations are proposed to be performed; and

receiving a response to the reconfiguration from each node in at least a subset of the first configuration and from each node in at least a subset of the second configuration, wherein each response received from a node rejecting the reconfiguration includes an indication of the rejection of the reconfiguration.

9. The computer-readable storage medium of claim 8 , wherein:

determining the identifier associated with the second configuration comprises:

selecting the identifier associated with the second configuration;

sending to each of a plurality of nodes in the first configuration and to each of a plurality of nodes in the second configuration a reconfiguration proposal comprising the selected identifier associated with the second configuration;

receiving a response to the reconfiguration proposal from each node in at least a subset of the first configuration and from each node in at least a subset of the second configuration, wherein each response received from a node contains an indication of an acceptance or a rejection of the reconfiguration proposal;

determining whether the received responses to the reconfiguration proposal indicate acceptance of the reconfiguration proposal by at least the quorum number of nodes; and

in response to a determination that the received responses to the reconfiguration proposal indicate acceptance of the reconfiguration proposal by at least the quorum number of nodes, determining that the selected identifier is to be associated with the second configuration.

10. The computer-readable storage medium of claim 8 , wherein the operations further comprise:

in response to a determination that received responses to a reconfiguration proposal indicate rejection of the reconfiguration proposal, restarting the operations.

11. The computer-readable storage medium of claim 8 , wherein updating transactions stored on the second configuration with transactions stored on the first configuration comprises:

instructing each node in the second configuration to update a replica of the at least a portion of the distributed system stored on the node, wherein the updated replica contains each transaction that has been stored as committed on at least a subset of the first configuration, the size of the subset of the first configuration representing a first quorum value associated with the first configuration; and

determining whether a number of nodes in the second configuration that have updated the replica is at least a second quorum value associated with the second configuration.

12. The computer-readable storage medium of claim 8 , wherein transitioning from the first configuration to the second configuration comprises:

sending a deactivation message to each of a plurality of nodes of the first configuration;

receiving a response to the deactivation message from each node in at least a subset of the plurality of nodes of the first configuration, wherein the responses to the deactivation message contain an indication of an acceptance or a rejection of the deactivation;

determining whether the received responses to the deactivation message indicate acceptance of the deactivation by at least the quorum number of nodes; and

in response to a determination that the received responses to the deactivation message indicate acceptance of the deactivation by at least the quorum number of nodes, deactivating the first configuration.

13. The computer-readable storage medium of claim 8 , wherein:

transitioning from the first configuration to the second configuration comprises:

sending an activation message to each of a plurality of nodes of the second configuration;

receiving a response to the activation message from each node in at least a subset of the plurality of nodes of the second configuration, wherein the responses to the activation messages contain an indication of an acceptance or a rejection of the activation;

determining whether the received responses to the activation message indicate acceptance of the activation by at least the quorum number of nodes; and

in response to a determination that the received responses to the activation message indicate acceptance of the activation by at least the quorum number of nodes, activating the second configuration.

14. A method of reconfiguring at least a portion of a distributed system including a first set of nodes in a first configuration to a second set of nodes in a second configuration, the method comprising operations including:

determining an identifier associated with the second configuration;

updating transactions stored on the second configuration with transactions stored on the first configuration;

transitioning from the first configuration to the second configuration;

committing the second configuration,

wherein the operations are performed in response to acceptance of the operations by at least a quorum number of nodes on which the operations are to be performed; and

receiving a response to the operations from each node in at least a subset of the first configuration and from each node in at least a subset of the second configuration, wherein each response received from a node rejecting the operations includes an indication of the rejection of the operations.

15. The method of claim 14 , wherein:

determining the identifier associated with the second configuration comprises:

selecting the identifier associated with the second configuration;

sending to each of a plurality of nodes in the first configuration and to each of a plurality of nodes in the second configuration a reconfiguration proposal comprising the selected identifier associated with the second configuration;

receiving a response to the reconfiguration proposal from each node in at least a subset of the first configuration and from each node in at least a subset of the second configuration, wherein each response received from a node contains an indication of an acceptance or a rejection of the reconfiguration proposal;

determining whether the received responses to the reconfiguration proposal indicate acceptance of the reconfiguration proposal by at least the quorum number of nodes; and

in response to a determination that the received responses to the reconfiguration proposal indicate acceptance of the reconfiguration proposal by at least the quorum number of nodes, determining that the selected identifier is to be associated with the second configuration.

16. The method of claim 14 , wherein the operations further comprise:

in response to a determination that received responses to a reconfiguration proposal indicate rejection of the reconfiguration proposal, restarting the operations.

17. The method of claim 14 , wherein:

transitioning from the first configuration to the second configuration comprises:

sending a deactivation request to each of a plurality of nodes of the first configuration;

receiving a response to the deactivation request from each node in at least a subset of the plurality of nodes of the first configuration, wherein the responses to the deactivation request contain an indication of an acceptance or a rejection of the deactivation request;

determining whether the received responses to the deactivation request indicate acceptance of the deactivation request by at least the quorum number of nodes; and

in response to a determination that the received responses to the deactivation request indicate acceptance of the deactivation request by at least the quorum number of nodes, deactivating the first configuration.

18. The method of claim 14 , wherein the operations further comprise:

in response to a determination that deactivation request was rejected, restarting the operations.

19. The method of claim 14 , wherein:

transitioning from the first configuration to the second configuration comprises:

sending an activation request to each of a plurality of nodes of the second configuration;

receiving a response to the activation request from each node in at least a subset of the plurality of nodes of the second configuration, wherein the responses to the activation request contain an indication of an acceptance or a rejection of the activation request;

determining whether the received responses to the activation request indicate acceptance of the activation request by at least the quorum number of nodes; and

in response to a determination that the received responses to the activation request indicate acceptance of the activation request by at least the quorum number of nodes, activating the second configuration.

20. The method of claim 14 , wherein:

updating transactions stored on the second configuration with transactions stored on the first configuration comprises:

instructing each node in the second configuration to update a replica of the at least a portion of the distributed system stored on the node, wherein the updated replica contains each transaction that has been stored as being committed on at least a subset of the first configuration, the size of the subset of the first configuration representing a first quorum value associated with the first configuration; and

determining whether a number of nodes in the second configuration that have updated the replica is at least a second quorum value associated with the second configuration.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 1, 2019
From: VOUTILAINEN, SANTERI OLAVI; KAKIVAYA, GOPALA KRISHNA REDDY; KALHAN, AJAY; XUN, LU
To: MICROSOFT CORPORATION
Reel/Frame 050895/0005 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 1, 2019
From: MICROSOFT CORPORATION
To: MICROSOFT TECHNOLOGY LICENSING, LLC
Reel/Frame 050895/0012 →
Continuity (5)
Continuation 15401012 · Jan 7, 2017
Continuation 13861448 · Apr 12, 2013
Continuation 12511525 · Jul 29, 2009
Provisional Application 61107938 · Oct 23, 2008
Related Publication 20200050495A1 · Feb 13, 2020