IP Library Granted Patent US 12,386,856
Granted Patent B2
US 12,386,856 · App. 18/418,257 · Granted Aug 12, 2025

Distributed database configuration

Inventors: Alexander Shraer (Stanford, CA); Artyom Sharov (Haifa, IL); Arif Abdulhusein Merchant (Los Altos, CA); Brian F. Cooper (San Jose, CA)
Assignee: Google LLC
G06F16/27G06F16/25G06F16/955H04L67/1023H04L67/1025H04L67/1097H04L67/52
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,386,856
App. No.
18/418,257
Granted
Aug 12, 2025
Kind
B2
Abstract

Replicas are selected in a large distributed network, and the roles for these replicas are identified. In one example, a leader is selected from among candidate computing clusters. To make this selection, an activity monitor predicts or monitors the workload of one or more clients. Different activities of the workload are given corresponding weights. The delay in performing requested activities, modified by these weights is found, and the candidate leader with the lowest weighted delay is selected as the leader.

Claims (43)

1. A computer-implemented method comprising:

generating, by data processing hardware, predicted workload data indicating a predicted weighted communication delay for each respective computing cluster of a plurality of computing clusters in performing one or more operations for a client application;

monitoring, by the data processing hardware, interactions between the client application and a distributed database executing the plurality of computing clusters, one of the plurality of computing clusters assigned a leader computing cluster role responsible for proposing operations for the plurality of computing clusters to perform on the distributed database based on the predicted weighted communication delay;

generating, by the data processing hardware and based on the interactions between the client application and the distributed database, historical workload data indicating a historical weighted communication delay for each respective computing cluster of the plurality of computing clusters in performing one or more operations for the client application;

determining, by the data processing hardware, that a different computing cluster of the plurality of computing clusters has less historical weighted communication delay than the leader computing cluster based on the historical workload data; and

based on determining that the different computing cluster of the plurality of computing clusters has less historical weighted communication delay than the leader computing cluster, assigning, by the data processing hardware, the different computing cluster the leader computing cluster role.

2. The computer-implemented method of claim 1 , wherein the distributed database comprises one or more voters that are configured to accept or deny the proposed operations of the one of the plurality of computing clusters assigned the leader computing cluster role.

3. The computer-implemented method of claim 2 , wherein the distributed database comprises one or more replicas that replicate at least a portion of the distributed database.

4. The computer-implemented method of claim 3 , wherein the operations further comprise, in response to assigning the different computing cluster the leader computing cluster role, assigning, by the data processing hardware, one or more new voters and one or more new replicas for the distributed database.

5. The computer-implemented method of claim 4 , wherein a quantity of the one or more new voters is equal in number to a quantity of the one or more voters.

6. The computer-implemented method of claim 4 , wherein a quantity of the one or more new replicas is equal in number to a quantity of the one or more replicas.

7. The computer-implemented method of claim 1 , wherein the operations further comprise, for each respective computing cluster:

determining, by the data processing hardware and based on the historical workload data, a communication delay for the respective computing cluster interacting with the client application;

identifying, by the data processing hardware, a frequency of the interactions; and

determining, by the data processing hardware, the historical weighted communication delay based on the communication delay and the frequency of the interactions.

8. The computer-implemented method of claim 1 , wherein the historical workload data comprises logs of the interactions between the client application and the distributed database.

9. 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:

generate predicted workload data indicating a predicted weighted communication delay for each respective computing cluster of a plurality of computing clusters in performing one or more operations for a client application;

monitor interactions between the client application and a distributed database executing the plurality of computing clusters, one of the plurality of computing clusters assigned a leader computing cluster role responsible for proposing operations for the plurality of computing clusters to perform on the distributed database based on the predicted weighted communication delay;

generate, based on the monitored interactions between the client application and the distributed database, historical workload data indicating a historical weighted communication delay for each respective computing cluster of the plurality of computing clusters in performing one or more operations for the client application;

determine that a different computing cluster of the plurality of computing clusters has less historical weighted communication delay than the leader computing cluster based on the historical workload data; and

