IP Library Granted Patent US 11,106,489
Granted Patent B2
US 11,106,489 · App. 17/068,576 · Granted Aug 31, 2021

Efficient, time-based leader node election in a distributed computing system

Inventors: Zhenkun Yang (Hangzhou, CN); Jinliang Xiao (Hangzhou, CN)
Assignee: ANT FINANCIAL (HANG ZHOU) NETWORK TECHNOLOGY CO., LTD.
G06F9/4806G06F11/1425G06F11/187G06F16/27
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 11,106,489
App. No.
17/068,576
Granted
Aug 31, 2021
Kind
B2
Abstract

A same voting time, a same vote counting time, and a same leader node tenure is configured by a host for all nodes. Time configuration information including the same configured voting time, the same vote counting time, and the same leader node tenure, is sent to all the nodes. The nodes are operable to vote during the same voting time, count the number of votes during the same vote counting time, and elect a leader node according to a vote counting result. The nodes are enabled to perform periodic node election according to the same leader node tenure.

Claims (91)

1. A computer-implemented method, comprising:

sending, by a first node, a first voting request to a first plurality of nodes during a first voting time, wherein the first voting time is comprised in time configuration information, wherein the first node and the first plurality of nodes are processing devices in a distributed system;

receiving, by the first node and from one or more node of the first plurality of nodes, first vote information generated based on the first voting request;

counting, by the first node and in a first vote counting time, a first number of votes based on the first vote information, wherein the first vote counting time is comprised in the time configuration information; and

determining, based on the first number of votes, a leader node to process a first specified service.

2. The computer-implemented method of claim 1 , wherein the determining comprises:

in response to determining that the first number of votes exceeds a preset threshold, determining the first node as the leader node, wherein the preset threshold is one-half of a number of a second plurality of nodes, wherein the second plurality of nodes comprise the first node and the first plurality of nodes.

3. The computer-implemented method of claim 2 , further comprising:

determining, by the first node, node information comprised in the first voting request;

determining, based on the node information, a priority of the second plurality of nodes;

selecting a particular node of the second plurality of nodes having a maximum priority;

generating vote information for the particular node; and

sending the vote information to the particular node.

4. The computer-implemented method of claim 1 , further comprising:

determining, by the leader node, that the leader node itself meets a preset criteria, sending, by the leader node, a reappointment request to the first plurality of nodes during a second voting time, wherein the second voting time is determined based on a leader node tenure comprised in the time configuration information;

receiving, by the leader node, second vote information generated by one or more node of the first plurality of nodes according to the reappointment request;

counting, by the leader node and in a second vote counting time, a second number of votes based on the second vote information;

determining, by the leader node and based on the second number of votes, a reappointment state; and

processing a second specified service after reappointment succeeds, wherein the preset criteria comprises at least one of a load processing criterion or a data updating degree.

5. The computer-implemented method of claim 1 , further comprising:

determining, by the leader node, that the leader node itself does not meet a preset criteria;

determining, by the leader node and based on a leader node tenure in the time configuration information, a particular node of the first plurality of nodes that meets the preset criteria to serve as a successor node; and

sending, by the leader node, a second voting request for the successor node to the first plurality of nodes during a second voting time, wherein the second voting time is determined based on the time configuration information, wherein the successor node receives second vote information from one or more node of the first plurality of nodes, wherein the successor node counts, based on the second vote information, a second number of votes in a second vote counting time, wherein the successor node is determined as a new leader node.

6. The computer-implemented method of claim 1 , further comprising:

receiving, by the first node and from the leader node during a second voting time, a reappointment request, wherein the second voting time is determined based on the time configuration information;

generating, by the first node and based on the reappointment request, second vote information; and

sending the second vote information to the leader node in the second voting time, wherein the leader node counts, based on the second vote information and in a second vote counting time, a second number of votes to determine a reappointment state.

7. The computer-implemented method of claim 1 , further comprising:

