IP Library Granted Patent US 12,461,937
Granted Patent B2
US 12,461,937 · App. 18/336,932 · Granted Nov 4, 2025

Change data capture state tracking for multi-region multi-master noSQL database

Inventors: Zijie Li (New York, NY); Shitanshu Verma (Livingston, NJ); Can Tang (Brooklyn, NY); Gary Elliott (Larchmont, NY); Gregory Allen Morris (Hanover, NH); Thomas Robert Magrino (Mamaroneck, NY); Jack Timothy Dingilian (Brooklyn, NY); Teng Zhong (Mountain View, CA); Andrii Shyshkalov (Munich, DE); Siu Man Yau (Plainview, NY); Yijie Bu (Sunnyvale, CA)
Assignee: Google LLC
G06F16/273G06F16/2365G06F16/256G06F16/285
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,461,937
App. No.
18/336,932
Granted
Nov 4, 2025
Kind
B2
Abstract

A method for change data capture state tracking includes accessing a distributed database including a plurality of clusters, each cluster of the plurality of clusters including a respective plurality of partitions, each cluster of the plurality of clusters configured to receive read and write operation. The method includes receiving, at a second cluster, a plurality of changes for a second table and storing the plurality of changes at a replication log. The method also includes asynchronously replicating the plurality of changes from the second table to a first table and maintaining a respective change stream position tracking a respective position in the replication log indicating one or more changes of the plurality of changes that have been replicated. The method includes receiving a read request at the first cluster requesting one or more rows of the first table and returning the respective change stream position.

Claims (52)

1 . A computer-implemented method executed by data processing hardware that causes the data processing hardware to perform operations comprising:

accessing a distributed database comprising a plurality of clusters, each cluster of the plurality of clusters comprising a respective plurality of partitions, each cluster of the plurality of clusters configured to receive read and write operations, wherein:

the respective plurality of partitions of a first cluster of the plurality of clusters stores a first table comprising a first respective plurality of rows, each respective partition of the respective plurality of partitions of the first cluster comprising a respective portion of the first respective plurality of rows; and

the respective plurality of partitions of a second cluster of the plurality of clusters stores a second table comprising a second respective plurality of rows, each respective partition of the respective plurality of partitions of the second cluster comprising a respective portion of the second respective plurality of rows;

receiving, at the second cluster, a plurality of changes for the second table;

storing the plurality of changes at a replication log;

asynchronously replicating the plurality of changes from the second table to the first table;

while asynchronously replicating the plurality of changes from the second table to the first table, for each row of the first respective plurality of rows of each respective partition of the respective plurality of partitions of the first cluster, maintaining a respective change stream position, the respective change stream position tracking a respective position in the replication log indicating one or more changes of the plurality of changes that have been replicated;

receiving a read request at the first cluster requesting one or more rows of the first respective plurality of rows, the one or more rows associated with a first subset of partitions of the respective plurality of partitions of the first cluster; and

in response to the read request, returning, for each respective partition of the first subset of partitions, the respective change stream position.

2 . The method of claim 1 , wherein returning the respective change stream position comprises returning an encrypted blob comprising the respective change stream position.

3 . The method of claim 1 , wherein the respective plurality of partitions of the first cluster is different from the respective plurality of partitions of the second cluster.

4 . The method of claim 3 , wherein each respective change of the plurality of changes for the second table is associated with a respective partition of the respective plurality of partitions of the second cluster.

5 . The method of claim 4 , wherein the operations further comprise, for each respective change of the plurality of changes for the second table:

determining a different partition of the respective plurality of partitions of the first cluster that aligns with the respective partition of the respective plurality of partitions of the second cluster associated with the respective change; and

adding the respective change to a respective partition replication log associated with the respective partition of the respective plurality of partitions of the first cluster.

6 . The method of claim 4 , wherein the operations further comprise, for each respective change of the plurality of changes for the second table, determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is active.

7 . The method of claim 4 , wherein the operations further comprise, for each respective change of the plurality of changes for the second table, determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is inactive.

8 . The method of claim 7 , wherein the operations further comprise, in response to determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is inactive, determining that a respective incarnation associated with the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is closed.

9 . The method of claim 4 , wherein the operations further comprise, for a respective change of the plurality of changes for the second table:

determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is new; and

in response to determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is new, determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change belongs to a new incarnation.

10 . The method of claim 1 , wherein:

the first cluster comprises a plurality of nodes; and

each partition of the respective plurality of partitions of the first cluster is stored on a respective node of the plurality of nodes.

11 . A system comprising:

data processing hardware; and

memory hardware in communication with the data processing hardware, the memory hardware storing instructions that when executed on the data processing hardware cause the data processing hardware to perform operations comprising:

accessing a distributed database comprising a plurality of clusters, each cluster of the plurality of clusters comprising a respective plurality of partitions, each cluster of the plurality of clusters configured to receive read and write operations, wherein:

the respective plurality of partitions of a first cluster of the plurality of clusters stores a first table comprising a first respective plurality of rows, each respective partition of the respective plurality of partitions of the first cluster comprising a respective portion of the first respective plurality of rows; and

