IP Library › Granted Patent US 10,678,794
Granted Patent B2
US 10,678,794 · App. 15/858,489 · Granted Jun 9, 2020

Skew detection and handling in a parallel processing relational database system

Inventors: Bhashyam Ramesh (West Marredpally, IN); Suresh Kumar Jami (Medak, IN)
Assignee: Teradata US, Inc.
G06F16/24552G06F11/1435G06F16/211G06F16/215G06F16/2255G06F16/2456G06F16/24537G06F16/24544G06F16/278G06F16/284G06F16/9014G06F9/3887
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,678,794
App. No.
15/858,489
Granted
Jun 9, 2020
Kind
B2
Abstract

A method for detecting and handling skew and spillover in in-memory hash join operations. To improve the detection of skew and spillover in parallel processing systems, a Poisson distribution of unique hash values to Units of Parallelism (UoPs) is employed to determine the number of rows per UoP and in turn, the potential of spillover at a UoP. Hash join plan options can be selected or adjusted to reduce the likelihood of spillover.

Claims (36)

1. A method for reducing spillover during a hash join for joining a small database table and a large database table in a parallel processing relational database system, said system comprising a plurality of Units of Parallelism (UoPs), said method comprising the steps of:

evaluating, by a processor, multiple plan options for said hash join, wherein said multiple plan options include alternatives for distributing rows from said large database table across said plurality of UoPs, and wherein evaluating comprises for each one of said multiple plan options:

determining a highest number of said rows to be distributed to a UoP using a Poisson distribution of the number of unique hash values to the number of UoPs; and

determining a probability of spillover at said UoP from said highest number of rows;

selecting, by said processor, from said multiple plan options, a plan option having the lowest probability for spillover; and

executing, by said processor, said hash join in accordance with said selected plan option to distribute said rows to said plurality of UoPs.

2. A method for reducing cost of a hash join for joining a small database table and a large database table in a parallel processing relational database system, said system comprising a plurality of Units of Parallelism (UoPs), said method comprising the steps of:

evaluating, by a processor, multiple plan options for said hash join, wherein said multiple plan options include alternatives for distributing rows from said large database table across said plurality of UoPs, and wherein evaluating comprises for each one of said multiple plan options:

determining a highest number of rows to be distributed to a UoP using a Poisson distribution of the number of unique hash values to the number of UoPs; and

determining a cost associated with said highest number of rows;

selecting, by said processor, from said multiple plan options, a plan option having the lowest cost; and

executing, by said processor, said hash join in accordance with said selected plan option to distribute said rows to said plurality of UoPs.

3. The method in accordance with claim 1 , wherein:

said Poisson distribution comprises the formula:

E =ceiling ( M *(1−exp(−1*( U/M )))); and

C=U/E;

where:

M=maximum number of UoPs,

U=number of unique hash values,

E=UoP to which data is distributed, and

C=computed unique values; and

said highest number of rows to be distributed is determined using the relation R=((C−1)*rows per value from the original data)+HighModeFrequency, where R=rows per UoP.

4. The method in accordance with claim 1 , wherein:

said hash join comprises an in-memory hash join.

5. The method in accordance with claim 2 , wherein:

said Poisson distribution comprises the formula:

E =ceiling ( M *(1−exp(−1*( U/M )))); and

C=U/E;

where:

M=maximum number of UoPs,

U=number of unique hash values,

E=UoP to which data is distributed, and

C=computed unique values; and

said highest number of rows to be distributed is determined using the relation R=((C−1)*rows per value from the original data)+HighModeFrequency, where R=rows per UoP.

6. The method in accordance with claim 2 , wherein:

said hash join comprises an in-memory hash join.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 5, 2018
From: RAMESH, BHASHYAM; JAMI, SURESH KUMAR
To: TERADATA US, INC.
Reel/Frame 045013/0306 →
Continuity (4)
Division 15631224 · Jun 23, 2017
Provisional Application 62354288 · Jun 24, 2016
Provisional Application 62354262 · Jun 24, 2016
Related Publication 20180121563A1 · May 3, 2018