IP Library › Granted Patent US 11,727,009
Granted Patent B2
US 11,727,009 · App. 17/036,644 · Granted Aug 15, 2023

System and method for processing skewed datasets

Inventor: Avnish Kumar Rastogi (Noida, IN)
G06F16/2456G06F9/542
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 11,727,009
App. No.
17/036,644
Granted
Aug 15, 2023
Kind
B2
Abstract

Disclosed is a method and system for processing skewed datasets. The processor 202 is configured to capture a broadcast size of non-skewed datasets to be loaded onto a memory associated with one or more nodes in a distributed system. The skewed dataset is identified from two or more datasets to be joined. Each of the non-skewed dataset is divided into a plurality of non-skewed data chunks at the node and each of the non-skewed data chunk is broadcasted to one or more nodes having the skewed dataset. The joining operation is then performed between each of the skewed dataset and the non-skewed data chunk till all the non-skewed data chunks are consumed in the join operation. Resultant joined dataset is then collected as a single joined dataset from the nodes involved in the joining operation.

Claims (29)

1. A method of processing skewed datasets in a distributed computing environment, the method comprising:

capturing, at a first node, a broadcast size of non-skewed datasets to be loaded onto a memory associated with one or more other nodes in a distributed system;

identifying, at first node, a skewed dataset from two or more datasets to be joined, wherein the dataset comprises one or more non-skewed datasets and the skewed dataset;

dividing, at the first node, each of the non-skewed datasets into a plurality of non-skewed data chunks, wherein each of the non-skewed data chunks comprises a broadcast size chunk, wherein the broadcast size defines a maximum size of the non-skewed data chunk;

broadcasting, at the first node, each of the non-skewed data chunks to one or more other nodes having the skewed dataset, wherein the one or more other nodes are used for performing joining operation of the two or more datasets;

performing, over each of the one or more nodes, the joining operation between each of the skewed dataset and the non-skewed data chunks received as a result of the broadcasting, for obtaining a resultant joined dataset;

storing, each of the resultant joined dataset over each of the one or more nodes, wherein each of the broadcasting and the performing of the joining operation is repeated till skewed dataset is joined with the non-skewed data chunk; and

collecting, from the one or more nodes involved in the joining operation, the resultant joined dataset as a single joined dataset.

2. The method as claimed in claim 1 , wherein the skewed dataset is a larger dataset in a set of two datasets to be joined.

3. The method as claimed in claim 1 , wherein the identifying comprises detecting the skewed dataset based on data distribution on joining keys used for joining the two or more datasets, wherein the data distribution is calculated by using any predefined methodology.

4. The method as claimed in claim 1 , comprising:

capturing, a broadcast size of the two or more datasets to be loaded onto a memory associated with one or more other nodes in the distribution system, wherein the broadcast size is captured from at least one of configuration file of the one or more datasets, or command line option or from another computer based system.

5. The method as claimed in claim 1 , wherein the broadcasting comprises of a serializing each of the non-skewed data chunks, transmitting the data chunk to all the nodes and deserializing each of the non-skewed data chunks received on each node of the other nodes.

6. A system processing skewed datasets in a distributed computing environment, the system comprising:

a memory; and

a processor coupled to the memory, wherein the processor is configured to execute a set of instructions stored in the memory, wherein the processor is configured to:

capture, at a first node, a broadcast size of non-skewed datasets to be loaded onto a memory associated with one or more other nodes in a distributed system;

identify, at the first node, a skewed dataset from two or more datasets to be joined, wherein the dataset comprises one of a one or more non-skewed datasets and the skewed dataset;

divide, at the first node, each of the non-skewed dataset into a plurality of non-skewed data chunks, wherein each of the non-skewed data chunks comprises a broadcast size chunk, wherein the broadcast size defines a maximum size of the non-skewed data chunk;

broadcast, by the first node, each of the non-skewed data chunks to one or more other nodes having the skewed dataset, wherein the one or more other nodes are used for performing joining operation of the two or more datasets;

perform, over each of the one or more nodes, the joining operation between each of the skewed dataset and the non-skewed data chunk received as a result of the broadcasting, for obtaining a resultant joint dataset; and

store, each of the resultant joint dataset over each of the one or more nodes, wherein each of the broadcasting and the performing of the joining operation is repeated till a last skewed dataset of one or more skewed datasets is joined with the non-skewed data chunk,

collect, from the one or more nodes involved in the joining operation, the resultant joined dataset as a single joined dataset.

7. The system as claimed in claim 6 , wherein the skewed dataset is a larger dataset in a set of two datasets to be joined.

8. The system as claimed in claim 6 , wherein the processor is configured to:

detect, the skewed dataset based on data distribution on joining keys used for joining the two or more datasets, wherein the data distribution is calculated by using a predefined methodology.

9. The system as claimed in claim 6 , wherein the processor is configured to:

capture, a broadcast size of the two or more datasets to be loaded onto a memory associated with one or more other nodes in the distribution system, wherein the broadcast size is captured from at least one of configuration file of the one or more datasets, or command line option of the one or more datasets.

10. The system as claimed in claim 6 , wherein the broadcasting comprises one of a serializing each of the non-skewed data chunks, or deserializing each of the non-skewed data chunks on each node of the one or more other nodes.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 2, 2020
From: RASTOGI, AVNISH KUMAR
To: HCL TECHNOLOGIES LIMITED
Reel/Frame 053963/0476 →
Continuity (1)
Related Publication 20220100752A1 · Mar 31, 2022