IP Library Granted Patent US 10,185,743
Granted Patent B2
US 10,185,743 · App. 14/553,786 · Granted Jan 22, 2019

Method and system for optimizing reduce-side join operation in a map-reduce framework

Inventors: Srikanth Sundarrajan (Velachery, IN); Shwetha G. Shivalingamurthy (Bangalore, IN)
Assignee: InMobi PTE Ltd.
G06F17/30442
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,185,743
App. No.
14/553,786
Granted
Jan 22, 2019
Kind
B2
Abstract

The present invention provides a system and method for optimizing reduce-side join operation in a map-reduce framework. The system and method executing one or more map operations on the second data structure, grouping the data tuples to a single region of the second data structure, providing the grouped data to a single reducer and, selecting one of scan approach and a look-up approach by one or more reducers based on region key count value and pre-determined conditions of the user.

Claims (62)

1. A computer system for optimizing reduce-side join operation in a Map-reduce framework between a first data structure and a second data structure, the first data structure being sorted and divided into one or more regions, the system comprising:

one or more processors; and

a non-transitory memory that includes modules that are executable by said one or more processors, wherein the modules include:

an executing module to execute one or more map operations by one or more processors, wherein to execute one or more map operations by one or more processors comprises to:

fetch input data of the second data structure;

partition the data of the second data structure according to key-value pair;

project the key-value pairs of the second data structure to a partitioner;

maintain one or more region key counters; wherein the region key counter being used for registering key count value of one or more regions of the second data structure; and

emit the key count value of one or more regions and corresponding data, wherein the key count values are emitted prior to the corresponding data;

a grouping module to group mapped data corresponding to a single region of the second data structure;

an accumulating module to provide the grouped data to a reducer;

a fetching module to retrieve descriptive metadata of one or more regions of the first data structure; and

a selecting module to select one of a look-up approach and a scan approach to perform the join operation by one or more reducers based on associated key count value and predefined criteria by the reducer, to perform the join operation.

2. A method for optimizing reduce-side join operation in a Map-reduce framework between a first data structure and a second data structure, the first data structure being sorted and divided into one or more regions, the method comprising:

executing instructions, stored in a memory, by one or more processors to perform:

executing one or more map operations, wherein executing one or more map operations comprises:

fetching input data of the second data structure;

partitioning the data of the second data structure according to key-value pair;

projecting the key-value pairs of the second data structure to a partitioner;

maintaining one or more region key counters; wherein the region key counter being used for registering key count value of one or more regions of the second data structure; and

emitting the key count value of one or more regions and corresponding data, wherein the key count values are emitted prior to the corresponding data; and

grouping mapped data corresponding to a single region of the second data structure;

providing the grouped data to a reducer;

retrieving descriptive metadata of one or more regions of the first data structure; and

selecting one of a look-up approach and a scan approach to perform the join operation by one or more reducers based on associated key count value and predefined criteria by the reducer, for performing the join operation.

3. The method as claimed in claim 2 , wherein the descriptive metadata comprises region key count value of one or more regions of the first data structure.

4. The method as claimed in claim 2 , wherein each set of mapped data includes a set of tuples, each tuple characterized by key/value pair, wherein the keys and values are sets of attributes.

5. The method as claimed in claim 2 , wherein:

the join operation is carried out by a plurality of reducers; and

the data that is not intermediate data, for a particular reducer, includes data that is associated with another reducer.

6. The method as claimed in claim 2 , wherein the join operation includes relating the data among the plurality of the data structures.

7. The method as claimed in claim 2 , wherein executing the instructions is implemented using a cluster of machines.

8. The method as claimed in claim 2 , wherein the first data structure being sorted and divided into data tuples is stored in key count format.

9. The method as claimed in claim 2 , wherein executing instructions comprises executing instructions by processors of a cluster of computers, the method further comprising persisting memory cache across the cluster of computers wherein:

a failure of one of the computers results in replacing the failed computer with a different computer;

the replacing is performed as a single transaction; and

a redundant copy of the data is obtained from one or more remaining computers in the cluster.

10. The system as claimed in claim 1 , wherein the descriptive metadata comprises region key count value of one or more regions of the first data structure.

11. The system as claimed in claim 1 , wherein each set of mapped data includes a set of tuples, each tuple characterized by key/value pair, wherein the keys and values are sets of attributes.

12. The system as claimed in claim 1 , wherein:

the join operation is carried out by a plurality of reducers; and

the data that is not intermediate data, for a particular reducer, includes data that is associated with another reducer.

13. The system as claimed in claim 1 , wherein the join operation includes relating the data among the plurality of the data structures.

14. The system as claimed in claim 1 , wherein the one or more processors comprises multiple processors included in a cluster of computers.

15. The system as claimed in claim 1 , wherein the first data structure being sorted and divided into data tuples is stored in key count format.

