Determining cluster membership in a distributed computer system
Abstract
(57) [Summary] Cluster membership in a distributed computer system is determined by determining which other nodes each node is communicating with and by distributing its connectivity information through the nodes of the system. Therefore, each node can determine a new optimized cluster based on connectivity information. In particular, each node has information about the node with which it is communicating and has similar information about each other node in the system. Therefore, each node has complete information about the connectivity of all nodes that are directly or indirectly connected. Each node applies optimization criteria to its connectivity information to determine the best new cluster. Data representing the optimal new cluster is broadcast by each node. In addition, the optimal new cluster determined by many nodes is collected at each node. Each node will have data representing the proposed new cluster that was found to be optimal by each node. Each node uses that information to elect a new cluster from many proposed new clusters. For example, a new cluster that is proposed more than any other cluster is elected as the new cluster. Each node receives the same proposed new cluster from a node that is a potential member of the new cluster, so new cluster membership reaches unanimous. Since each node has complete more information about the nodes of potential members of the new cluster, the resulting cluster is undoubtedly in a relatively optimal configuration.
Term
Term ended
Projected expiry passed 20 October 2018, 7.9 years ago.
- Priority
- Filed
- Published
- Projected expiry
- Today
1 claim: 1 independent, 0 dependent
- 1【特許請求の範囲】 【請求項1】 分散型コンピュータ・システムにおけるノードのメンバーシップを決定する方法において、 (a) 分散型コンピュータ・システムのノードの接続性を表す接続データを決定するステップと、 (b) プロポーズされた新しいクラスタのプロポーズされたメンバーシップ・リストを形成するために接続データに最適化基準を適用するステップと、 (c) プロポーズされたメンバーシップ・リストを接続されているノードに同報通信するステップと、 (d) 接続されているノードから他のプロポーズされたメンバーシップ・リストを受け取るステップと、 (e) 他のプロポーズされたメンバーシップ・リストから選出されたプロポーズされたメンバーシップ・リストを選択するステップと を有する方法。 【請求項2】 接続データを決定する(a)ステップが 選択されたノードが分散型コンピュータ・システムの他のノードのどのノードと通信しているかを決定するステップと、 接続されているノードを特定するデータを他のノードに同報通信するステップと、 接続されているノードからノード接続データを受け取るステップと、 接続されているノードからのノード接続データと接続されているノードを特定するデータを接続データを形成するために組み合わせるステップと を有する請求項1記載の方法。 【請求項3】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)ステップが プロポーズされたメンバーシップ・リストと他の全てのプロポーズされたメンバーシップ・リストとが一致すことを確認するステップ を含む請求項1記載の方法。 【請求項4】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)ステップがさらに プロポーズされたメンバーシップ・リストと他の全てのプロポーズされたメンバーシップ・リストとの間の不一致を検出するステップを有し、 その不一致に応じて(a)~(d)を繰り返す請求項1記載の方法。 【請求項5】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)ステップが 選出されたプロポーズされたメンバーシップ・リストがまとまってクオーラムを形成するノードを表すことを決定するステップを 含む請求項1記載の方法。 【請求項6】 選出されたプロポーズされたメンバーシップ・リストがまとまってクオーラムを形成するノードを表すことを決定するステップが、 分散型コンピュータ・システムの動作しているノードの数を推測するステップを 含む請求項5記載の方法。 【請求項7】 分散型コンピュータ・システムの動作しているノードの数を推測するステップが 最初に述べたプロポーズされたメンバーシップ・リストに表されているノードの数を決定するステップと、 加入ノードの数を追加するステップと、 自発に離脱するノードの数を減算するステップと を有する請求項6記載の方法。 【請求項8】 プロセッサとメモリを含むコンピュータで使用するコンピュータ可読媒体であって、 (a) 分散型コンピュータ・システムのノードの接続性を表す接続データを決定し、 (b) プロポーズされた新しいクラスタのプロポーズされたメンバーシップ・リストを形成するために接続データに最適化基準を適用し、 (c) プロポーズされたメンバーシップ・リストを接続されているノードに同報通信し、 (d) 接続されているノードから他のプロポーズされたメンバーシップ・リストを受け取り、 (e) 他のプロポーズされたメンバーシップ・リストから選出されたプロポーズされたメンバーシップ・リストを選択する ことによって、分散型コンピュータ・システムのノードのメンバーシップをコンピュータに決めさせるコンピュータ命令を含むコンピュータ可読媒体。 【請求項9】 接続データを決定(a)する際に 選択されたノードがbpの他のノードのどのノードと通信しているかを決定し、 接続されているノードを特定するデータを他のノードに同報通信し、 接続されているノードからノード接続データを受け取り、 接続されているノードからのノード接続データと接続されているノードを特定するデータを接続データを形成するために組み合わせる 請求項8記載のコンピュータ可読媒体。 【請求項10】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)際に プロポーズされたメンバーシップ・リストと他の全てのプロポーズされたメンバーシップ・リストとが一致すことを確認し、 を含む請求項8記載のコンピュータ可読媒体。 【請求項11】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)際に、さらに プロポーズされたメンバーシップ・リストと他の全てのプロポーズされたメンバーシップ・リストとの間の不一致を検出し、 その不一致に応じて(a)~(d)を繰り返す請求項8記載のコンピュータ可読媒体。 【請求項12】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)際に 選出されたプロポーズされたメンバーシップ・リストがまとまってクオーラムを形成するノードを表すことを決定する請求項8記載のコンピュータ可読媒体。 【請求項13】 選出されたプロポーズされたメンバーシップ・リストがまとまってクオーラムを形成するノードを表すことを決定する際に、 分散型コンピュータ・システムの動作しているノードの数を推測する請求項12記載のコンピュータ可読媒体。 【請求項14】 分散型コンピュータ・システムの動作しているノードの数を推測する際に 最初に述べたプロポーズされたメンバーシップ・リストに表されているノードの数を決定し、 加入ノードの数を追加し、 自発に離脱するノードの数を減算する 請求項13記載のコンピュータ可読媒体。 【請求項15】 プロセッサと、 プロセッサに接続されたメモリと、 (i)メモリーからプロセッサで実行し、かつ(ii)プロセッサによって実行されたとき、コンピュータに分散型コンピュータ・システムのノードのメンバーシップをであって (a) 分散型コンピュータ・システムのノードの接続性を表す接続データを決定し、 (b) プロポーズされた新しいクラスタのプロポーズされたメンバーシップ・リストを形成するために接続データに最適化基準を適用し、 (c) プロポーズされたメンバーシップ・リストを接続されているノードに同報通信し (d) 接続されているノードから他のプロポーズされたメンバーシップ・リストを受け取り、 (e) 他のプロポーズされたメンバーシップ・リストから選出されたプロポーズされたメンバーシップ・リストを選択する ことによって決定させる欠陥検出モジュールと を有するコンピュータ・システム。 【請求項16】 接続データを決定(a)する際に 選択されたノードがbpの他のノードのどのノードと通信しているかを決定し、 接続されているノードを特定するデータを他のノードに同報通信し、 接続されているノードからノード接続データを受け取り、 接続されているノードからのノード接続データと接続されているノードを特定するデータを接続データを形成するために組み合わせる 請求項15記載のコンピュータ・システム。 【請求項17】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)際に プロポーズされたメンバーシップ・リストと他の全てのプロポーズされたメンバーシップ・リストとが一致すことを確認し、 を含む請求項15記載のコンピュータ・システム。 【請求項18】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)際に、さらに プロポーズされたメンバーシップ・リストと他の全てのプロポーズされたメンバーシップ・リストとの間の不一致を検出し、 その不一致に応じて(a)~(d)を繰り返す請求項15記載のコンピュータ可読媒体。 【請求項19】 選出されたプロポーズされたメンバーシップ・リストを選択する(e)際に 選出されたプロポーズされたメンバーシップ・リストがまとまってクオーラムを形成するノードを表すことを決定する請求項15記載のコンピュータ・システム。 【請求項20】 選出されたプロポーズされたメンバーシップ・リストがまとまってクオーラムを形成するノードを表すことを決定する際に、 分散型コンピュータ・システムの動作しているノードの数を推測する請求項19記載のコンピュータ・システム。 【請求項21】 分散型コンピュータ・システムの動作しているノードの数を推測する際に 最初に述べたプロポーズされたメンバーシップ・リストに表されているノードの数を決定し、 加入ノードの数を追加し、 自発に離脱するノードの数を減算する 請求項20記載のコンピュータ・システム。
124 paragraphs, as filed
Description: TECHNICAL FIELD [Detailed description of the invention]
【0001】
Background of the invention The present invention relates to defect tolerance of distributed computer systems, especially to a robust mechanism for determining which nodes form a cluster and access shared resources in a failed distributed computer system. [0002]
The issue of providing membership services in distributed computer systems has become of great academic and industrial interest. A distributed system, a Parallel Database (PDB) system available from Sun Microsystems, Palo Alto, Calif., Uses the Cluster Membership Monitor to keep track of member nodes as they change cluster membership. It also provides a mechanism for coordinating the reconfiguration of cluster applications and services. Here we define a common problem of membership in a cluster of computers as some nodes in the cluster not being fully connected, and here we propose a solution. [0003]
The general issue of membership is encapsulated by the design goals for the membership algorithm outlined below. Further describe the problem of attempting an address after stating those goals. [0004]
1. A uniform and robust membership algorithm that does not involve a system architecture that can tolerate continuous defects in nodes, links, storage devices or communication media. In other words, a single flaw does not render the cluster unusable. 2. Data integrity is not compromised by a number of synchronized defects. This is achieved by the following points. (a) Have only one cluster with a majority quorum that operates at any given time. (b) Clusters with a majority quorum do not reach inconsistent agreements. (c) Remove isolated and defective nodes from the cluster within a certain period of time. (d) Timely shield non-member nodes from shared resources. [0005]
The hardware architecture of traditional distributed computer systems has been confused by the inherent problems with membership algorithms. For example, consider the configuration shown in FIG. In this figure, each node 100A-D is connected to two switches 101-102. However, the two links have failed, preventing nodes 100A and 100D from communicating with each other. Traditional membership algorithms are incapable of dealing with such deficiencies and will not reach an agreement on the surviving majority quorum. These algorithms assume that all nodes are connected and cannot handle the problem of delimited networks. We need a general algorithm that handles the problems of partitioned networks as well as unsegmented networks. [0006]
Split-brain or possible split-brain status Complexity arises when it is necessary to make a decision about. For example, consider the configuration shown in FIG. In this configuration, the current quorum algorithm shuts down the entire cluster when there is no communication between node {200A, 200B} and node {200C, 200D} so that there are subclusters of the same number of nodes. become. Other situations that the current algorithm cannot handle include when there are two nodes in the system that cannot share external devices. [0007]
The above example illustrates a new combination of problems with membership and quorum algorithms that is not possible under the simpler architecture of traditional distributed computer systems that assume that the network is fully connected. doing. The way to solve this new problem is to integrate membership and quorum algorithms closer together to provide a flexible algorithm that maximizes cluster applicability and performance for the user to see. [0008]
The impact of external device configuration is a matter of occlusion failure. Shared resources (mostly disks) in a clustered system are blocked by the intervention of nodes that are not part of the cluster. In some distributed computer systems, this blocking problem is simpled by the fact that the cluster has only two nodes, which are connected to all shared resources. Nodes that remain in the cluster maintain all shared resources and do not allow non-member nodes to access those resources until they become part of the cluster. Such a simple operation cannot be performed on an architecture in which not all disks are connected to all nodes. Given that SPARC storage arrays (SSAs) are doubly ported, there is a need for new ways to effectively block non-member nodes from shared resources. [0009]
The cluster membership monitor, or CMM, responsible for membership, quorum, and algorithm occlusion failures, handles state transitions that lead to membership changes. The transitions are listed below. [0010]
Node defect: When a node fails, it initiates a cluster reconfiguration that does not include the defective node in its membership. · Join node: After reconfiguration, the node restarts and can join the cluster after being accepted as a new member by other members of the cluster. · Spontaneous leave: A node can leave the cluster at any time, and the remaining members of the cluster reconfigure the next generation of the cluster. Communication Defects: Handles communication defects where the Cluster Membership Monitor separates one or more nodes from nodes with a majority quorum. The detection of communication defects, that is, the detection that the communication graph is not fully connected, is the responsibility of the communication monitor, which is not part of the membership monitor. The communication monitor informs the membership monitor of a communication defect and the membership monitor handles it via reconfiguration. [0011]
It is important to know that the CMM does not guarantee that the entire system is healthy and that the application is given to any given node. The only guarantee made by CMM is that the system hardware is booted and working, and that the system is in existence and working. [0012]
Exactly define what defects are considered in the design of this system. There are three defects to consider. Node defects, communication defects, device defects. It should be noted that defects in client nodes, terminal connectors, and management workstations are not considered by the journal system to be defects. [0013]
Node defect: Suppose a node fails when it stops sending periodic heart-beat messages (SCI or CMM) to other members of the cluster. In addition, the node is believed to behave in a non-malicious manner, and it is believed that a node identified as defective by the system will not send conflicting information to other members of the cluster. Nodes can fail intermittently, as in the case of a primary deadlock, and can be seen as defective by what remains in the system, such as in the case of defective adapters or switches. .. The cluster member monitor must be able to handle all of these cases and must remove the defective node from the system in a given amount of time. Communication Defects: Private communication media fail due to switch defects, adapter card defects, cable defects, and many software layer defects. These defects are masked by the Cluster Communication Monitor (CCM or CIS) so that the Cluster Membership Monitor does not handle certain defects. In addition, the Cluster Membership Monitor sends the message through an available link on the medium. Defects in individual links do not affect the correct operation of the CMM. The only communication flaw that affects the operation of the CMM is the overall impairment of communication with member nodes. This is essentially the same as a node defect, as there is no physical path to send a heartbeat message over a private communication medium. In switch architectures such as the Energizer 2 release, failure of all switches is logically equivalent to simultaneous failure of n-1 nodes. Where n is the number of nodes in the system. Device Deficiency: The device that affects the operation of the Cluster Membership Monitor is a quorum device. Traditionally, they are Sparc Stptage It was a disk controller for Array (SSA). However, in some distributed computer systems, the disk can be used as a quorum device. It should be noted that quorum device flaws are equivalent to node flaws, and that CMMs in some traditional systems cannot use quorum devices unless they are run on two node clusters. .. [0014]
Some distributed computer systems are said to have no single point of flaw. Therefore, a single node defect must be tolerated as well as a continuous defect of n-1 nodes in the system. Given the above discussion of communication defects, this specification shows that the overall impairment of the communication medium in the system is unacceptable. It should be possible to tolerate more defects than a single defect at any given time, if it is not possible or desirable to tolerate the overall impairment of the communication medium. First, what is a cluster and the following defines how various flaws affect it. [0015]
A cluster is defined as having N nodes, a private communication medium, and a quorum mechanism, the overall flaw in the private communication medium is equivalent to the flaw in the N-1 node, and the flaw in the quorum mechanism is the flaw in one node. Is equivalent to. [0016]
The following fault tolerance goals for the cluster membership monitor are shown. A cluster of N nodes with N 3, but partially lacking [N / 2] -1 nodes, private communication media, quorum mechanisms should be able to serve and access data. For two clusters, the cluster can tolerate only one of the following flaws: [0017]
-Impairment of one of the nodes. -Impairment of private communication media. In this case, logically one node is impaired Is equivalent to. -Impairment of quorum devices. -Impairment of one of the nodes and private communication media. In this case, logically one Equivalent to the impairment of a node in. [0018]
Total loss of communication media in a system with three or more nodes is a double defect (because both exchanges are inactive) and the system is not required to tolerate such defects. [0019]
Abstract of the invention According to the present invention, cluster membership in a distributed computer system can be achieved by determining which other nodes each node is communicating with and by distributing its connectivity information through the nodes of the system. It is determined. Therefore, each node can determine a new cluster optimized based on connectivity information. Each node has information about the node with which it is communicating and similar information about each other node in the system. Therefore, each node has complete information about the connectivity of all nodes that are directly or articulated. [0020]
Each node applies optimization criteria to connectivity information to determine the best new cluster. Data representing the optimal new cluster is broadcast to each node. The optimal new cluster determined by the various clusters is collected by each node. Each node has data representing the proposed new class that was found to be optimal for each node. Each node uses that data to elect a new cluster from various proposed new clusters. For example, if there are more new clusters proposed than others, they will be elected as new clusters. Each node receives the same proposed new cluster from a node that is a potential member of the new cluster, so the new cluster membership unanimously reaches. In addition, each node has more complete information about the nodes of potential members of the new cluster, so the new cluster obtained is undoubtedly in a relatively optimal configuration. [0021] [0021]
Agreements between processors in distributed systems of which processors are members are a fundamental issue in the design of highly useful distributed systems. Processors shut down, are defective, come back, and membership changes occur when new processors are added. There is currently no agreed definition of processor membership issues. And existing membership protocols provide substantially different guarantees for those services. The protocol of interest occurs when the current membership processor agrees on a set of member nodes and the change in membership is logically equivalent on different nodes. [0022]
Due to the flaws mentioned above, cluster membership is divided into two or more fully connected subsets of nodes that have a voting majority, a voting minority, or an exact half of a vote. The first two cases are resolved by allowing a subset of the cluster with the next generation of forming majority votes. In the last case, the tiebreak mechanism must be adopted. Some cluster membership algorithms have the advantage of the limitations imposed by the two node architectures that solve those problems. Generalized for an architecture containing three or more nodes, the following new problems are solved by the algorithm according to the invention. [0023]
1. Resolve quorum and membership when not all pairs of nodes share a common external device. [0024]
Integration of quorum and membership algorithms may be required for systems with three or more nodes. With more than two nodes in a distributed system, external devices don't really need to solve membership and quorum issues. A system with only two nodes does not require an external quorum mechanism. [0025]
In some distributed computer systems, this external device is a controller that resides on a disk or SSA. This choice of quorum device has unfavorable properties that adversely affect the overall capacity of the cluster, especially for disks. [0026]
It is more complicated for a four-node system with an architecture where not all nodes are allowed to connect to all external devices. In that architecture, some combinations of nodes forming a cluster do not share any external device other than the communication medium, and therefore allowing such a cluster to exist would allow other quorum mechanisms. I need. Public networks have serious security shortcomings, leaving only a last resort, human intervention. If the majority of votes are not automatically determined, use this method as fully described below. The new user interface is described in detail below. [0027]
2. Authorization of the majority quorum request adopted to change membership. When composed of three or more nodes, requiring more than half of the total votes cast for the majority quorum limits user flexibility. In a four-node system, two unequal nodes form a cluster. The modified algorithm is based on the quorum requirements for current membership and voting of joining nodes. [0028]
3. Deal with "voluntary withdrawal" of cluster members with a hint of lower majority quorum requirements. The original algorithm believed that simultaneous cluster shutdown of more than half of the nodes would be the part that excluded those nodes. To avoid resulting quorum loss and complete cluster shutdown, the new algorithm uses explicit shutdown notification by the node to reduce quorum requests. [0029]
4. Joining process when a node is partitioned With a two-node configuration and a tiebreak quorum device, it is not possible for the two nodes to form an independent cluster when communication between the two nodes is broken. In a dynamic quorum request with three or more nodes and item 2, the opposite state (two or more independent clusters) because a subset of all fully connected nodes form a cluster with the quorum. ) Is possible. The algorithm according to the invention distinguishes between the initial formation of a cluster and its subsequent joining. With the exception of this initial join, the joins cannot form a cluster independently, the nodes only join an existing cluster. The user interface for that is discussed below. [0030]
5. Handling defects that occur during membership algorithms In dynamic quorum requests, inconsistencies between nodes in the number of votes requested for quorum can occur when defects occur during reconstruction. To avoid the possibility that two or more subsets have quorum and form independent clusters, the transformation algorithm imposes restrictions on what they join. Those who join join the existing cluster completely as they are. [0031]
6. Handling of the partial connection status shown in Figure 1. In such a scenario, the original algorithm does not reach an agreement. The algorithm converges when a set of nodes agree on the same membership proposal, but that condition is never met. In the algorithms according to the invention, when this condition is suspected (using a timeout), a node transforms their membership proposals into the most connected subset. [0032]
In the following sections, we will discuss the format of the messages exchanged by the cluster daemon, decide what is the best membership and how to choose it, and in addition to what was done above, the members Identify the assumptions made by the ship algorithm, describe how changes in membership have occurred, describe the membership algorithm, and how the CMM suspends and resumes the set of registered processes. Explain what to do, discuss how the CMM checks the consistency of the configuration database, and identify the new user interface needed. [0033]
4.1 CMM message Membership monitors for different nodes in the cluster exchange messages with each other to indicate that they are alive, ie exchange heartbeats and initiate a cluster reconfiguration. It is possible to distinguish between those two types of messages, but in reality they are the same message, called the RECONF_msg message, which is reconfigured on the receiving node. [0034]
Each RECONF_msg contains the following fields: · Seq_num, a sequence number that distinguishes between different reconstructions. -Vector containing membership voting for node i, M<sub>i</sub>.. Vector S containing a view of the most recent stable membership node i<sub>i</sub>.. Vector V containing connectivity information for node i<sub>i</sub>.. Vector SD containing a view of node i of the node that voluntarily left the cluster when the most recent stable membership was established. Node St<sub>i</sub>State. -The node id of the original node. · Vector J containing a view of node i of the node trying to join<sub>i</sub>.. A flag that indicates whether the original node is itself seen as a subscribed node. [0035]
4.2 Definitions and assumptions The membership algorithm assumes that the clusters are nodes of equal value, that is, homogeneous clusters. The membership algorithm is based on a set of rules described in the following priorities and used in the development of the membership algorithm. 1. The node contains itself in the proposed set. 2. Nodes vote for nodes already in the cluster for the node they are trying to join. 3. The node proposes the set that has the most fully connected nodes, including itself. 4. All nodes agree on statically determined priorities among the nodes. That is, the node with the lower number takes precedence over the node with the higher number. [0036]
The above set of rules defines a hierarchy of rules with statically determined priorities at the bottom of the hierarchy. The set of rules above is the optimal membership set, ie [Number 1]
<img file="JP2001521222A_D0001.tif" />To define. [0037]
Optimal membership set in a cluster with one or more flaws, [Number 2]
<img file="JP2001521222A_D0002.tif" />Finding is a computer-intensive task. This problem is addressed from the standpoint of selecting the optimal subset of the set of nodes according to the definition of optimality derived from the above rules. Assuming the cluster consists of N nodes [Number 3]
<img file="JP2001521222A_D0003.tif" />Finding is equivalent to finding the optimal matrix size for MxM from an NxN size matrix. Here M <N. This problem is a well-known problem of "N-selection M", which is well known as a binomial coefficient. The solution to this problem is 0 (2), assuming the system is homogeneous and each node can be represented by either 0,1 or -1.<sup>N</sup>) It is a complex number. Find the optimal subset The cost is so high that you have to stop it when N is large, but if N <20, this cost is not enough to stop. Therefore, for systems with 16 or less nodes, it is recommended to use a comprehensive search method to find the optimal set. For systems with more than 20 nodes, a discovery algorithm suitable for the optimal solution is desirable. [0038]
Assume that the failing node broadcasts RECONF_msg to all other members of the current cluster. Also assume that the node trying to join the cluster does it in the initial state and its sequence number is reset to 0. Similarly, messages with a sequence number higher than or equal to their own sequence number and with a rank value at most one behind are processed. However, there are serious exceptions. When a message comes from a node that is flagged as'subscriber', it will be processed even if its state is not fresh (two or more behind). They are the nodes that are going to join the cluster and must receive their initial message. All these assumptions are made by embodiments of the membership algorithm according to the invention. [0039]
4.3 Change of membership There are several ways nodes can be reconfigured to result in membership changes according to the algorithms provided in the next section. Below is a list of them. 1. Join: This is when a node forms a new cluster or joins an existing cluster. (a) First join: Performed only on the first node of the cluster and via the new command pdbadmin startcluster. The command signals the CMM running on the node. There are no nodes in the cluster, so I don't plan to be asked by other nodes. This command is issued only once at the beginning of the life of the cluster. The lifetime of a cluster is the period from when the pdbadmin start cluster is issued until there are no members in the cluster. If an additional pdbadmin startcluster command is issued, in the worst case the system will compromise data consistency, the nodes will be isolated, or in the most likely case an error will occur and this command will be issued incorrectly. It is to suspend the node. (b) Subscription following the first subscription: This subscription is a common pdbadmin Made by the startnode command, a node or set of nodes joins the cluster. Nodes trying to join the cluster communicate with nodes that are already members of the cluster and try to find out if they can join with the membership algorithm. [0040]
2. Departure: This occurs when a node that was a member of the cluster voluntarily or involuntarily leaves the cluster. (a) Spontaneous withdrawal: The operator issues the pdbadmin stop node command to the node. This causes the node to finish the stop sequence. As a result, the node sends a message to all the nodes forming the cluster indicating that the cluster is about to leave. This information can and is used by membership algorithms for membership optimization. (b) Involuntary withdrawal: There are two different cases. i. A node can complete its stop or stop sequence and then "clean up". More importantly, as far as the CMM is concerned, it is possible to send the same message as when the node spontaneously left, that is, to inform other members of the cluster that the node will no longer belong to the cluster. The optimizations performed for the spontaneous withdrawal of the node are actually performed here as well. The node leaves the cluster at the request of an application program with unique privileges. ii. The system panics if the node does not complete the abort sequence. This is the most difficult flaw to handle and usually detects the loss of heartbeat messages from the flawed node. This flaw is indistinguishable from a network flaw in an asynchronous distributed system. [0041]
4.4 Algorithm This section describes the membership algorithm based on the following assumptions and definitions: The user interface used in this algorithm will be described later to make the algorithm flow "clean". Before going into the explanation of the algorithm, I will describe the rules required to realize the membership algorithm. [0042]
Each node can only vote once, regardless of whether it is already a member of the cluster or is about to join the cluster. [0043]
-Each node i updates its "connectivity state matrix Ci" as soon as it receives it from another node. Matrix Ci is an understanding of node i with respect to all connectivity of the system. If node i does not receive from node j, or node j is down or unreachable, then element e of Ci<sub>ij</sub>Is marked to zero. In addition, mark all elements in line j as "null". This implies that node i has no information about the connectivity of node j. For the other rows of the matrix, node i takes the kth row of its connectivity matrix to the connectivity vector V.<sub>k</sub>Update by replacing with I do. [0044]
-Each node i is initially C in its "RECONF_msg"<sub>i</sub>I th Proposal membership set voting for the line [Number 4]
<img file="JP2001521222A_D0004.tif" />Include as. Set proposed by node i [Number 5]
<img file="JP2001521222A_D0005.tif" />Is the vector V<sub>i</sub>Note that this is not the case.
[Number 6]
<img file="JP2001521222A_D0006.tif" />Is a proposed set that describes voting for other nodes in binary form, while V<sub>i</sub>Is a "state" vector, no in the system Handles connectivity.
[Number 7]
<img file="JP2001521222A_D0007.tif" />Node disagrees with stable membership and V<sub>i</sub>Subset of new V when you need to propose as a membership set<sub>i</sub>Then It has different elements, each of which is a node id and a binary voting value. [0045]
· Each node i sets the total number of existing nodes in the current view of cluster membership, whether agreed or proposed, to the local variable N.<sub>i</sub>Keep inside. Local variable N<sub>i</sub>Follows the following rules during the execution of the membership algorithm. (a) N<sub>i</sub>Is [Number 8]
<img file="JP2001521222A_D0008.tif" />It is initialized to the cardinality of. (b) N<sub>i</sub>Increments (increments) for each node trying to join the cluster (Every time the node id embedded in the message is checked by the receiver thread, an increment of 1 is forced for each node.) (c) N<sub>i</sub>Is each no, as defined in 2 (b) i of Section 4.3. Decrement (decrease) for withdrawal or voluntary withdrawal. (Execution by receiver thread) (d) At the end of the membership algorithm, quorum, N<sub></sub><sub></sub><sub>i</sub>It is decided about this concept of. [0046]
· At the end of the membership algorithm, the nodes forming the cluster are a new set of member nodes. [Number 9]
<img file="JP2001521222A_D0009.tif" />I agree with. This set was used in the next run of the membership algorithm, when all the nodes that were previously part of it matched. [Number 10]
<img file="JP2001521222A_D0010.tif" />Is assumed to have a set of. [0047]
Each node i will have the same sequence number seq_num for all nodes in the current cluster before entering the membership algorithm. In addition, each node has its connectivity state matrix C.<sub>i</sub>Will have .. C<sub>i</sub>Is an n × n matrix, where n is the current cluster configuration file ( That is, the maximum number of nodes defined by the current cdb file). [0048]
-Each joining node sets the variable joining_node to true (TRUE). Node once [Number 11]
<img file="JP2001521222A_D0011.tif" />When you become a member of, you are no longer a joiner and you are joining_ node is set to FALSE. [0049]
-The node achieving the first subscription sets the variable start_cluster to true (TRUE). Nodes trying to join the cluster set the variable start_cluster to false (FALSE). [0050]
All nodes in the cluster get information about various timeout values from the configuration file. The symbols T1, T2, .... are used to indicate the different possible timeout values. All of these values must match for all nodes and are set to reasonable values to introduce communication and queue delays. [0051]
The algorithm can be described for each node i as follows.
[table 1]
<img file="JP2001521222A_D0012.tif" />[Table 2]
<img file="JP2001521222A_D0013.tif" />[Table 3]
<img file="JP2001521222A_D0014.tif" />[Table 4]
<img file="JP2001521222A_D0015.tif" />[Table 5]
<img file="JP2001521222A_D0016.tif" />[Table 6]
<img file="JP2001521222A_D0017.tif" /> 【0052】
In the above algorithm, it is assumed that there is a way to send a message to all the nodes that are part of the cluster. If the node is down or unreachable, it is considered processed by the previous configuration and the matrix [Number 12]
<img file="JP2001521222A_D0018.tif" />It is reflected in. [0053]
In the above algorithm, the function membership_proposal () contains all nodes that are not in the DOWN state.<sub>i</sub>Members based on -Return a ship proposal. The function is in the proposal, [Number 13]
<img file="JP2001521222A_D0019.tif" />If all of the above are not included, all applicants are excluded from the proposal. The important function is the stable_proposal () function. This function is a proposed set [Number 14]
<img file="JP2001521222A_D0020.tif" />Determines if is agreed by all other members of the set. To count the number of votes, node i is a proposed set from another node (ie, [Number 15]
<img file="JP2001521222A_D0021.tif" />And self-set [Number 16]
<img file="JP2001521222A_D0022.tif" />(However, j i) Need to be compared with. The function share_quorum_dev (), implemented by using a CCD dynamic file, tells the membership algorithm when two nodes share a quorum device, such as a cluster of two nodes. The binary function reserve_quorum () returns false only if the device is already reserved by one other node. The function wait_for_user_input () is described in detail below. [0054]
The function propose_new_membership () is called and is the optimal subset according to the optimality conditions described above. [Number 17]
<img file="JP2001521222A_D0023.tif" />Find out. The function is [Number 18]
<img file="JP2001521222A_D0024.tif" />Thoroughly test a subset of combinations until you find the first well-connected set. The fully_connected (prop) function returns true if the candidate proposal prop is included in all of the proposals that are members of the prop. if, [Number 19]
<img file="JP2001521222A_D0025.tif" />Note that the proposal does not change if is already well connected. if, [Number 20]
<img file="JP2001521222A_D0026.tif" />If is not fully connected, the joiner is present in the proposal Also note that we do not. Finally, the find_optimal_proposal () and get_next_proposal () functions do a thorough search. [0055]
4.5 User interface So far, we haven't yet explained how to tie-break with user input in a potential split-brain situation. This subsection reveals how to achieve a tiebreaker. [0056]
If the situation where both set X of nodes and set Y of different nodes have exactly N / 2 (N is the number of nodes in the previous cluster) votes requires operator assistance. Is. If both X and Y candidates are 1 and they share a quorum device, there is no need to ask for input from the operator. In both situations, the node waits for user input by executing a wait_for_user_input () call in the membership algorithm. A call to wait_for_user_input () spawns a "print" thread that will continuously print a message telling the operator that a potential tie break should be made. The message identifies a set X or Y that must be notified to the appropriate node that it must be shut down or kept alive. The operator issues the pdbadmin stopnode command on a set of nodes and the new command pdbadmin A tiebreaker is made by issuing continue for another set. The set that receives the stop command aborts, the other set stops printing the message and continues its reconstruction. Alternatively, the operator can issue the crustm reconfigure command. This command is a valid option if a communication breakdown occurs but the operator repairs it. The occurrence of the clustm reconfigure command causes the generation of a new reconfiguration. If the operator issues any other command at this time in addition to pdbadmin stopnode, clustm reconfigure or pdbadmin continue, the command reader thread will be waiting for one of those commands. Do not signal the thread and simply ignore the command. The print thread, on the other hand, keeps printing those messages once every few seconds to inform the operator that some action should be taken immediately. [0057]
The function wait_for_user_input () executed by the transitions thread is executed as follows. [0058] [0058]
[Table 7]
<img file="JP2001521222A_D0027.tif" /> 【0059】
The above sequence of actions puts the transitions thread to sleep on the condition variable state_change_cv. The condition variable is flagged under the following conditions. -The user issues the continue command. -The user issues the stopnode command. -The user forces the reconfiguration. The node receives a message indicating that the remote node in its current membership set is down. -The node has not received a message for node_down_timeout from a remote node in its current membership set. [0060]
All of these actions are appropriate for flagging the transitions thread, allowing the user to generate the correct set of commands, ensuring that only one primary group stays operational within the cluster. You can make it. [0061]
5. Failure Fencing and Resource Migration One of the other components of the system that comes from the new architecture and needs to be modified is the defect shielding mechanism that is sometimes adopted in distributed computer systems. This section discusses solutions to the general problem of resource transfer and the specific problem of defect shielding. The resulting solutions are different array topologies (cascaded, n + 1, cross-connected (cross-connected), other topologies), as well as different software configurations (CVM with Netdisk, standalone). It is common in the sense that it deals with VxVM, or other configurations). The solution also deals with array configurations connected at the intersection of two nodes, rather than as a special case. [0062]
Assumptions and general solutions are discussed next. It is followed by a brief commentary on how to solve the resource transfer problem (highly available disk groups, HF / NFS file systems, logical IP addresses on public networks). [0063]
5.1 Assumptions In a shared disk configuration with CVM and Netdisk, the master node for all NetDisk devices and its backup nodes provide direct physical access to the underlying physical device. [0064]
In a non-shared configuration with VxVM, it is assumed that each node with the primary ownership of a set of disk groups has direct physical access to the devices belonging to those disk groups. Ru. More specifically, if node N has the primary ownership of disk group G, then all of the disks belonging to disk group G are connected to N and of the storage device labeled D. Can be found in the set. [0065]
It is assumed that information about the primary and backup ownership of NetDisk devices or other resources is maintained in the Cluster Configuration Database (CCD) and is consistently available to all nodes. .. This assumption can be enforced by using the dynamic part of the CCD. In particular, the CCD may be asked to obtain the above information when the steps for defect shielding and resource transfer are performed during the reconstruction process outlined in the next subsection. These steps are only performed after cluster membership has been determined and a quorum has been obtained. [0066]
5.2 Defect shielding Some distributed computer systems have backup nodes on every node. The node (main) and its backup nodes share a common set of devices to which they are connected. This is B (N<sub>i</sub>) = N<sub>j</sub>Indicated by. In a CVM plus Netdisk configuration, the backup node becomes the master of the set of NetDisk devices owned by the defective node. In a VxVM configuration, the backup node is the primary owner of the set of disk group resources owned by the defective node. [0067]
Ni indicates the node of the cluster, D<sub>i</sub>In (1 or more SSA and / or Mult It shall indicate a storage device (consisting of ipacks). Suppose there are four nodes in the cluster, and assume that there is the following relationship for the cascaded configuration: That is, we assume the relationship of B (N1) = N4, B (N2) = N1, B (N3) = N2, B (N4) = N3. For the n + 1 configuration, B (N1) = The relationship N4, B (N2) = N4, B (N3) = N4 is given. Now N4 is Does not have a backup node. Finally, in cross-connections, the relationship between backup and primary is given by B (N1) = N2, B (N2) = N1, B (N3) = N4, B (N4) = N3. In the case of a two-node crossed connection, it is simply simplified to B (N1) = N2, B (N2) = N1. [0068]
5.3 General solution Suppose node i becomes a defect and as a result of this defect all other nodes are reconfigured. The surviving node j performs the following simple steps after membership and quorum have been determined. [0069]
[Table 8]
<img file="JP2001521222A_D0028.tif" /> 【0070】
Take over these devices implies that any of the interfaces provided by NetDisk for acquiring ownership of NetDisk devices can be used to achieve it. [0071]
Without mentioning the details of the syntax, it would be sufficient to describe how the node accurately determines which NetDisk device was mastering the defective node. .. That is, it would be sufficient to state that the CCD holds the information in its database and is asked to retrieve that information and whether the node currently being reconfigured is a backup of the defective node. [0072]
Depending on the distributed computer system, information about the primary ownership of the disk group is stored in a cdb file in the following format: [0073]
cluster.node. (). cdg: dg1dg2 cluster.nodc. (). Cdg: dg3dg4 [0074]
In a CCD, finding an equivalent representation of this information should be simple in order to make it available to all nodes in exactly the same way as a NetDisk device configuration. The spare representation to add for each node is, of course, a backup node in the same way as a NetDisk device configuration. For example: cdg: dg1, dg2: 0.1. Cluster Disk-The primary owners of groups dg1 and dg2 are node 0 and its backup node 1. [0075]
You can also ask the CCD or Volume Manager to find out the set of physical devices that accompany each particular NetDisk virtual device or particular disk group. [0076]
Finally, these are the steps that are taken when the defective node i is ready to join the cluster. Each of the other nodes j i executes this sequence in an undetermined step k of the reconstruction process. [0077]
[Table 9]
<img file="JP2001521222A_D0029.tif" /> 【0078】
The following sequence is performed by Node Ni, who wants to join, in step k + 1 of the reconstruction process. [0079]
[Table 10]
<img file="JP2001521222A_D0030.tif" /> 【0080】
At node j i, it is not possible to determine whether i is about to join the cluster or is about to undergo a reconfiguration, and whether it was already part of the cluster. This is not a problem. This is because it is a simple matter to determine whether Node Nj owns a resource that has part of its cluster membership as its primary owner, which causes it to take appropriate action. If the algorithm is implemented correctly [0081]
The approach in this section can be used to solve common resource transfer problems in certain distributed computer systems. Resources that should be highly available are transferred from defective nodes to surviving nodes. Examples of such resources are disk groups in non-shared database environments, disk groups for HA-NFS file systems and logical IP addresses. Any resource can be designated as a master and backup node in the CCD. For example, a logical IP address can be transferred from a defective node to a surviving node. For the switch execution, the backup node must release the resources of the join node one step before the join node takes over its resources. [0082]
5.4 Restrictions on disk groups (restriction) Some distributed computer systems do not allow any connection from the nodes of the cluster to the array. This is because disk groups distributed across several arrays cannot be relocated to multiple other nodes in the cluster and must be relocated to one node as a whole. For example, consider a configuration in which four nodes N1, ... N4 and four array devices D1 ... D4 exist, and assume that node Ni has a physical connection to array Dj. Further assume that N2 is physically connected to D1 and D3, N3 is physically connected to D2 and D4, and N1 and N4 have no other connections. [0083]
Suppose node N2 has a disk group G distributed across multiple disks in arrays D1 and D3. If N2 becomes defective, G cannot import the entire disc into N1 or N3 because not all of its discs are visible on N1 or N3. Such configurations are not supported by some distributed computer systems. If a node owns a disk group and that node becomes defective, the entire disk group should be taken over by one of the surviving nodes. This does not limit the topology of the array, but it does limit how the data is distributed across the array. [0084]
5.5 Relocation with minimal burden One of the most time-consuming operations in a system is the layout of data. We propose to minimize this for upgrading from an existing 2-node cluster to a 3-node cluster or from an existing 3-node cluster to a 4-node cluster. It doesn't do this dynamically. The cluster is shut down and restarted. The only criterion is to allow access to mirrors and primary copies of data from the same node without relaying the entire volume and / or disk group. This requires the addition of an adapter card. [0085]
The above description is for illustrative purposes only and is not intended to be limiting. The present invention is therefore described only by the claims and equivalents.
[Simple explanation of drawings]
[Figure 1]
It is a block diagram of a distributed computer system in which communication between two nodes and two respective exchanges have failed. [Figure 2]
It is a block diagram of a distributed computer system including a dual port device.
Every citation, both ways
| Document | Relation | Office | Cited during |
|---|---|---|---|
| JP2001109726A | Cited by | Japan | Search report |
| JP2006004434A | Cited by | Japan | Search report |
| KR20150058280A | Cited by | Republic of Korea | Search report |
| JP2015057696A | Cited by | Japan | Search report |
| JP2010186472A | Cited by | Japan | Search report |
| JP2006004433A | Cited by | Japan | Search report |
| JP2015535970A | Cited by | Japan | Search report |
| JP2002049601A | Cited by | Japan | Search report |
| JP2015535970A | Cited by | Japan | Search report |
| JP2015057696A | Cited by | Japan | Search report |
| JPH02125544A | Cites | Japan | Search report |
| JPH05216845A | Cites | Japan | Search report |
| JPH0793265A | Cites | Japan | Search report |
| JPH0844690A | Cites | Japan | Search report |
| JPH09171502A | Cites | Japan | Search report |
| JPH10301905A | Cites | Japan | Search report |
| JPH1093655A | Cites | Japan | Search report |
| JPS62197860A | Cites | Japan | Search report |
17 members in 7 offices
Priority claims9
| Document | Office | Kind | Date |
|---|---|---|---|
| 08955885 | United States of America | – | |
| 95588597 | United States of America | A | |
| 95588597 | United States of America | A | |
| 9822161 | United States of America | W | |
| 9822161 | United States of America | W | |
| 1997955885 | – | – | – |
| 199822161 | – | – | – |
| US19970955885 | – | – | – |
| WO1998US22161 | – | – | – |
Members17
| Document | Office | Kind | |
|---|---|---|---|
| CA2306718A1 | Canada | A1 | |
| WO9921098A2 | World Intellectual Property Organization (WIPO) | A2 | |
| AU1105499A | Australia | A | |
| WO9921098A3 | World Intellectual Property Organization (WIPO) | A3 | |
| US5999712A | United States of America | A | |
| EP1025506A2 | European Patent Office (EPO) | A2 | |
| WO0054152A2 | World Intellectual Property Organization (WIPO) | A2 | |
| AU3868100A | Australia | A | |
| WO0054152A3 | World Intellectual Property Organization (WIPO) | A3 | |
| US6192401B1 | United States of America | B1 | |
| JP2001521222AThis record | Japan | A | |
| EP1159681A2 | European Patent Office (EPO) | A2 | |
| US6449641B1 | United States of America | B1 | |
| JP2002539522A | Japan | A | |
| EP1159681B1 | European Patent Office (EPO) | B1 | |
| DE60006499D1 | Germany | D1 | |
| DE60006499T2 | Germany | T2 |
3 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Decision of refusalJAPANESE INTERMEDIATE CODE: A02A02 | A02 | |
| Notification of reasons for refusalJAPANESE INTERMEDIATE CODE: A131A131 | A131 | |
| Written request for application examinationJAPANESE INTERMEDIATE CODE: A621A621 | A621 |
Numbers
- Publication
- 2001-521222
- Publication, DOCDB
- 2001521222
- Publication, EPODOC
- JP2001521222
- Application
- 2000517348
- Application, DOCDB
- 2000517348
- Application, EPODOC
- JP20000517348
Titles2
- Japanese
- 【発明の名称】分散型コンピュータ・システムにおいてクラスタ・メンバーシップを決定する方法
- English
- INDUSTRIAL APPLICABILITY A method for determining cluster membership in a distributed computer system.
Classification
- CPC, 11
- G06F9/5061
- G06F11/2033
- G06F11/2038
- G06F11/2046
- H04L41/06
- H04L41/0823
- H04L41/0873
- H04L43/0811
- G06F11/1425
- G06F2209/505
- H04L41/0893
- IPC, 4
- G06F11 20
- G06F9 50
- G06F11 00
- G06F13 00