IP Library › Granted Patent US 12,277,140
Granted Patent B2
US 12,277,140 · App. 18/372,002 · Granted Apr 15, 2025

Consensus protocol for asynchronous database transaction replication with fast, automatic failover, zero data loss, strong consistency, full SQL support and horizontal scalability

Inventors: Lik Wong (Palo Alto, CA); Leonid Novak (Castro Valley, CA); Sampanna Salunke (San Carlos, CA); Mark Dilman (Sunnyvale, CA); Wei-Ming Hu (Palo Alto, CA)
Assignee: Oracle International Corporation
G06F16/273G06F11/1469G06F16/2282G06F16/2379G06F2201/80
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 12,277,140
App. No.
18/372,002
Granted
Apr 15, 2025
Kind
B2
Abstract

A consensus protocol-based replication approach is provided. For each change operation performed by a leader server on a copy of the database, the leader server creates a replication log record and returns a result to the client. The leader does not wait for consensus for the change operation from the followers. For a commit, the leader creates a commit log record and waits for consensus. Thus, the leader executes database transactions asynchronously, performs replication of change operations asynchronously, and performs replication of transaction commits synchronously.

Claims (93)

1. A computer-implemented method comprising:

receiving, at a leader server from a client, a command to perform a change operation on a row of a table of a database, wherein:

the row of the table is replicated on a replication group of servers such that each server within the replication group of servers stores a respective copy of the row of the table,

the replication group of servers includes the leader server and one or more follower servers, and

the leader server is configured to perform data manipulation language (DML) operations on the row of the table and replicate the DML operations to the one or more follower servers;

performing, by the leader server, the change operation on the copy of the row of the table stored at the leader server;

in response to the leader server performing the change operation, creating, by the leader server, a replication log record for the change operation in a replication pipeline to be replicated to the one or more follower servers asynchronously; and

in response to the leader server creating the replication log record for the change operation in the replication pipeline, returning a result of the change operation from the leader server to the client;

subsequent to returning the result to the client, tracking, by the leader server, whether the replication log record is replicated to a replication pipeline of a consensus number of follower servers;

wherein the method is performed by one or more computing devices.

2. The method of claim 1 , wherein:

the command is associated with a particular transaction,

the method further comprises:

receiving, from the client at the leader server, a commit command to perform a commit operation on the particular transaction;

creating, by the leader server, a replication log record for the commit operation in the replication pipeline;

in response to the leader server receiving acknowledgement that the replication log record for the commit operation has been appended to a replication log of a consensus number of the one or more follower servers, performing the commit operation on the particular transaction on the copy of the row of the table at the leader server; and

returning a result of the commit operation to the client.

3. The method of claim 2 , wherein a given acknowledgement received from a given follower server includes multiple transaction commits.

4. The method of claim 2 , the method further comprising:

in response to the leader server receiving acknowledgement that the replication log record for the commit operation has been appended to the replication log of the consensus number of the one or more follower servers, advancing, by the leader server, a commit index.

5. The method of claim 2 , wherein performing the commit operation comprises writing log records for the particular transaction to disk.

6. The method of claim 1 , wherein:

the command is associated with a particular transaction,

the method further comprises:

receiving, from the client at the leader server, a commit command to perform a commit operation on the particular transaction;

preparing the particular transaction for commit at the leader server;

marking the particular transaction as in-doubt;

creating, by the leader server, a replication log record for the commit operation in the replication pipeline;

in response to the leader server failing and becoming a new follower server, determining whether there is consensus that the replication log record for the commit operation has been appended to a replication log of a consensus number of the one or more follower servers;

in response to determining there is consensus, performing the commit operation on the particular transaction on the copy of the row of the table at the new follower server; and

in response to determining there is no consensus, rolling back the particular transaction at the new follower server.

7. The method of claim 1 , wherein:

the command is associated with a particular transaction,

the method comprises:

receiving, from the client at the leader server, a commit command to perform a commit operation on the particular transaction;

creating, by the leader server, a replication log record for a pre-commit operation in the replication pipeline; and

in response to determining that there is consensus for the replication log record for the pre-commit operation, performing the commit operation on the particular transaction on the copy of the row of the table at the leader server and creating a replication log record for the commit operation in the replication pipeline.

8. The method of claim 7 , the method comprising:

in response to the leader server failing and becoming a new follower server, determining whether there is consensus for the replication log record for the pre-commit operation;

in response to determining there is consensus, performing the commit operation on the particular transaction on the copy of the row of the table at the new follower server; and

in response to determining there is no consensus, rolling back the particular transaction at the new follower server.

9. The method of claim 1 , wherein:

the command is associated with a particular transaction,

the method further comprises:

receiving, from the client at the leader server, a rollback command to perform a rollback operation on the particular transaction; and

creating, by the leader server, a replication log record for the rollback operation in the replication pipeline.

10. The method of claim 1 , wherein:

the replication log of each follower contains interleaved replication log records for uncommitted transactions, and

the replication log records in the replication log for each follower server have strictly increasing log indices.

11. The method of claim 1 , wherein each given follower server within the one or more follower servers eagerly performs the change operation on its respective copy of the row of the table and appends the replication log record to its respective replication log in parallel.

12. The method of claim 1 , wherein each given follower server within the one or more follower servers is configured to redirect commands from the client to the leader server.

