IP Library › Granted Patent US 12,748,633
Granted Patent B2
US 12,748,633 · App. 18/500,966 · Granted Sep 29, 2026

Weighted auto-sharding

Inventors: Alexander Shraer (Stanford, CA); Kfir Lev-Ari (Kfar Saba, IL); Arif Abdulhusein Merchant (Los Altos, CA); Vishesh Khemani (Seattle, WA); Atul Adya (Palo Alto, CA)
Assignee: Google LLC
G06F9/5066G06F9/5083G06F9/5088G06F16/00G06F16/278H04L67/1001H04L67/148G06F2209/5017
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,748,633
App. No.
18/500,966
Granted
Sep 29, 2026
Kind
B2
Abstract

Methods, systems, and apparatus for automatic sharding and load balancing in a distributed data processing system. In one aspect, a method includes determining workload distribution for an application across worker computers and in response to determining a load balancing operation is required: selecting a first worker computer having a highest load measure relative to respective load measure of the other work computers; determining one or more move operations for a partition of data assigned to the first worker computer and a weight for each move operation; and selecting the move operation with a highest weight the selected move operation.

Claims (42)

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

partitioning a data set for an application job into a plurality of partitions based on a key;

assigning, to each worker computer in a set of worker computers, one or more partitions of the plurality of partitions;

receiving, from each respective worker computer in the set of worker computers, a respective load measure indicating a computational load of the respective worker computer;

selecting, based on the respective load measure of each respective worker computer, a first worker computer in the set of worker computers for a replication operation for replicating at least one partition from the first worker computer to a second worker computer in the set of worker computers;

replicating, the at least one partition from the first worker computer to the second worker computer;

after replicating the at least one partition from the first worker computer to the second worker computer, distributing requests for a workload associated with the at least one partition between the first worker computer and the second worker computer such that the first worker computer and the second worker computer share the workload for the at least one partition; and

in response to distributing the requests for the workload associated with the at least one partition between the first worker computer and the second worker computer, de-replicating the at least one partition from the second worker computer by deleting the at least one partition from the second worker computer.

2 . The method of claim 1 , wherein the key comprises an atomic unit of work placement.

3 . The method of claim 1 , wherein each worker computer in the set of worker computers receives a different one or more partitions of the plurality of partitions.

4 . The method of claim 1 , wherein the operations further comprise determining an initial workload distribution for the application job to each worker computer in the set of worker computers.

5 . The method of claim 1 , wherein the respective load measure of the first worker computer is higher than the respective load measure of each other respective worker computer in the set of worker computers.

6 . The method of claim 1 , wherein the operations further comprise:

determining, for each partition, a constituent load measure for the partition;

determining pairs of adjacent partitions, each pair comprising two partitions that collectively have a contiguous range of key values; and

for each pair of adjacent partitions for which a sum of the constituent load measures of the partition does not meet a load measure merger threshold, merging the adjacent partitions into a single partition.

7 . The method of claim 6 , wherein merging the adjacent partitions into a single partition occurs prior to receiving, from each respective worker computer in the set of worker computers, the respective load measure.

8 . The method of claim 1 , wherein the respective load measure of the second worker computer is lower than the respective load measure of each other respective worker computer in the set of worker computers.

9 . The method of claim 1 , wherein the first worker computer of the set of worker computers is assigned a different number of partitions than the second worker computer of the set of worker computers.

10 . The method of claim 1 , wherein the operations further comprise determining, based on the respective load measure of the first worker computer, a weight of a move operation.

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:

partitioning a data set for an application job into a plurality of partitions based on a key;

assigning, to each worker computer in a set of worker computers, one or more partitions of the plurality of partitions;

receiving, from each respective worker computer in the set of worker computers, a respective load measure indicating a computational load of the respective worker computer;

selecting, based on the respective load measure of each respective worker computer, a first worker computer in the set of worker computers for a replication operation for replicating at least one partition from the first worker computer to a second worker computer in the set of worker computers;

replicating, the at least one partition from the first worker computer to the second worker computer;

