IP Library Granted Patent US 12,450,259
Granted Patent B2
US 12,450,259 · App. 18/805,383 · Granted Oct 21, 2025

Performance optimization in raft-based asynchronous database transaction replication

Inventors: Lik Wong (Palo Alto, CA); Sampanna Salunke (Dublin, CA); Leonid Novak (Castro Valley, CA); Mark Dilman (Sunnyvale, CA); Wei-Ming Hu (Palo Alto, CA)
Assignee: Oracle International Corporation
G06F16/273
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,450,259
App. No.
18/805,383
Granted
Oct 21, 2025
Kind
B2
Abstract

Replication is improved in a globally distributed database, such as a replicated sharded database, which uses raft-based asynchronous database replication. Improvements include Raft log persistence, coordination of followers' processing speed, transaction outcome determination, and column name compression, and improved failover time through heartbeat consolidation and keeping apply processes of followers running across failovers.

Claims (73)

1. A computer-implemented method comprising:

within a replication group that replicates a replication unit of a database, a leader for said replication unit sending via an in-memory queue a stream of LCRs (Logical Change Record) to a plurality of followers for said replication unit, wherein LCRs sent in said stream reflect transactions executed by said leader that change database data in said replication unit;

said leader storing LCRs in said stream in a persistent Raft log;

determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching said lagging follower as a subscriber of said in-memory queue; and

in response to determining that said lagging follower meets said one or more criteria for detaching said lagging follower, detaching said lagging follower from said in-memory queue, wherein determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching comprises at least one of:

determining whether a number of LCRs in said in-memory queue, which are not yet read by said lagging follower but have been read by one or more other followers of said plurality of followers, is greater than a threshold number, or

determining whether the lagging follower has remained behind reading LCRs from said in-memory queue for at least a threshold period of time.

2. The method of claim 1 , further including said lagging follower reading LCRs from said persistent Raft log in response to detaching said lagging follower.

3. The method of claim 2 , wherein determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching includes determining there are sufficient followers for consensus without said lagging follower.

4. The method of claim 1 , further including:

said leader storing in a commit queue records indicating a number of pending commits; sending an indication of said number of pending commits to a particular follower of said plurality of followers; and

said particular follower flushing to persistent storage a batch of LCRs from another in memory queue for said particular follower based on said indication of said number of pending commits.

5. The method of claim 1 , further including:

in response to a particular follower of said plurality of followers determining that the particular follower has not received a post-commit LCR for a database transaction after a threshold period of time, said particular follower invoking a remote procedure call to obtain an outcome for said database transaction; and

said particular follower committing or rolling back a corresponding apply database transaction based on the outcome.

6. The method of claim 1 , wherein:

wherein a particular LCR in said stream specifies a change to a plurality of columns in a database table; and

wherein for each column of said plurality of columns, said particular LCR includes a hash value in lieu of a column name of said each column, said hash value being calculated according to a hash algorithm.

7. The method of claim 1 , wherein:

said leader runs on a first shard server that hosts a plurality of leaders that includes said leader, and said followers of said plurality of leaders run a respective shard server of a plurality of shard servers;

each leader of said plurality of leaders stores respective liveness data in a shared memory area of said shard server; and

a single process on said first shard server sending to each shard server of said plurality of shard servers a consolidated heartbeat message that consolidates liveness data stored by each leader of said plurality of leaders in said shared memory area.

8. The method of claim 1 , wherein during a failover for the replication group in which a follower of said plurality of followers becomes a new leader, on said lagging follower:

a network receiver receiving a stream of LCRs from the new leader;

said network receiver storing the stream of LCRs to a local persistently stored local Raft log of said lagging follower; and

in lieu of receiving said stream of LCRs from said network receiver, an apply process of said lagging follower reading a stream of LCRs from said local Raft log to apply.

9. The method of claim 8 , after completing said failover, said network receiver forwarding LCRs to said apply process.