based on determining that the different computing cluster of the plurality of computing clusters has less historical weighted communication delay than the leader computing cluster, assign the different computing cluster the leader computing cluster role.

10. The system of claim 9 , wherein the distributed database comprises one or more voters that are configured to accept or deny the proposed operations of the one of the plurality of computing clusters assigned the leader computing cluster role.

11. The system of claim 10 , wherein the distributed database comprises one or more replicas that replicate at least a portion of the distributed database.

12. The system of claim 11 , wherein the instructions further cause the data processing hardware to, in response to assigning the different computing cluster the leader computing cluster role, assign one or more new voters and one or more new replicas for the distributed database.

13. The system of claim 12 , wherein a quantity of the one or more new voters is equal in number to a quantity of the one or more voters.

14. The system of claim 12 , wherein a quantity of the one or more new replicas is equal in number to a quantity of the one or more replicas.

15. The system of claim 9 , wherein the instructions further cause the data processing hardware to, for each respective computing cluster:

determine, based on the historical workload data, a communication delay for the respective computing cluster interacting with the client application;

identify a frequency of the interactions; and

determine the historical weighted communication delay based on the communication delay and the frequency of the interactions.

16. The system of claim 9 , wherein the historical workload data comprises logs of the interactions between the client application and the distributed database.

17. A non-transitory computer-readable storage medium encoded with instructions that, when executed by one or more processors of a computing system, cause the one or more processors to:

generate predicted workload data indicating a predicted weighted communication delay for each respective computing cluster of a plurality of computing clusters in performing one or more operations for a client application;

monitor interactions between the client application and a distributed database executing the plurality of computing clusters, one of the plurality of computing clusters assigned a leader computing cluster role responsible for proposing operations for the plurality of computing clusters to perform on the distributed database based on the predicted weighted communication delay;

generate, based on the monitored interactions between the client application and the distributed database, historical workload data indicating a historical weighted communication delay for each respective computing cluster of the plurality of computing clusters in performing one or more operations for the client application;

determine that a different computing cluster of the plurality of computing clusters has less historical weighted communication delay than the leader computing cluster based on the historical workload data; and

responsive to determining that the different computing cluster of the plurality of computing clusters has less historical weighted communication delay than the leader computing cluster, assign the different computing cluster the leader computing cluster role.

18. The non-transitory computer-readable storage medium of claim 17 , wherein the distributed database comprises one or more voters that are configured to accept or deny the proposed operations of the one of the plurality of computing clusters assigned the leader computing cluster role.

19. The non-transitory computer-readable storage medium of claim 18 , wherein the distributed database comprises one or more replicas that replicate at least a portion of the distributed database.