receiving, by the first node and from the leader node during a second voting time, a second voting request for a successor node, wherein the second voting time is determined based on the time configuration information;

generating, by the first node and based on the second voting request, second vote information; and

sending the second vote information to the successor node during the second voting time, wherein the successor node counts, based on the second vote information, a second number of votes during a second vote counting time to determine a new leader node.

8. A non-transitory, computer-readable medium storing one or more instructions executable by a computer system to perform operations comprising:

sending, by a first node, a first voting request to a first plurality of nodes during a first voting time, wherein the first voting time is comprised in time configuration information, wherein the first node and the first plurality of nodes are processing devices in a distributed system;

receiving, by the first node and from one or more node of the first plurality of nodes, first vote information generated based on the first voting request;

counting, by the first node and in a first vote counting time, a first number of votes based on the first vote information, wherein the first vote counting time is comprised in the time configuration information; and

determining, based on the first number of votes, a leader node to process a first specified service.

9. The non-transitory, computer-readable medium of claim 8 , wherein the determining comprises:

in response to determining that the first number of votes exceeds a preset threshold, determining the first node as the leader node, wherein the preset threshold is one-half of a number of a second plurality of nodes, wherein the second plurality of nodes comprise the first node and the first plurality of nodes.

10. The non-transitory, computer-readable medium of claim 9 , the operations further comprising:

determining, by the first node, node information comprised in the first voting request;

determining, based on the node information, a priority of the second plurality of nodes;

selecting a particular node of the second plurality of nodes having a maximum priority;

generating vote information for the particular node; and

sending the vote information to the particular node.

11. The non-transitory, computer-readable medium of claim 8 , the operations further comprising:

determining, by the leader node, that the leader node itself meets a preset criteria, sending, by the leader node, a reappointment request to the first plurality of nodes during a second voting time, wherein the second voting time is determined based on a leader node tenure comprised in the time configuration information;

receiving, by the leader node, second vote information generated by one or more node of the first plurality of nodes according to the reappointment request;

counting, by the leader node and in a second vote counting time, a second number of votes based on the second vote information;

determining, by the leader node and based on the second number of votes, a reappointment state; and

processing a second specified service after reappointment succeeds, wherein the preset criteria comprises at least one of a load processing criterion or a data updating degree.

12. The non-transitory, computer-readable medium of claim 8 , the operations further comprising:

determining, by the leader node, that the leader node itself does not meet a preset criteria;

determining, by the leader node and based on a leader node tenure in the time configuration information, a particular node of the first plurality of nodes that meets the preset criteria to serve as a successor node; and

sending, by the leader node, a second voting request for the successor node to the first plurality of nodes during a second voting time, wherein the second voting time is determined based on the time configuration information, wherein the successor node receives second vote information from one or more node of the first plurality of nodes, wherein the successor node counts, based on the second vote information, a second number of votes in a second vote counting time, wherein the successor node is determined as a new leader node.

13. The non-transitory, computer-readable medium of claim 8 , the operations further comprising:

receiving, by the first node and from the leader node during a second voting time, a reappointment request, wherein the second voting time is determined based on the time configuration information;

generating, by the first node and based on the reappointment request, second vote information; and

sending the second vote information to the leader node in the second voting time, wherein the leader node counts, based on the second vote information and in a second vote counting time, a second number of votes to determine a reappointment state.

14. The non-transitory, computer-readable medium of claim 8 , the operations further comprising:

receiving, by the first node and from the leader node during a second voting time, a second voting request for a successor node, wherein the second voting time is determined based on the time configuration information;

generating, by the first node and based on the second voting request, second vote information; and

sending the second vote information to the successor node during the second voting time, wherein the successor node counts, based on the second vote information, a second number of votes during a second vote counting time to determine a new leader node.

15. A computer-implemented system, comprising:

one or more computers; and

one or more computer memory devices interoperably coupled with the one or more computers and having tangible, non-transitory, machine-readable media storing one or more instructions that, when executed by the one or more computers, perform one or more operations comprising:

