IP Library Granted Patent US 9,330,129
Granted Patent B2
US 9,330,129 · App. 14/536,220 · Granted May 3, 2016

Organizing, joining, and performing statistical calculations on massive sets of data

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 9,330,129
App. No.
14/536,220
Granted
May 3, 2016
Kind
B2
Abstract

A system, method, and apparatus are provided for organizing and joining massive sets of data (e.g., tens or hundreds of millions of event records). A dataset is Blocked by first identifying a partition key, which comprises one or more columns of the data. Each Block will contain all dataset records that have partition key values assigned to that Block. A cost constraint (e.g., a maximum size, a maximum number of records) may also be applied to the Blocks. A Block index is generated to identify all Blocks, their corresponding (sequential) partition key values, and their locations. A second dataset that includes the partition key column(s) and that must be correlated with the first dataset may then be Blocked according to the same ranges of partition key values (but without the cost constraint). Corresponding Blocks of the datasets may then be Joined/Aggregated, and analyzed as necessary.

Claims (98)

1. A method of correlating multi-dimensional datasets, the method comprising:

partitioning a first dataset into a first set of blocks by:

selecting a partition key comprising one or more dimensions common to the first dataset and a second dataset; and

defining each block in the first set of blocks to include records of the first dataset having a partition key value associated with the block, wherein each block is associated with a different set of partition key values;

for each block in the first set of blocks, updating an index to identify:

the block;

the set of partition key values associated with the block; and

a storage location of the block;

partitioning the second dataset into a second set of blocks, wherein each block in the second set of blocks corresponds to a block in the first set of blocks and includes records of the second dataset having partition key values associated with the corresponding block; and

for each pair of corresponding blocks, comprising a first block in the first set of blocks and a second block in the second set of blocks:

storing the first block in memory;

for each of multiple sub-blocks of the second block, correlating the subblock with the first block; and

aggregating the correlations between the first block and each of the multiple sub-blocks of the second block.

2. The method of claim 1 , wherein each set of partition key values is configured such that each block of the first dataset can be entirely stored in memory allocated to a computer process that performs said correlating.

3. The method of claim 1 , further comprising:

sorting one or more blocks of one or both datasets.

4. The method of claim 1 , wherein multiple blocks in the first set of blocks are stored in a single file.

5. The method of claim 1 , further comprising incrementally updating the first set of blocks by:

receiving an incremental update to the first dataset;

dividing the incremental update according to the index to form incremental blocks corresponding to one or more blocks of the first set of blocks; and

merging the incremental blocks with corresponding blocks of the first set of blocks.

6. The method of claim 1 , wherein partitioning the first dataset comprises:

executing multiple mapper processes, wherein each mapper process handles multiple records of the first dataset and, for each record, emits the record's partition key value.

7. The method of claim 6 , wherein partitioning the first dataset further comprises:

executing multiple reducer processes, wherein each reducer process is associated with a different set of partition key values and receives records of the first dataset that have partition key values that match the associated set of partition key values.

8. The method of claim 1 , wherein correlating a sub-block of the second block with the first block comprises:

storing the sub-block in memory;

joining the sub-block with each of a plurality of sub-blocks of the first block; and

aggregating the plurality of joins.

9. The method of claim 1 , further comprising:

prior to said correlating:

assembling a daily update to the first dataset after partitioning the first dataset into the first set of blocks;

dividing the daily update into an update set of blocks based on the index; and

storing the update set of blocks in memory; and

only after said aggregating:

physically merging the update blocks with corresponding blocks of the first set of blocks.

10. The method of claim 1 , wherein:

the first dataset comprises metrics of users of an online service; and

the partition key comprises a user identifier dimension.

11. An apparatus comprising:

one or more processors; and

memory storing instructions that, when executed by the one or more processors, cause the apparatus to:

partition a first dataset into a first set of blocks by:

selecting a partition key comprising one or more dimensions common to the first dataset and a second dataset; and

defining each block in the first set of blocks to include records of the first dataset having a partition key value associated with the block, wherein each block is associated with a different set of partition key values;

for each block in the first set of blocks, update an index to identify:

the block;

the set of partition key values associated with the block; and

a storage location of the block;