the respective plurality of partitions of a second cluster of the plurality of clusters stores a second table comprising a second respective plurality of rows, each respective partition of the respective plurality of partitions of the second cluster comprising a respective portion of the second respective plurality of rows;

receiving, at the second cluster, a plurality of changes for the second table;

storing the plurality of changes at a replication log;

asynchronously replicating the plurality of changes from the second table to the first table;

while asynchronously replicating the plurality of changes from the second table to the first table, for each row of the first respective plurality of rows of each respective partition of the respective plurality of partitions of the first cluster, maintaining a respective change stream position, the respective change stream position tracking a respective position in the replication log indicating one or more changes of the plurality of changes that have been replicated;

receiving a read request at the first cluster requesting one or more rows of the first respective plurality of rows, the one or more rows associated with a first subset of partitions of the respective plurality of partitions of the first cluster; and

in response to the read request, returning, for each respective partition of the first subset of partitions, the respective change stream position.

12 . The system of claim 11 , wherein returning the respective change stream position comprises returning an encrypted blob comprising the respective change stream position.

13 . The system of claim 11 , wherein the respective plurality of partitions of the first cluster is different from the respective plurality of partitions of the second cluster.

14 . The system of claim 13 , wherein each respective change of the plurality of changes for the second table is associated with a respective partition of the respective plurality of partitions of the second cluster.

15 . The system of claim 14 , wherein the operations further comprise, for each respective change of the plurality of changes for the second table:

determining a different partition of the respective plurality of partitions of the first cluster that aligns with the respective partition of the respective plurality of partitions of the second cluster associated with the respective change; and

adding the respective change to a respective partition replication log associated with the respective partition of the respective plurality of partitions of the first cluster.

16 . The system of claim 14 , wherein the operations further comprise, for each respective change of the plurality of changes for the second table, determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is active.

17 . The system of claim 14 , wherein the operations further comprise, for each change of the plurality of changes for the second table, determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is inactive.

18 . The system of claim 17 , wherein the operations further comprise, in response to determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is inactive, determining that a respective incarnation associated with the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is closed.

19 . The system of claim 14 , wherein the operations further comprise, for a respective change of the plurality of changes for the second table:

determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is new; and

in response to determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change is new, determining that the respective partition of the respective plurality of partitions of the second cluster associated with the respective change belongs to a new incarnation.

20 . The system of claim 11 , wherein:

the first cluster comprises a plurality of nodes; and

each partition of the respective plurality of partitions of the first cluster is stored on a respective node of the plurality of nodes.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jun 22, 2023
From: LI, ZIJIE; VERMA, SHITANSHU; TANG, CAN; ELLIOTT, GARY; MORRIS, GREGORY ALLEN; MAGRINO, THOMAS ROBERT; DINGILIAN, JACK TIMOTHY; ZHONG, TENG; SHYSHKALOV, ANDRII; YAU, SIU MAN; BU, YIJIE
To: GOOGLE LLC
Reel/Frame 064025/0287 →
Continuity (1)
Related Publication 20240419686A1 · Dec 19, 2024
References Cited (18)
US 10068002B1 · Wilczynski · 2018 [cited by examiner]
US 12182105B1 · Holenstein · 2024 [cited by examiner]
US 20040230619A1 · Blanco · 2004 [cited by examiner]
US 20060155945A1 · McGarvey · 2006 [cited by examiner]
US 20060168120A1 · Parham · 2006 [cited by examiner]
US 20180268044A1 · Barber · 2018 [cited by examiner]
US 20190102418A1 · Vasudevan et al. · 2019 [cited by applicant]
US 20190391957A1 · Ye · 2019 [cited by applicant]
US 20200301947A1 · Botev et al. · 2020 [cited by applicant]
US 20200409566A1 · Demoor et al. · 2020 [cited by applicant]
US 20220188196A1 · Vig et al. · 2022 [cited by applicant]
US 20220207036A1 · Ou · 2022 [cited by examiner]
US 20240028580A1 · Zhang · 2024 [cited by examiner]
House, Daniel, et al. “Toward fast and reliable active-active geo-replication for a distributed data caching service in the mobile cloud.” Procedia Computer Science 191 (2021) (Year: 2021). [cited by examiner]
Bhaskaran, S. V. “Resilient real-time data delivery for ai summarization in conversational platforms: Ensuring low latency, high availability, and disaster recovery.” Journal of Intelligent Connectivity and Emerging Tec… [cited by examiner]
Zhang, Irene, et al. “Building consistent transactions with inconsistent replication.” ACM Transactions on Computer Systems (TOCS) 35.4 (2018): 1-37. (Year: 2018). [cited by examiner]
Zhao, Weibin, and Henning G. Schulzrinne. “A Flexible and Efficient Protocol for Multi-Scope Service Registry Replication.” (2002). (Year: 2002). [cited by examiner]
International Search Report and Written Opinion issued in related PCT Application No. PCT/US2024/033724, dated Sep. 18, 2024. [cited by applicant]