10. The method of claim 1 , wherein in response said leader determining that a file of said Raft log selected for writing a set of LCRs includes a minimum required log index for said plurality of followers, said leader waiting to write said set of LCRs at least until said file no longer includes said minimum required log index.

11. One or more non-transitory storage media storing sequences of instructions that, when executed by one or more computing devices, cause:

within a replication group that replicates a replication unit of a database, a leader for said replication unit sending via an in-memory queue a stream of LCRs (Logical Change Record) to a plurality of followers for said replication unit, wherein LCRs sent in said stream reflect transactions executed by said leader that change database data in said replication unit;

said leader storing LCRs in said stream in a persistent Raft log;

determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching said lagging follower as a subscriber of said in-memory queue; and

in response to determining that said lagging follower meets said one or more criteria for detaching said lagging follower, detaching said lagging follower from said in-memory queue, wherein determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching comprises at least one of:

determining whether a number of LCRs in said in-memory queue, which are not yet read by said lagging follower but have been read by one or more other followers of said plurality of followers, is greater than a threshold number, or

determining whether the lagging follower has remained behind reading LCRs from said in-memory queue for at least a threshold period of time.

12. The one or more non-transitory storage media of claim 11 , wherein the one or more sequences of instructions include instructions, that when executed by one or more computing devices, cause said lagging follower reading LCRs from said persistent Raft log in response to detaching said lagging follower.

13. The one or more non-transitory storage media of claim 12 , wherein determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching includes determining there are sufficient followers for consensus without said lagging follower.

14. The one or more non-transitory storage media of claim 11 , wherein the one or more sequences of instructions include instructions, that when executed by one or more computing devices, cause:

said leader storing in a commit queue records indicating a number of pending commits; sending an indication of said number of pending commits to a particular follower of said plurality of followers; and

said particular follower flushing to persistent storage a batch of LCRs from another in-memory queue for said particular follower based on said indication of said number of pending commits.

15. The one or more non-transitory storage media of claim 11 , wherein the one or more sequences of instructions include instructions, that when executed by one or more computing devices, cause:

in response to a particular follower of said plurality of followers determining that the particular follower has not received a post-commit LCR for a database transaction after a threshold period of time, said particular follower invoking a remote procedure call to obtain an outcome for said database transaction; and

said particular follower committing or rolling back a corresponding apply database transaction based on the outcome.

16. The one or more non-transitory storage media of claim 11 , wherein the one or more sequences of instructions include instructions, that when executed by one or more computing devices, cause:

wherein a particular LCR in said stream specifies a change to a plurality of columns in a database table; and

wherein for each column of said plurality of columns, said particular LCR includes a hash value in lieu of a column name of said each column, said hash value being calculated according to a hash algorithm.

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

said leader runs on a first shard server that hosts a plurality of leaders that includes said leader, and said followers of said plurality of leaders run a respective shard server of a plurality of shard servers;

each leader of said plurality of leaders stores respective liveness data in a shared memory area of said shard server; and

a single process on said first shard server sending to each shard server of said plurality of shard servers a consolidated heartbeat message that consolidates liveness data stored by each leader of said plurality of leaders in said shared memory area.

18. The one or more non-transitory storage media of claim 11 , wherein during a failover for the replication group in which a follower of said plurality of followers becomes a new leader, on said lagging follower:

a network receiver receiving a stream of LCRs from the new leader;

said network receiver storing the stream of LCRs to a local persistently stored local Raft log of said lagging follower; and

in lieu of receiving said stream of LCRs from said network receiver, an apply process of said lagging follower reading a stream of LCRs from said local Raft log to apply.

19. The one or more non-transitory storage media of claim 18 , wherein the one or more sequences of instructions include instructions, that when executed by one or more computing devices, cause after completing said failover, said network receiver forwarding LCRs to said apply process.

20. The one or more non-transitory storage media of claim 11 , wherein in response said leader determining that a file of said Raft log selected for writing a set of LCRs includes a minimum required log index for said plurality of followers, said leader waiting to write said set of LCRs at least until said file no longer includes said minimum required log index.