partition the second dataset into a second set of blocks, wherein each block in the second set of blocks corresponds to a block in the first set of blocks and includes records of the second dataset having partition key values associated with the corresponding block; and

for each pair of corresponding blocks, comprising a first block in the first set of blocks and a second block in the second set of blocks:

store the first block in memory;

for each of multiple sub-blocks of the second block, correlate the sub-block with the first block; and

aggregate the correlations between the first block and each of the multiple sub-blocks of the second block.

12. A system, comprising:

a first multi-dimensional dataset;

a second multi-dimensional dataset;

one or more processors;

a partition module comprising a non-transitory computer-readable medium storing instructions that, when executed by the one or more processors, cause the system to:

partition a first dataset into a first set of blocks by:

selecting a partition key comprising one or more dimensions common to the first dataset and a second dataset; and

defining each block in the first set of blocks to include records of the first dataset having a partition key value associated with the block, wherein each block is associated with a different set of partition key values; and

partition the second dataset into a second set of blocks, wherein each block in the second set of blocks corresponds to a block in the first set of blocks and includes records of the second dataset having partition key values associated with the corresponding block;

an update module comprising a non-transitory computer-readable medium storing instructions that, when executed by the one or more processors, cause the system to, for each block in the first set of blocks, update an index to identify:

the block;

the set of partition key values associated with the block; and

a storage location of the block; and

a correlation module comprising a non-transitory computer-readable medium storing instructions that, when executed by the one or more processors, cause the system to:

for each pair of corresponding blocks, comprising a first block in the first set of blocks and a second block in the second set of blocks:

store the first block in memory;

for each of multiple sub-blocks of the second block, correlate the sub-block with the first block; and

aggregate the correlations between the first block and each of the multiple sub-blocks of the second block.

13. The system of claim 12 , wherein each set of partition key values is configured such that each block of the first dataset can be entirely stored in memory allocated to a computer process that performs said correlating.

14. The system of claim 12 , further comprising a sort module comprising a non-transitory computer-readable medium storing instructions that, when executed by the one or more processors, cause the system to:

sort one or more blocks of one or both datasets.

15. The system of claim 12 , wherein multiple blocks in the first set of blocks are stored in a single file.

16. The system of claim 12 , wherein the non-transitory computer-readable medium of the update module further stores instructions that, when executed by the one or more processors, cause the system to incrementally update the first set of blocks by:

receiving an incremental update to the first dataset;

dividing the incremental update according to the index to form incremental blocks corresponding to one or more blocks of the first set of blocks; and

merging the incremental blocks with corresponding blocks of the first set of blocks.

17. The system of claim 12 , wherein partitioning the first dataset comprises:

executing multiple mapper processes, wherein each mapper process handles multiple records of the first dataset and, for each record, emits the record's partition key value.

18. The system of claim 17 , wherein partitioning the first dataset further comprises:

executing multiple reducer processes, wherein each reducer process is associated with a different set of partition key values and receives records of the first dataset that have partition key values that match the associated set of partition key values.

19. The system of claim 12 , wherein correlating a sub-block of the second block with the first block comprises:

storing the sub-block in memory;

joining the sub-block with each of a plurality of sub-blocks of the first block; and

aggregating the plurality of joins.

20. The system of claim 12 , wherein the non-transitory computer-readable medium of the update module further stores instructions that, when executed by the one or more processors, cause the system to:

prior to said correlating:

assemble a daily update to the first dataset after partitioning the first dataset into the first set of blocks;

divide the daily update into an update set of blocks based on the index; and

store the update set of blocks in memory; and

only after said aggregating:

physically merge the update blocks with corresponding blocks of the first set of blocks.

21. The system of claim 12 , wherein:

the first dataset comprises metrics of users of an online service; and

the partition key comprises a user identifier dimension.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 1, 2017
From: LINKEDIN CORPORATION
To: MICROSOFT TECHNOLOGY LICENSING, LLC
Reel/Frame 044746/0001 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Feb 5, 2015
From: VEMURI, SRINIVAS S.; VARSHNEY, MANEESH; PUTTASWAMY NAGA, KRISHNA P.; LIU, RUI
To: LINKEDIN CORPORATION
Reel/Frame 034895/0008 →