13. The method of claim 1 , wherein the replication log record comprises a logical change record for the change operation, a valid log index, and a term.

14. The method of claim 1 , further comprising calling an append entries Remote Procedure Call (RPC) to propagate replication log records in the replication log to the one or more follower servers.

15. One or more non-transitory storage media storing instructions which, when executed by one or more computing devices, cause:

receiving, at a leader server from a client, a command to perform a change operation on a row of a table of a database, wherein:

the row of the table is replicated on a replication group of servers such that each server within the replication group of servers stores a respective copy of the row of the table,

the replication group of servers includes the leader server and one or more follower servers, and

the leader server is configured to perform data manipulation language (DML) operations on the row of the table and replicate the DML operations to the one or more follower servers;

performing, by the leader server, the change operation on the copy of the row of the table stored at the leader server;

in response to the leader server performing the change operation, creating, by the leader server, a replication log record for the change operation in a replication pipeline to be replicated to the one or more follower servers asynchronously; and

in response to the leader server creating the replication log record for the change operation in the replication pipeline, returning a result of the change operation from the leader server to the client

subsequent to returning the result to the client, tracking, by the leader server, whether the replication log record is replicated to a replication pipeline of a consensus number of follower servers.

16. The one or more non-transitory storage media of claim 15 , wherein:

the command is associated with a particular transaction,

the instructions further cause:

receiving, from the client at the leader server, a commit command to perform a commit operation on the particular transaction;

creating, by the leader server, a replication log record for the commit operation in the replication pipeline;

in response to the leader server receiving acknowledgement that the replication log record for the commit operation has been appended to a replication log of a consensus number of the one or more follower servers, performing the commit operation on the particular transaction on the copy of the row of the table at the leader server; and

returning a result of the commit operation to the client.

17. The one or more non-transitory storage media of claim 15 , wherein:

the command is associated with a particular transaction,

the instructions further cause:

receiving, from the client at the leader server, a commit command to perform a commit operation on the particular transaction;

preparing the particular transaction for commit at the leader server;

marking the particular transaction as in-doubt;

creating, by the leader server, a replication log record for the commit operation in the replication pipeline;

in response to the leader server failing and becoming a new follower server, determining whether there is consensus that the replication log record for the commit operation has been appended to a replication log of a consensus number of the one or more follower servers;

in response to determining there is consensus, performing the commit operation on the particular transaction on the copy of the row of the table at the new follower server; and

in response to determining there is no consensus, rolling back the particular transaction at the new follower server.

18. The one or more non-transitory storage media of claim 15 , wherein:

the command is associated with a particular transaction,

the instructions further cause:

receiving, from the client at the leader server, a commit command to perform a commit operation on the particular transaction;

creating, by the leader server, a replication log record for a pre-commit operation in the replication pipeline; and

in response to determining that there is consensus for the replication log record for the pre-commit operation, performing the commit operation on the particular transaction on the copy of the row of the table at the leader server and creating a replication log record for the commit operation in the replication pipeline.

19. The one or more non-transitory storage media of claim 15 , wherein:

the command is associated with a particular transaction,

the instructions further cause:

receiving, from the client at the leader server, a rollback command to perform a rollback operation on the particular transaction; and

creating, by the leader server, a replication log record for the rollback operation in the replication pipeline.

20. The one or more non-transitory storage media of claim 15 , wherein:

the replication log of each follower contains interleaved replication log records for uncommitted transactions, and

the replication log records in the replication log for each follower server have strictly increasing log indices.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Sep 28, 2023
From: WONG, LIK; NOVAK, LEONID; SALUNKE, SAMPANNA; DILMAN, MARK; HU, WEI-MING
To: ORACLE INTERNATIONAL CORPORATION
Reel/Frame 065068/0041 →
Continuity (2)
Provisional Application 63415466 · Oct 12, 2022
Related Publication 20240126781A1 · Apr 18, 2024
References Cited (18)
US 20050165858A1 · Tom et al. · 2005 [cited by applicant]
US 20130290249A1 · Merriman et al. · 2013 [cited by applicant]
US 20160371358A1 · Lee · 2016 [cited by examiner]
US 20170103092A1 · Hu et al. · 2017 [cited by applicant]
US 20190155705A1 · Chavan · 2019 [cited by examiner]
US 20190163545A1 · Singh et al. · 2019 [cited by applicant]
US 20190325055A1 · Lee et al. · 2019 [cited by applicant]
US 20200034257A1 · Mahmood et al. · 2020 [cited by applicant]
US 20200364239A1 · Kumar · 2020 [cited by examiner]
US 20220100710A1 · Camargos et al. · 2022 [cited by applicant]
US 20220114058A1 · Mylavarapu et al. · 2022 [cited by applicant]
US 20220114164A1 · Krishnaswamy et al. · 2022 [cited by applicant]
US 20240045887A1 · VanBenschoten · 2024 [cited by examiner]
EP 1876788A1 · 2008 [cited by applicant]
EP 3182300A1 · 2017 [cited by applicant]
WO WO2021021757A1 · 2021 [cited by applicant]
Ongaro et al., “In Search of an Understandable Consensus Algorithm”, Jun. 19, 2014, pp. 1-16. [cited by applicant]
Cao et al., “PolarDB-X: An Elastic Distributed Relational Database for Cloud-Native Applications”, 2022 IEEE 38th ICDE, pp. 1-14. [cited by applicant]