sending, by a first node, a first voting request to a first plurality of nodes during a first voting time, wherein the first voting time is comprised in time configuration information, wherein the first node and the first plurality of nodes are processing devices in a distributed system;

receiving, by the first node and from one or more node of the first plurality of nodes, first vote information generated based on the first voting request;

counting, by the first node and in a first vote counting time, a first number of votes based on the first vote information, wherein the first vote counting time is comprised in the time configuration information; and

determining, based on the first number of votes, a leader node to process a first specified service.

16. The computer-implemented system of claim 15 , wherein the determining comprises:

in response to determining that the first number of votes exceeds a preset threshold, determining the first node as the leader node, wherein the preset threshold is one-half of a number of a second plurality of nodes, wherein the second plurality of nodes comprise the first node and the first plurality of nodes.

17. The computer-implemented system of claim 16 , the operations further comprising:

determining, by the first node, node information comprised in the first voting request;

determining, based on the node information, a priority of the second plurality of nodes;

selecting a particular node of the second plurality of nodes having a maximum priority;

generating vote information for the particular node; and

sending the vote information to the particular node.

18. The computer-implemented system of claim 15 , the operations further comprising:

determining, by the leader node, that the leader node itself meets a preset criteria, sending, by the leader node, a reappointment request to the first plurality of nodes during a second voting time, wherein the second voting time is determined based on a leader node tenure comprised in the time configuration information;

receiving, by the leader node, second vote information generated by one or more node of the first plurality of nodes according to the reappointment request;

counting, by the leader node and in a second vote counting time, a second number of votes based on the second vote information;

determining, by the leader node and based on the second number of votes, a reappointment state; and

processing a second specified service after reappointment succeeds, wherein the preset criteria comprises at least one of a load processing criterion or a data updating degree.

19. The computer-implemented system of claim 15 , the operations further comprising:

determining, by the leader node, that the leader node itself does not meet a preset criteria;

determining, by the leader node and based on a leader node tenure in the time configuration information, a particular node of the first plurality of nodes that meets the preset criteria to serve as a successor node; and

sending, by the leader node, a second voting request for the successor node to the first plurality of nodes during a second voting time, wherein the second voting time is determined based on the time configuration information, wherein the successor node receives second vote information from one or more node of the first plurality of nodes, wherein the successor node counts, based on the second vote information, a second number of votes in a second vote counting time, wherein the successor node is determined as a new leader node.

20. The computer-implemented system of claim 15 , the operations further comprising:

receiving, by the first node and from the leader node during a second voting time, a reappointment request, wherein the second voting time is determined based on the time configuration information;

generating, by the first node and based on the reappointment request, second vote information; and

sending the second vote information to the leader node in the second voting time, wherein the leader node counts, based on the second vote information and in a second vote counting time, a second number of votes to determine a reappointment state.

Assignments (5)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 27, 2021
From: ANT FINANCIAL (HANG ZHOU) NETWORK TECHNOLOGY CO., LTD.
To: BEIJING OCEANBASE TECHNOLOGY CO., LTD.
Reel/Frame 057349/0070 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Feb 11, 2021
From: ADVANCED NEW TECHNOLOGIES CO., LTD.
To: ANT FINANCIAL (HANG ZHOU) NETWORK TECHNOLOGY CO., LTD.
Reel/Frame 055237/0137 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 2, 2020
From: YANG, ZHENKUN; XIAO, JINLIANG
To: ALIBABA GROUP HOLDING LIMITED
Reel/Frame 054238/0705 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 2, 2020
From: ALIBABA GROUP HOLDING LIMITED
To: ADVANTAGEOUS NEW TECHNOLOGIES CO., LTD.
Reel/Frame 054275/0388 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 2, 2020
From: ADVANTAGEOUS NEW TECHNOLOGIES CO., LTD.
To: ADVANCED NEW TECHNOLOGIES CO., LTD.
Reel/Frame 054275/0524 →