21. A data processing system, comprising:

a database implemented on one or more computing devices;

a replication group that replicates a replication unit of the database;

a leader for said replication unit sending, via an in-memory queue, a stream of LCRs (Logical Change Record) to a plurality of followers for said replication unit, wherein:

LCRs sent in said stream reflect transactions executed by said leader that change database data in said replication unit,

said leader stores LCRs in said stream in a persistent Raft log, and

the computing system comprises a processor and a memory coupled to the processor, wherein the memory comprises instructions that are executed by the processor to cause the processor to execute the process comprising:

determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching said lagging follower as a subscriber of said in-memory queue; and

in response to determining that said lagging follower meets said one or more criteria for detaching said lagging follower, detaching said lagging follower from said in-memory queue, wherein determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching comprises at least one of:

determining whether a number of LCRs in said in-memory queue, which are not yet read by said lagging follower but have been read by one or more other followers of said plurality of followers, is greater than a threshold number, or

determining whether the lagging follower has remained behind reading LCRs from said in-memory queue for at least a threshold period of time.

22. The data processing system of claim 21 , further including said lagging follower reading LCRs from said persistent Raft log in response to detaching said lagging follower.

23. The data processing system of claim 22 , wherein determining whether a lagging follower of said plurality of followers meets one or more criteria for detaching includes determining there are sufficient followers for consensus without said lagging follower.

24. The data processing system of claim 1 , further including:

said leader storing in a commit queue records indicating a number of pending commits;

sending an indication of said number of pending commits to a particular follower of said plurality of followers; and

said particular follower flushing to persistent storage a batch of LCRs from another in memory queue for said particular follower based on said indication of said number of pending commits.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Aug 14, 2024
From: WONG, LIK; SALUNKE, SAMPANNA; NOVAK, LEONID; DILMAN, MARK; HU, WEI-MING
To: ORACLE INTERNATIONAL CORPORATION
Reel/Frame 068289/0210 →
Continuity (2)
Provisional Application 63582681 · Sep 14, 2023
Related Publication 20250094445A1 · Mar 20, 2025
References Cited (28)
US 10936573B2 · Ducott · 2021 [cited by applicant]
US 20040243558A1 · Nelson · 2004 [cited by applicant]
US 20050165858A1 · Tom et al. · 2005 [cited by applicant]
US 20100145909A1 · Ngo · 2010 [cited by applicant]
US 20130006933A1 · Holden · 2013 [cited by applicant]
US 20130290249A1 · Merriman et al. · 2013 [cited by applicant]
US 20140337529A1 · Antony · 2014 [cited by applicant]
US 20160371358A1 · Lee · 2016 [cited by applicant]
US 20170103092A1 · Hu et al. · 2017 [cited by applicant]
US 20180260125A1 · Botes · 2018 [cited by examiner]
US 20190155705A1 · Chavan · 2019 [cited by applicant]
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 applicant]
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 applicant]
EP 1876788A1 · 2008 [cited by applicant]
EP 3182300A1 · 2017 [cited by applicant]
WO WO2021021757A1 · 2021 [cited by applicant]
“In Search of an Understandable Consensus Algorithm”; By: Diego Ongaro, Published 2014 https://www.usenix.org/system/files/conference/atc14/atc14-paper-ongaro.pdf (Year: 2014). [cited by examiner]
Ongaro et al., “In Search of an Understandable Consensus Algorithm”, Jun. 19, 2014, pp. 1-16. [cited by applicant]
International Searching Authority, “International Search Report and Written Opinion”, in International Application No. PCT/US2023/034464 dated Jan. 3, 2024, p. 14. [cited by applicant]
International Searching Authority, “International Search Report and Written Opinion”, in International Application No. PCT/US 2023/034465 dated Jan. 9, 2024, pp. 13. [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]
Ongaro, “In Search of an Understandable Consensus Algorithm”, Published 2014 https://www.usenix.org/system/files/conference/atc14/atc14-paper-ongaro.pdf. [cited by applicant]