EP1617331B1

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

Abstract

This record has no abstract on file.

EP1617331B1, drawing sheet 1
Sheet 1 of 8

Term

Term ended

Expired 14 June 2025, 1.3 years ago.

  1. Priority
  2. Filed
  3. Granted
  4. Expired
  5. Today

19 claims: 3 independent, 16 dependent

  1. 1
    In a fault-tolerant distributed computing system (10) comprising a first set of computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404), 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 during a view identified by a view ID, wherein operations proposed by the leader are performed by the state machine only if a quorum of computing devices in the first set agree to the proposal, a method of creating a second set of computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404), each computing device in the second set executing a replica of the state machine, each of the computing devices of the first set having an execution module for executing a first set of operations and one or more agreement modules, each having an own epoch number, for performing the following method steps while said operations of the first set of operations are executed:requesting, by some computing device in the first set, to create the second set having a second leader, by notifying the computing devices of the second set to each start up an agreement module with a new epoch number and a new view ID, wherein the epoch number is increased when a switch occurs between the first set and the second set and the view ID is increased when a change occurs between the first leader and the second leader;agreeing, among a quorum of the devices in the first set, that the computing devices comprised in the second set will execute the state machine after a given step;and 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, wherein the second set of computing devices executes a second set of operations after said given step has been reached.
  2. 5
    A computer-readable medium (130, 140, 141, 150, 152, 156, 180) including computer-executable instructions for execution on a first computing device in a fault-tolerant distributed computing system (10) comprising a plurality of computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404), a first set of the computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404) determining a first set of operations to be performed by a state machine in a first sequence of steps, the operations of the first set of operations being proposed by a leader of the first set of computing devices during a view identified by a view ID, wherein operations proposed by the leader are performed by the state machine only if a quorum of computing devices in the first set agree to the proposal, and a second set of the computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404) determining a second set of operations to be performed later by the state machine in a second sequence of steps, the operations of the second set of operations being proposed by a leader of the second set of computing devices, where the first sequence and second sequence are mutually exclusive and where a creation of the second set is requested by some computing device in the first set, by notifying the computing devices of the second set to each start up an agreement module with a new epoch number and a new view ID, the computer-executable instructions facilitating performance by:an execution module (307, 407) for executing operations in steps of the state machine, wherein the execution module is arranged to execute a sufficient number of null operations proposed by the leader of the first set in order to arrive at a given step;and a first agreement module (306, 406) having an epoch number and adapted to coordinate with the computing devices in the first set to determine the operations in the first sequence, and to agree, among a quorum of the devices in the first set, that the computing devices comprised in the second set will execute the state machine after a given step, wherein the second set of computing devices executes the second set of operations after said given step has been reached, and wherein the epoch number is increased when a switch occurs between the first set and the second set and the view ID is increased when a change occurs between the first leader and the second leader.
  3. 13
    A computing device in a first set of computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404) in a fault-tolerant distributed computing system (10) comprising a plurality of computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404), the first set determining a first set of operations to be performed by a state machine in a first sequence of steps, the operations of the first set of operations being proposed by a leader of the first set of computing devices during a view identified by a view ID, wherein operations proposed by the leader are performed by the state machine only if a quorum of computing devices in the first set agree to the proposal, and a second set of computing devices (11-15, 20, 30, 31, 100, 180, 301-304, 401-404) determining a second set of operations to be performed by the state machine in a second sequence of steps, the operations of the second set of operations being proposed by a leader of the second set of computing devices, where the first sequence and second sequence are mutually exclusive and where a creation of the second set is requested by some computing device in the first set, by notifying the computing device of the second set to each start up an agreement module with a new epoch number and a new view ID, the computing device comprising:an execution module (307, 407) for executing operations of the state machine, wherein the execution module is arranged to execute a sufficient number of null operations proposed by the leader of the first set in order to arrive at a given step;and a first agreement module (306, 406) having an epoch number, and adapted to coordinate with the other computing devices in the first set to determine operations in the first sequence, and to agree, among a quorum of the devices in the first set, that the computing devices comprised in the second set will execute the state machine after a given step, wherein the second set of computing devices executes the second set of operations after said given step has been reached, and wherein the epoch number is increased when a switch occurs between the first set and the second set and the view ID is increased when a change occurs between the first leader and the second leader.