IP Library Patent Application 12511644
Patent Application
App. No. 12/511,644

CONFIGURATION MANAGEMENT IN DISTRIBUTED DATA SYSTEMS

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 None
App. No.
12/511,644
Abstract

Systems and methods for managing configurations of data nodes in a distributed environment A configuration manager is implemented as a set of distributed master nodes that may use quorum-based processing to enable reliable identification of master nodes storing current configuration information, even if some of the master nodes fail. If a quorum of master nodes cannot be achieved or some other event occurs that precludes identification of current configuration information, the configuration manager may be rebuilt by analyzing reports from read/write quorums of nodes associated with a configuration, allowing automatic recovery of data partitions.

Claims (53)

1 . A method of obtaining configuration information defining a current configuration of a plurality of data nodes storing replicas of a partition of a database, the method comprising:

operating at least one processor to perform acts comprising:

receiving a plurality of messages, each message generated by a data node of the plurality of data nodes and indicating a version of the configuration of the database for which the data node is configured and a set of data nodes configured in accordance with the indicated configuration to replicate the partition stored on the data node;

identifying, based on the received messages, a selected set of data nodes, the selected set of data nodes being a set identified in at least one of the plurality of messages for which a quorum of the data nodes in the set each generated a message indicating the same configuration version and the selected set of data nodes; and

storing as a portion of the configuration information an indication that each data node of the selected set is a data node storing a replica of the partition.

2 . The method of claim 1 , wherein the plurality of messages comprise messages from at least half of the data nodes configured to store the partition, and the data nodes forming the quorum comprise at least half of the data nodes storing the partition.

3 . The method of claim 1 , further comprising:

sending a request to the plurality of data nodes storing the database for each to provide a respective message among the plurality of messages.

4 . The method of claim 3 , wherein the storing comprises:

storing the configuration information in a configuration manager, the configuration manager comprising a plurality of master nodes in a master cluster.

5 . The method of claim 4 , further comprising:

in response to detecting an event indicating a loss of integrity of the configuration information stored in the master cluster:

deleting the configuration information from master nodes of the master cluster; and

selecting a master node among the plurality of master nodes as a new primary master node.

6 . The method of claim 1 , wherein a second message among the plurality of messages generated by a second node indicates the second node has a second partition with a first configuration version for said second partition, and identifying data nodes for the second partition, the method further comprising:

inspecting any messages among the plurality of messages from the data nodes for the second partition; and

determining a quorum of data nodes for the second partition does not exit.

7 . The method of claim 1 , further comprising activating the partition in the configuration information.

8 . The method of claim 7 , wherein the identified quorum of data nodes for the partition comprises all of the data nodes for the partition.

9 . A database system storing a database comprising a plurality of partitions, the system comprising:

a plurality of computing nodes; and

a network communicably interconnecting the plurality of computing nodes,

wherein, the plurality of computing nodes comprise:

a plurality of data nodes organized as a plurality of sets, each set comprising nodes of the plurality of data nodes storing a replication of a partition of the plurality of partitions; and

a plurality of master nodes, each master node storing a replication of configuration information, the configuration information identifying the data nodes in each of the plurality of sets and a partition of the plurality of partitions replicated on the nodes of each of the plurality of sets.

10 . The system of claim 9 , wherein

the data nodes for a first partition of the plurality or partitions is each configured to generate a first message identifying the first partition as being replicated on said node, a configuration version of the first partition, and identifying each of the data nodes for the configuration version of the first partition, and

the plurality of master nodes is configured to perform a method in response to a reconfiguration triggering event, the method comprising:

receiving a plurality of the first messages generated by the data nodes for the first partition;

identifying a quorum of data nodes for the first partition, the data nodes forming said quorum each having a same configuration version for said first partition; and

updating the configuration information to indicate, for said first partition, the configuration version of the first partition and the data nodes for said first partition.

11 . The system of claim 10 , wherein the reconfiguration triggering event is a loss of quorum among the plurality of master nodes.

12 . The system of claim 10 , wherein the reconfiguration triggering event is a loss of a primary master node among the plurality of master nodes.

13 . The system of claim 12 , wherein:

each of the plurality of master nodes are assigned a token on a communications ring; and

a new primary master node among a plurality of master nodes remaining after the loss of the primary master node is identified as a master node having a token spanning a predetermined value.

14 . The system of claim 13 , wherein the new primary master node performs the method.

15 . The system of claim 10 , wherein identifying the quorum by the plurality of master nodes comprises:

comparing the configuration version of the first partition identified by the first message from one of the data nodes to the configuration version indicated by the first messages from one or more other data nodes, the one or more other data nodes being the data nodes identified by the first message from the one of the data nodes as being the data nodes for the configuration version of the first partition.

16 . A computer-readable storage medium comprising computer-executable instructions that, when executed by a computer system, perform a method, the method comprising:

identifying the computer system as a primary node for a master partition;

deleting any existing data for a global partition map;

receiving a plurality of messages from at least a subset of a federation of nodes, each message generated by a node among the subset and indicating for said node a partition replicated on said node, a configuration version of the partition, and data nodes for the partition;

identifying a quorum of data nodes for a first partition, the data nodes forming said quorum each having a same configuration version for said first partition; and

updating the global partition map to indicate, for said first partition, the configuration version of the first partition and the data nodes for said first partition.

17 . The computer-readable storage medium of claim 16 , wherein the method further comprises:

analyzing each of the plurality of messages to determine if the configuration version for the partition replicated by the respective node is part of a quorum for said partition.

18 . The computer-readable storage medium of claim 16 , wherein identifying the computer system as the primary node comprises determining the computer system has a token spanning a predetermined value.

19 . The computer-readable storage medium of claim 16 , wherein the method further comprises:

sending a request to a plurality of nodes in the federation for each to provide a respective message among the plurality of messages.

20 . The computer-readable storage medium of claim 16 , wherein identifying the quorum of data nodes comprises:

comparing the configuration version of the first partition identified by the first message from one of the data nodes to the configuration version indicated by the first messages from one or more other data nodes, the one or more other data nodes being the data nodes identified by the first message from the one of the data nodes as being the data nodes for the configuration version of the first partition; and

determining that at least half of the data nodes for the configuration version of the first partition have the same configuration version for said first partition.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 15, 2015
From: MICROSOFT CORPORATION
To: MICROSOFT TECHNOLOGY LICENSING, LLC
Reel/Frame 034766/0509 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 4, 2009
From: VOUTILAINEN, SANTERI OLAVI; KAKIVAYA, GOPALA KRISHNA REDDY; KALHAN, AJAY; XUN, LU; BENVENUTO, MARK C.; SINHA, RISHI RAKESH; SRIKANTH, RADHAKRISHNAN
To: MICROSOFT CORPORATION
Reel/Frame 023049/0178 →