20. The non-transitory computer-readable storage medium of claim 19 , wherein the instructions further cause the one or more processors to, in response to assigning the different computing cluster the leader computing cluster role, assign assigning one or more new voters and one or more new replicas for the distributed database.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 20, 2024
From: SHRAER, ALEXANDER; SHAROV, ARTYOM; MERCHANT, ARIF ABDULHUSEIN; COOPER, BRIAN F.
To: GOOGLE INC.
Reel/Frame 066189/0814 →
CHANGE OF NAME Recorded Jan 20, 2024
From: GOOGLE INC.
To: GOOGLE LLC
Reel/Frame 068436/0196 →
Continuity (5)
Continuation 18090453 · Dec 28, 2022
Continuation 17074578 · Oct 19, 2020
Continuation 15200939 · Jul 1, 2016
Provisional Application 62188076 · Jul 2, 2015
Related Publication 20240160641A1 · May 16, 2024
References Cited (93)
US 5710915A · McElhiney · 1998 [cited by applicant]
US 6122264A · Kaufman et al. · 2000 [cited by applicant]
US 6324654B1 · Wahl et al. · 2001 [cited by applicant]
US 6401120B1 · Gamache et al. · 2002 [cited by applicant]
US 6591272B1 · Williams · 2003 [cited by applicant]
US 7478263B1 · Kownacki et al. · 2009 [cited by applicant]
US 7555516B2 · Lamport · 2009 [cited by applicant]
US 7558883B1 · Lamport · 2009 [cited by applicant]
US 7698465B2 · Lamport · 2010 [cited by applicant]
US 7797457B2 · Lamport · 2010 [cited by applicant]
US 7840662B1 · Natanzon · 2010 [cited by applicant]
US 7987152B1 · Gadir · 2011 [cited by applicant]
US 8005888B2 · Lamport · 2011 [cited by applicant]
US 8126848B2 · Wagner · 2012 [cited by applicant]
US 8380846B1 · Abu-Ghazaleh et al. · 2013 [cited by applicant]
US 8392482B1 · McAlister et al. · 2013 [cited by applicant]
US 8694647B2 · Bolosky et al. · 2014 [cited by applicant]
US 8843441B1 · Rath et al. · 2014 [cited by applicant]
US 9230000B1 · Hsieh et al. · 2016 [cited by applicant]
US 9294558B1 · Vincent et al. · 2016 [cited by applicant]
US 9740472B1 · Sohi et al. · 2017 [cited by applicant]
US 9971785B1 · Byrne et al. · 2018 [cited by applicant]
US 10831777B2 · Shraer et al. · 2020 [cited by applicant]
US 11556561B2 · Shraer et al. · 2023 [cited by applicant]
US 11907258B2 · Shraer et al. · 2024 [cited by applicant]
US 20020032883A1 · Kampe et al. · 2002 [cited by applicant]
US 20020107934A1 · Lowery et al. · 2002 [cited by applicant]
US 20020143798A1 · Lisiecki et al. · 2002 [cited by applicant]
US 20040143607A1 · Beck · 2004 [cited by applicant]
US 20040210673A1 · Cruciani et al. · 2004 [cited by applicant]
US 20040249904A1 · Moore et al. · 2004 [cited by applicant]
US 20040254984A1 · Dinker · 2004 [cited by applicant]
US 20050132154A1 · Rao et al. · 2005 [cited by applicant]
US 20050256824A1 · Vingralek · 2005 [cited by applicant]
US 20070022122A1 · Bahar et al. · 2007 [cited by applicant]
US 20070297374A1 · El-Damhougy · 2007 [cited by examiner]
US 20090063356A1 · Heise et al. · 2009 [cited by applicant]
US 20090083390A1 · Abu-Ghazaleh et al. · 2009 [cited by applicant]
US 20090089365A1 · Serghi et al. · 2009 [cited by applicant]
US 20100017495A1 · Lamport · 2010 [cited by applicant]
US 20100017644A1 · Butterworth · 2010 [cited by applicant]
US 20100131545A1 · Srivastava et al. · 2010 [cited by applicant]
US 20110178984A1 · Talius et al. · 2011 [cited by applicant]
US 20110178985A1 · San Martin Arribas et al. · 2011 [cited by applicant]
US 20120011398A1 · Eckhardt et al. · 2012 [cited by applicant]
US 20120078848A1 · Jennas, II et al. · 2012 [cited by applicant]
US 20120166390A1 · Merriman et al. · 2012 [cited by applicant]
US 20120239722A1 · Bolosky et al. · 2012 [cited by applicant]
US 20120254175A1 · Horowitz et al. · 2012 [cited by applicant]
US 20120254342A1 · Evans · 2012 [cited by applicant]
US 20120271795A1 · Rao et al. · 2012 [cited by applicant]
US 20130006687A1 · Knapp · 2013 [cited by applicant]
US 20130031048A1 · Asai et al. · 2013 [cited by applicant]
US 20130066875A1 · Combet et al. · 2013 [cited by applicant]
US 20130290249A1 · Merriman et al. · 2013 [cited by applicant]
US 20130311441A1 · Erdogan et al. · 2013 [cited by applicant]
US 20140095505A1 · Blanchflower et al. · 2014 [cited by applicant]
US 20140101100A1 · Hu et al. · 2014 [cited by applicant]
US 20140164329A1 · Guo et al. · 2014 [cited by applicant]
US 20140164831A1 · Merriman et al. · 2014 [cited by applicant]
US 20140223135A1 · Parham et al. · 2014 [cited by applicant]
US 20140344331A1 · Johns et al. · 2014 [cited by applicant]
US 20150058306A1 · Earl et al. · 2015 [cited by applicant]
US 20150161016A1 · Bulkowski et al. · 2015 [cited by applicant]
US 20150220584A1 · Isaacson · 2015 [cited by examiner]
US 20150254325A1 · Stringham · 2015 [cited by applicant]
US 20150363124A1 · Rath et al. · 2015 [cited by applicant]
US 20160034555A1 · Rahut et al. · 2016 [cited by applicant]
US 20160036924A1 · Koppolu et al. · 2016 [cited by applicant]
US 20160077936A1 · Tang et al. · 2016 [cited by applicant]
US 20160094356A1 · Xiang et al. · 2016 [cited by applicant]
US 20160366220A1 · Gottlieb et al. · 2016 [cited by applicant]
US 20170004193A1 · Shraer et al. · 2017 [cited by applicant]
US 20170006105A1 · Shraer et al. · 2017 [cited by applicant]
US 20170083410A1 · Anglin et al. · 2017 [cited by applicant]
US 20170154091A1 · Vig et al. · 2017 [cited by applicant]
US 20170264493A1 · Cencini et al. · 2017 [cited by applicant]
US 20170329798A1 · Srivas et al. · 2017 [cited by applicant]
US 20180027048A1 · Zhang · 2018 [cited by applicant]
Ardekani et al., “A Self-Configurable Geo-Replicated Cloud Storage System,” Proceedings of the 11th USENIX Symposium on Operating Systems Design and Implementation, Oct. 6, 2014, Retrieved from the Internet: URL:https:/… [cited by applicant]
Becker et al., “Leader Election for Replicated Services Using Application Scores,” Correct System Design; [Lecture Notes in Computer Science; Lect.Notes Computer], Springer International Publishing, Cham, Dec. 12, 2011,… [cited by applicant]
Ejaz et al., “Improving Wide-Area Replication Performance through Informed Leader Election and Overlay Construction,” 2013 IEEE Sixth International Conference on Cloud Computing, IEEE, Jun. 28, 2013, pp. 422-429. [cited by applicant]
International Search Report and Written Opinion in International Application No. PCT/US2016/040741, dated Oct. 13, 2016, 16 pages. [cited by applicant]
Santos et al., “Latency-aware leader election,” Applied Computing, Mar. 8, 2009, pp. 1056-1061. [cited by applicant]
Baker et al., “Megastore: Providing Scalable, Highly Available Storage for Interactive Services,” 5th Biennial Conference on Innovative Data Systems Research (CIDR '11), Jan. 9-12, 2011, pp. 223-234. [cited by applicant]
Cooper et al., “PNUTS: Yahoo!'s Hosted Data Serving Platform,” Proceedings of the VLDB Endowment, 1(2):1277-1288, Aug. 2008. [cited by applicant]
Kadambi et al., “Where in the World is My Data?” Proceedings of the VLDB Endowment, 4(11):1040-1050, Sep. 2011. [cited by applicant]
Sharov et al., “Optimizing Spanner Replica Roles,” Research Reports, Sep. 2014, 1 page. [cited by applicant]
Wolfson et al., “An Adaptive Data Replication Algorithm,” ACM Transactions on Database Systems (TODS), 22(2):255-314, Jun. 1997. [cited by applicant]
USPTO. Office Action relating to U.S. Appl. No. 18/090,453, dated May 22, 2023. [cited by applicant]
Prosecution History from U.S. Appl. No. 15/200,939, now issued U.S. Pat. No. 10,831,777, dated May 2, 2018 through Jul. 1, 2020, 133 pp. [cited by applicant]
Prosecution History from U.S. Appl. No. 17/074,578, now issued U.S. Pat. No. 11,556,561, dated Nov. 9, 2021, through Sep. 20, 2022, 67 pp. [cited by applicant]
Prosecution History from U.S. Appl. No. 18/090,453, now issued U.S. Pat. No. 11,907,258, dated May 22, 2023, through Oct. 12, 2023, 39 pp. [cited by applicant]