EP1617331A2

Efficient changing of replica sets in distributed fault-tolerant computing system

Abstract

A distributed computing system can be operated in a fault tolerant manner using a set of computing devices. A set of computing devices can tolerate a number of failures by implementing identical replicas of a state machine and selecting proposals. The set of computing devices participating in the distributed computing system by hosting replicas can be modified by adding or removing a computing device from the set, or by specifying particular computing devices for participation. Changing the participating computing devices in the set increases fault tolerance by replacing defective devices with operational devices, or by increasing the amount of redundancy in the system.

EP1617331A2, drawing sheet 1
Sheet 1 of 9

Term

Term ended

Projected expiry passed 14 June 2025, 1.3 years ago.

  1. Priority
  2. Filed
  3. Published
  4. Projected expiry
  5. Today

25 claims: 5 independent, 20 dependent

  1. 1
    In a fault-tolerant distributed computing system comprising a first set of computing devices, the computing devices in the first set determining operations to be performed by a state machine as a sequence of steps, a method of allowing a second set of computing devices to determine operations in the sequence of steps, the method comprising:specifying, by some computing device in the first set, the second set of computing devices;designating one step to be the last step for which the first set determines an operation to be performed by the state machine;and determining, by the computing devices in the second set, operations to be performed by the state machine in steps subsequent to the designated step.
  2. 3
    The method of claim I wherein the method is initiated automatically in response to a policy.
  3. 5
    In a fault-tolerant distributed computing system comprising a first set of computing devices, each computing device in the first set executing a replica of a state machine, operations to be performed on the state machine being proposed by a leader of the first set, a method of creating a second set of computing devices, each computing device in the second set executing a replica of the state machine, the method comprising:requesting, by some computing device in the first set, to create the second set;agreeing, among a quorum of the devices in the first set, that the computing devices comprising the second set will execute the state machine after a given step;proposing, by the leader of the first set, a sufficient number of null operations for the state machine in order to arrive at the given step.
  4. 10
    In a fault-tolerant distributed computing system comprising a plurality of computing devices, a first set of the computing devices determining operations to be performed by a state machine in a first sequence of steps, and a second set of the computing devices determining operations to be performed later by the state machine in a second sequence of steps, where the first sequence and second sequence are mutually exclusive, a computer-readable medium including computer-executable instructions for execution on a first computing device, the computer-executable instructions facilitating performance by:an execution module for executing operations in steps of the state machine;and a first agreement module for coordinating with other computing devices in the first set to determine operations in the first sequence.
  5. 18
    In a fault-tolerant distributed computing system comprising a plurality of computing devices, a first set of the computing devices determining operations to be performed by a state machine in a first sequence of steps, and a second set of the computing devices determining operations to be performed by the state machine in a second sequence of steps, where the first sequence and second sequence are mutually exclusive and where the first sequence of steps precedes the second sequence of steps, a computing device comprising:an execution module for executing operations of the state machine;and a first agreement module for coordinating with other computing devices in the first set to determine operations in the first sequence.