IP Library Granted Patent US 10,642,866
Granted Patent B1
US 10,642,866 · App. 15/422,756 · Granted May 5, 2020

Automated load-balancing of partitions in arbitrarily imbalanced distributed mapreduce computations

Inventors: Wei Jiang (San Mateo, CA); Silvius V. Rus (Orinda, CA)
Assignee: Quantcast Corporation
G06F16/285G06F16/278
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 10,642,866
App. No.
15/422,756
Granted
May 5, 2020
Kind
B1
Abstract

A distributed computing system executes a MapReduce job on streamed data that includes an arbitrary amount of imbalance with respect to the frequency distribution of the data keys in the dataset. A map task module maps the dataset to a coarse partitioning, and generates a list of the top K keys with the highest frequency among the dataset. A sort task module employs a plurality of sorters to read the coarse partitioning and sort the data into buckets by data key. The values for the top K most frequent keys are separated into single-key buckets. The other less frequently occurring keys are assigned to buckets that each have multiple keys assigned to it. Then, more than one worker is assigned to each single-key bucket. The output of the multiple workers assigned to each respective single-key bucket is stitched together.

Claims (37)

1. A computer-implemented method of load-balancing in an arbitrarily imbalanced MapReduce job in a distributed computing system, the method comprising:

identifying K data keys with the highest frequency among received data, the received data comprising pairings of data keys and data values to be processed in the MapReduce job, wherein K is determined according to a data key frequency distribution of the received data and a threshold level of acceptable imbalance in reduce phase worker loads across reduce phase workers, and wherein K is a number selected to keep a maximum imbalance ratio under a threshold;

assigning data for each of the K data keys to a single-key bucket and other data keys to multiple-key buckets;

assigning one respective reduce phase worker to process data values corresponding to data keys of each multiple-key bucket, each multiple-key bucket comprising queued data items having several different keys;

assigning multiple reduce phase workers to process data values corresponding to the data key of each single-key bucket; and

stitching together output of the assigned multiple reduce phase workers on each respective single-key bucket.

2. The method of claim 1 , wherein the received data is a data stream.

3. The method of claim 1 , wherein the identified K keys with the highest frequency among the received data fluctuates as additional data are received.

4. The method of claim 1 , wherein K is a number selected so a frequency of a data key assigned to a single-key bucket with a lowest frequency among data keys assigned to single-key buckets exceeds a threshold.

5. The method of claim 1 , wherein a number of multiple reduce phase workers to assign to a single-key bucket is determined based on a respective frequency of the data key assigned to the single-key bucket and a threshold level of acceptable imbalance in reduce phase worker loads across reduce phase workers.

6. The method of claim 1 , further comprising reporting a frequency distribution of the K highest frequency data keys.

7. The method of claim 1 , further comprising combining output across all reduce phase workers to obtain a result of a MapReduce computation.

8. The method of claim 1 , wherein assigning each of the K data keys to a single-key bucket and other data keys to multiple-key buckets comprises a single sort step.

9. The method of claim 1 , wherein the single-key bucket for a first data key of the identified K data keys is arbitrarily large compared to other buckets.

10. A nontransitory computer readable storage medium including computer program instructions that, when executed, cause a computer processor to perform operations comprising:

identifying K data keys with the highest frequency among received data, the received data comprising pairings of data keys and data values to be processed in the MapReduce job, wherein K is determined according to a data key frequency distribution of the received data and a threshold level of acceptable imbalance in reduce phase worker loads across reduce phase workers, and wherein K is a number selected to keep a maximum imbalance ratio under a threshold;

assigning data for each of the K data keys to a single-key bucket and other data keys to multiple-key buckets;

assigning one respective reduce phase worker to process data values corresponding to data keys of each multiple-key bucket, each multiple-key bucket comprising queued data items having several different keys;

assigning multiple reduce phase workers to process data values corresponding to the data key of each single-key bucket; and

stitching together output of the assigned multiple reduce phase workers on each respective single-key bucket.

11. The medium of claim 10 , wherein the received data is a data stream.

12. The medium of claim 10 , wherein the identified K keys with the highest frequency among the received data fluctuates as additional data are received.

