IP Library Granted Patent US 7,620,845
Granted Patent B2
US 7,620,845 · App. 11/069,973 · Granted Nov 17, 2009

Distributed system and redundancy control method

Assignees: Kabushiki Kaisha Toshiba; Toshiba Solutions Corporation
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 7,620,845
App. No.
11/069,973
Granted
Nov 17, 2009
Kind
B2
Abstract

A distributed system using a quorum redundancy method in which a redundancy process is executed by at least Q processing elements of N processing elements communicable with each other, each of N processing elements includes a resynchronization determining unit for determining that an execution state of the processing element itself can be resynchronized with a latest execution state in the distributed system in the case where the processing element can communicate with at least F+1 elements (F=N−Q) already synchronized of the N processing elements at the time of rebooting the processing element, and a resynchronizing unit for resynchronizing the execution state of the processing element itself to the latest one of the execution states of the at least F+1 processing elements in accordance with the result of determination by the resynchronizing unit.

Claims (29)

1. A distributed system comprising N processing elements where N is an integer of 4 or more, the distributed system executing a redundancy process provided at least a quorum Q of the N processing elements are communicable with each other, at least one of the N processing elements comprising:

an execution state storage unit configured to store a latest execution state of the at least one processing element in a volatile memory;

a resynchronization determining unit configured to determine whether to resynchronize an execution state of the at least one processing element with a latest execution state of the distributed system upon rebooting the at least one processing element, the determination being to resynchronize provided the at least one processing element can communicate with at least F+1 of the N processing elements, where F+1>=2, F=N−Q, and F>=1; and

a resynchronizing unit configured to resynchronize the execution state of the at least one processing element to the latest execution state of the distributed system in accordance with the determination of the resynchronizing determining unit by:

comparing sequence numbers for the at least F+1 processing elements to determine which of the at least F+1 processing elements has a highest sequence number, the processing element determined to have the highest sequence number storing the latest execution state of the distributed system, and

copying the latest execution state from the processing element determined to have the highest sequence number.

2. The distributed system according to claim 1 ,

wherein the at least one processing element further comprises a progress information storing unit configured to store progress information comprising an indicator of progress of the redundancy process in the at least one processing element.

3. The distributed system according to claim 2 , wherein the progress information storing unit stores, as the progress information, the sequence number for the at least one processing element, and increments the sequence number by one each time the redundancy process executes an additional step.

4. The distributed system according to claim 1 ,

wherein Q is the ⅔ quorum coincident with the minimum integer more than ⅔ of N.

5. A method implemented in a distributed system comprising N processing elements where N is an integer of 4 or more, the distributed system executing a redundancy process provided at least a quorum Q of the N processing elements are communicable with each other, the method causing at least one of the N processing elements to:

store a latest execution state of the at least one processing element in a volatile memory;

determine whether to resynchronize an execution state of the at least one processing element with a latest execution state of the distributed system upon rebooting the at least one processing element, the determination being to resynchronize provided the at least one processing element can communicate with at least F+1 of the N processing elements, where F+1>=2, F=N−Q, and F>=1; and

resynchronize the execution state of the at least one processing element to the latest execution state of the distributed system in accordance with the determination of whether the processing element can be resynchronized by:

comparing sequence numbers for the at least F+1 processing elements to determine which of the at least F+1 processing elements has a highest sequence number, the processing element determined to have the highest sequence number storing the latest execution state of the distributed system, and

copying the latest execution state from the processing element determined to have the highest sequence number.

6. A computer-readable medium storing instructions for implementing a method in a distributed system comprising N processing elements where N is an integer of 4 or more, the distributed system executing a redundancy process provided at least a quorum Q of the N processing elements are communicable with each other, the method causing at least one of the N processing elements to:

store a latest execution state of the at least one processing element in a volatile memory;

determine whether to resynchronize an execution state of the at least one processing element with a latest execution state of the distributed system upon rebooting the at least one processing element, the determination being to resynchronize provided the at least one processing element can communicate with at least F+1 of the N processing elements, where F+1>=2, F=N−Q, and F>=1; and

resynchronize the execution state of the at least one processing element to the latest execution state of the distributed system in accordance with the determination of whether the processing element can be resynchronized by:

comparing sequence numbers for the at least F+1 processing elements to determine which of the at least F+1 processing elements has a highest sequence number, the processing element determined to have the highest sequence number storing the latest execution state of the distributed system, and

copying the latest execution state from the processing element determined to have the highest sequence number.

7. The method according to claim 5 , wherein the at least one processing element stores progress information comprising an indicator of progress of the redundancy process in the at least one processing element.

8. The method according to claim 7 , wherein the at least one processing element stores the progress information as the sequence number for the at least one processing element, and increments the sequence number by one each time the redundancy process executes an additional step.

9. The method according to claim 5 , wherein Q is the ⅔ quorum coincident with the minimum integer more than ⅔ of N.

10. The computer-readable medium according to claim 6 , wherein the at least one processing element stores progress information comprising an indicator of progress of the redundancy process in the at least one processing element.

11. The computer-readable medium according to claim 10 , wherein the at least one processing element stores the progress information as the sequence number for the at least one processing element, and increments the sequence number by one each time the redundancy process executes an additional step.

12. The computer-readable medium according to claim 6 , wherein Q is the ⅔ quorum coincident with the minimum integer more than ⅔ of N.

Assignments (6)
CHANGE OF CORPORATE NAME AND ADDRESS Recorded Feb 8, 2021
From: TOSHIBA SOLUTIONS CORPORATION
To: TOSHIBA DIGITAL SOLUTIONS CORPORATION
Reel/Frame 055259/0587 →
CORRECTIVE ASSIGNMENT TO CORRECT THE RECEIVING PARTY'S ADDRESS PREVIOUSLY RECORDED ON REEL 048547 FRAME 0098. ASSIGNOR(S) HEREBY CONFIRMS THE CHANGE OF ADDRESS. Recorded May 28, 2019
From: TOSHIBA SOLUTIONS CORPORATION
To: TOSHIBA SOLUTIONS CORPORATION
Reel/Frame 051297/0742 →
CHANGE OF ADDRESS Recorded Mar 8, 2019
From: TOSHIBA SOLUTIONS CORPORATION
To: TOSHIBA SOLUTIONS CORPORATION
Reel/Frame 048547/0098 →
CHANGE OF NAME Recorded Mar 8, 2019
From: TOSHIBA SOLUTIONS CORPORATION
To: TOSHIBA DIGITAL SOLUTIONS CORPORATION
Reel/Frame 048547/0215 →
ASSIGNMENT IN PART Recorded May 19, 2008
From: TOSHIBA SOLUTIONS CORPORATION
To: TOSHIBA SOLUTIONS CORPORATION; KABUSHIKI KAISHA TOSHIBA
Reel/Frame 020962/0724 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded May 31, 2005
From: ENDO, KOTARO
To: TOSHIBA SOLUTIONS CORPORATION
Reel/Frame 016609/0738 →
Priority Claims (1)
JP 2004-071494 · Mar 12, 2004 · national
Continuity (1)
Related Publication 20050204184A1 · Sep 15, 2005