16. The system as claimed in claim 1 , wherein the one or more processors comprises multiple processors included in a cluster of computers.

17. The system as claimed in claim 16 wherein the cluster of computers include a persisting memory cache across the cluster of computers, the system further comprising instructions in one or more of the cluster of computers that executable to:

replace a failed of one of the computers with a different computer as a single transaction; and

obtain a redundant copy of the data from one or more remaining computers in the cluster.

18. A non-transitory computer program product comprising instructions stored in the computer program product for optimizing reduce-side join operation in a Map-reduce framework between a first data structure and a second data structure, the first data structure being sorted and divided into one or more regions, wherein the instructions are executable by a processor to perform:

executing one or more map operations, wherein executing one or more map operations comprises:

fetching input data of the second data structure;

partitioning the data of the second data structure according to key-value pair;

projecting the key-value pairs of the second data structure to a partitioner;

maintaining one or more region key counters; wherein the region key counter being used for registering key count value of one or more regions of the second data structure; and

emitting the key count value of one or more regions and corresponding data, wherein the key count values are emitted prior to the corresponding data; and

grouping mapped data corresponding to a single region of the second data structure;

providing the grouped data to a reducer;

retrieving descriptive metadata of one or more regions of the first data structure; and

selecting one of a look-up approach and a scan approach to perform the join operation by one or more reducers based on associated key count value and predefined criteria by the reducer, for performing the join operation.

19. The non-transitory computer program product in claim 18 , wherein the descriptive metadata comprises region key count value of one or more regions of the first data structure.

20. The non-transitory computer program product in claim 18 , wherein each set of mapped data includes a set of tuples, each tuple characterized by key/value pair, wherein the keys and values are sets of attributes.

Assignments (10)
SECURITY INTEREST Recorded Apr 1, 2026
From: INMOBI TECHNOLOGY SERVICES PTE. LTD.
To: MADISON PACIFIC TRUST LIMITED
Reel/Frame 074244/0228 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 31, 2026
From: INMOBI PTE LTD.
To: INMOBI TECHNOLOGY SERVICES PTE. LTD.
Reel/Frame 074233/0395 →
RELEASE OF SECURITY INTEREST Recorded Dec 31, 2025
From: MARS GROWTH CAPITAL PRE-UNICORN FUND, L.P.
To: INMOBI PTE LTD.; INMOBI HOLDINGS PTE LTD.
Reel/Frame 073343/0448 →
RELEASE OF SECURITY INTEREST Recorded Dec 31, 2025
From: MARS GROWTH CAPITAL PRE-UNICORN FUND, L.P.
To: INMOBI PTE LTD.; INMOBI HOLDINGS PTE LTD.
Reel/Frame 073343/0481 →
SECURITY INTEREST Recorded Dec 31, 2025
From: INMOBI PTE LTD.
To: MADISON PACIFIC TRUST LIMITED
Reel/Frame 073343/0572 →
CORRECTIVE ASSIGNMENT TO CORRECT THE THE PROPERTY TYPE FOR NUMBERS 10725921, 11244354, 11455274, AND 11330398 FROM APPLICATION NUMBERS TO PATENT NUMBERS PREVIOUSLY RECORDED ON REEL 68126 FRAME 833. ASSIGNOR(S) HEREBY CONFIRMS THE SECURITY INTEREST. Recorded Aug 5, 2024
From: INMOBI PTE. LTD.; INMOBI HOLDINGS PTE. LTD.
To: MARS GROWTH CAPITAL PRE-UNICORN FUND, L.P.
Reel/Frame 068309/0178 →
SECURITY INTEREST Recorded Jul 30, 2024
From: INMOBI PTE. LTD.; INMOBI HOLDINGS PTE. LTD.
To: MARS GROWTH CAPITAL PRE-UNICORN FUND, L.P.
Reel/Frame 068126/0833 →
RELEASE OF SECURITY INTEREST IN PATENTS AT REEL 53147/FRAME 0341 Recorded Jul 30, 2024
From: CRESTLINE DIRECT FINANCE, L.P.
To: INMOBI PTE. LTD.
Reel/Frame 068202/0824 →
SECURITY INTEREST Recorded Jul 8, 2020
From: INMOBI PTE. LTD.
To: CRESTLINE DIRECT FINANCE, L.P., AS COLLATERAL AGENT FOR THE RATABLE BENEFIT OF THE SECURED PARTIES
Reel/Frame 053147/0341 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jun 16, 2016
From: SUNDARRAJAN, SRIKANTH; SHIVALINGAMURTHY, SHWETHA G
To: INMOBI PTE. LTD.
Reel/Frame 039050/0169 →
Priority Claims (1)
IN 5424/CHE/2013 · Nov 26, 2013 · national
Continuity (1)
Related Publication 20150149437A1 · May 28, 2015