13. The medium of claim 10 , wherein a number of multiple reduce phase workers to assign to a single-key bucket is determined based on a respective frequency of the data key assigned to the single-key bucket and a threshold level of acceptable imbalance in reduce phase worker loads across reduce phase workers.

14. The medium of claim 10 , wherein the operations further comprise reporting a frequency distribution of the K highest frequency data keys.

15. The medium of claim 10 , wherein the operations further comprise combining output across all reduce phase workers to obtain a result of a MapReduce computation.

16. The medium of claim 10 , wherein assigning each of the K data keys to a single-key bucket and other data keys to multiple-key buckets comprises a single sort step.

17. The medium of claim 10 , wherein the single-key bucket for a first data key of the identified K data keys is arbitrarily large compared to other buckets.

18. The medium of claim 10 , wherein K is a number selected so a frequency of a data key assigned to a single-key bucket with a lowest frequency among data keys assigned to single-key buckets exceeds a threshold.

19. A system comprising:

a computer processor; and

a computer readable storage medium storing processor-executable computer program instructions, the computer program instructions comprising instructions for:

identifying K data keys with the highest frequency among received data, the received data comprising pairings of data keys and data values to be processed in the MapReduce job, wherein K is determined according to a data key frequency distribution of the received data and a threshold level of acceptable imbalance in reduce phase worker loads across reduce phase workers, and wherein K is a number selected to keep a maximum imbalance ratio under a threshold;

assigning data for each of the K data keys to a single-key bucket and other data keys to multiple-key buckets;

assigning one respective reduce phase worker to process data values corresponding to data keys of each multiple-key bucket, each multiple-key bucket comprising queued data items having several different keys;

assigning multiple reduce phase workers to process data values corresponding to the data key of each single-key bucket; and

stitching together output of the assigned multiple reduce phase workers on each respective single-key bucket.

20. The system of claim 19 , wherein a number of multiple reduce phase workers to assign to a single-key bucket is determined based on a respective frequency of the data key assigned to the single-key bucket and a threshold level of acceptable imbalance in reduce phase worker loads across reduce phase workers.

Assignments (9)
RELEASE OF SECURITY INTEREST Recorded Jun 21, 2024
From: BANK OF AMERICA, N.A.
To: QUANTCAST CORPORATION
Reel/Frame 067807/0017 →
SECURITY INTEREST Recorded Jun 18, 2024
From: QUANTCAST CORPORATION
To: CRYSTAL FINANCIAL LLC D/B/A SLR CREDIT SOLUTIONS
Reel/Frame 067777/0613 →
SECURITY INTEREST Recorded Dec 5, 2022
From: QUANTCAST CORPORATION
To: VENTURE LENDING & LEASING IX, INC.; WTI FUND X, INC.
Reel/Frame 062066/0265 →
RELEASE OF SECURITY INTEREST Recorded Sep 30, 2021
From: WELLS FARGO BANK, NATIONAL ASSOCIATION
To: QUANTCST CORPORATION
Reel/Frame 057678/0832 →
SECURITY INTEREST Recorded Sep 30, 2021
From: QUANTCAST CORPORATION
To: BANK OF AMERICA, N.A., AS AGENT
Reel/Frame 057677/0297 →
RELEASE OF SECURITY INTEREST Recorded Mar 15, 2021
From: TRIPLEPOINT VENTURE GROWTH BDC CORP.
To: QUANTCAST CORPORATION
Reel/Frame 055599/0282 →
SECURITY INTEREST Recorded May 20, 2020
From: QUANTCAST CORPORATION
To: WELLS FARGO BANK, NATIONAL ASSOCIATION
Reel/Frame 052717/0840 →
SECURITY INTEREST Recorded Aug 7, 2018
From: QUANTCAST CORPORATION
To: TRIPLEPOINT VENTURE GROWTH BDC CORP.
Reel/Frame 046733/0305 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded May 16, 2017
From: RUS, SILVIUS V; JIANG, WEI
To: QUANTCAST CORP.
Reel/Frame 042401/0204 →
Continuity (1)
Continuation 14320373 · Jun 30, 2014