after replicating the at least one partition from the first worker computer to the second worker computer, distributing requests for a workload associated with the at least one partition between the first worker computer and the second worker computer such that the first worker computer and the second worker computer share the workload for the at least one partition; and

in response to distributing the requests for the workload associated with the at least one partition between the first worker computer and the second worker computer, de-replicating the at least one partition from the second worker computer by deleting the at least one partition from the second worker computer.

12 . The system of claim 11 , wherein the key comprises an atomic unit of work placement.

13 . The system of claim 11 , wherein each worker computer in the set of worker computers receives a different one or more partitions of the plurality of partitions.

14 . The system of claim 11 , wherein the operations further comprise determining an initial workload distribution for the application job to each worker computer in the set of worker computers.

15 . The system of claim 11 , wherein the respective load measure of the first worker computer is higher than the respective load measure of each other respective worker computer in the set of worker computers.

16 . The system of claim 11 , wherein the operations further comprise:

determining, for each partition, a constituent load measure for the partition;

determining pairs of adjacent partitions, each pair comprising two partitions that collectively have a contiguous range of key values; and

for each pair of adjacent partitions for which a sum of the constituent load measures of the partition does not meet a load measure merger threshold, merging the adjacent partitions into a single partition.

17 . The system of claim 16 , wherein merging the adjacent partitions into a single partition occurs prior to receiving, from each respective worker computer in the set of worker computers, the respective load measure.

18 . The system of claim 11 , wherein the respective load measure of the second worker computer is lower than the respective load measure of each other respective worker computer in the set of worker computers.

19 . The system of claim 11 , wherein the first worker computer of the set of worker computers is assigned a different number of partitions than the second worker computer of the set of worker computers.

20 . The system of claim 11 , wherein the operations further comprise determining, based on the respective load measure of the first worker computer, a weight of a move operation.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 2, 2023
From: SHRAER, ALEXANDER; LEV-ARI, KFIR; MERCHANT, ARIF ABDULHUSEIN; KHEMANI, VISHESH; ANYA, ATUL
To: GOOGLE INC.
Reel/Frame 065441/0507 →
CHANGE OF NAME Recorded Nov 2, 2023
From: GOOGLE INC.
To: GOOGLE LLC
Reel/Frame 065448/0072 →
Continuity (5)
Continuation 17663618 · May 16, 2022
Continuation 16725472 · Dec 23, 2019
Continuation 15428844 · Feb 9, 2017
Provisional Application 62345567 · Jun 3, 2016
Related Publication 20240064196A1 · Feb 22, 2024
References Cited (20)
US 7043621B2 · Merchant et al. · 2006 [cited by applicant]
US 7143170B2 · Swildens · 2006 [cited by examiner]
US 7421497B2 · Rosenbach et al. · 2008 [cited by applicant]
US 8849749B2 · Rishel · 2014 [cited by examiner]
US 8976636B1 · Martin et al. · 2015 [cited by applicant]
US 10530844B2 · Shraer · 2020 [cited by examiner]
US 11363096B2 · Shraer · 2022 [cited by examiner]
US 11838356B2 · Shraer · 2023 [cited by examiner]
US 20040230764A1 · Merchant et al. · 2004 [cited by applicant]
US 20040260684A1 · Agrawal · 2004 [cited by examiner]
US 20050114862A1 · Bisdikian · 2005 [cited by examiner]
US 20120254445A1 · Kawamoto et al. · 2012 [cited by applicant]
US 20120303791A1 · Calder · 2012 [cited by examiner]
US 20130204991A1 · Skjolsvold et al. · 2013 [cited by applicant]
US 20130290249A1 · Merriman et al. · 2013 [cited by applicant]
US 20140137134A1 · Roy · 2014 [cited by examiner]
US 20140164443A1 · Genc et al. · 2014 [cited by applicant]
US 20150249615A1 · Chen et al. · 2015 [cited by applicant]
US 20160371353A1 · Serafini et al. · 2016 [cited by applicant]
US 20170031908A1 · Liu et al. · 2017 [cited by applicant]