IP Library Granted Patent US 11,650,990
Granted Patent B2
US 11,650,990 · App. 16/084,529 · Granted May 16, 2023

Method, medium, and system for joining data tables

Inventors: Dong Xu (Hangzhou, CN); Weiguang Sun (Hangzhou, CN); Jiehong Lian (Hangzhou, CN); Longzhong Wang (Hangzhou, CN)
Assignee: Alibaba Group Holding Limited
G06F16/2456G06F16/24544G06F16/283
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,650,990
App. No.
16/084,529
Granted
May 16, 2023
Kind
B2
Abstract

Data tables that are located in a distributed data warehouse are joined with a target join-calculating algorithm that has been selected from a number of table joining algorithms which have been compared to each other based on the execution costs of each of the number of table joining algorithms.

Claims (139)

1. A method of joining data tables, the method comprising:

determining a plurality of to-be-joined data tables in a distributed data warehouse that has a plurality of computing nodes and a plurality of storage nodes;

setting, based on an environment of the distributed data warehouse where the plurality of to-be-joined data tables are located, a parameter list for cost estimation;

obtaining a plurality of table joining algorithms;

estimating a plurality of execution costs for the plurality of table joining algorithms such that each of the plurality of table joining algorithms has an estimated execution cost to join the plurality of to-be-joined data tables, wherein the estimating comprises:

determining, for each of the plurality of table joining algorithms, a number of execution steps and one or more key operations of a plurality of operations in each of the execution steps;

estimating an execution cost for each of the execution steps based on:

computing execution costs of the one or more key operations in each of the execution steps based on parameters selected from the parameter list for the one or more key operations; and

determining the execution cost for each of the execution steps based on at least one of: superposing the execution costs of the one or more key operations, and directly using [[the]]an execution cost of a selected key operation out of the one or more key operations; and

determining the estimated execution cost for each of the plurality of table joining algorithms based on the execution cost for each of the execution steps;

selecting a target algorithm from the plurality of table joining algorithms based on the estimated execution cost of each of the plurality of table joining algorithms; and

joining the plurality of to-be-joined data tables with the target algorithm.

2. The method according to claim 1 , wherein:

execution steps of a Partitioned Sort Join (PSJ) table joining algorithm include:

a re-partition execution step having a plurality of key operations, wherein the plurality of key operations in the re-partition execution step includes a local read operation, a network read operation, a local sort operation, and a local write operation; and

a sort join execution step having a key operation, wherein the key operation in the sort join execution step includes an output operation;

execution steps of a Broadcasted Hash Join (BHJ) table joining algorithm include:

a broadcast execution step having a key operation, wherein the key operation in the broadcast execution step includes the network read operation; and

a hash join execution step having a key operation, wherein the key operation in the hash join execution step includes the output operation; and

execution steps of a Blocked Hash Join (BKHJ) table joining algorithm include:

a broadcast distribution execution step having a plurality of key operations, wherein the plurality of key operations in the broadcast distribution execution step includes the local read operation, the network read operation, and the local write operation; and

a hash join execution step having a key operation, wherein the key operation in the hash join execution step includes the output operation.

3. The method according to claim 2 , further comprising obtaining target parameters needed for each of the execution steps, wherein obtaining the target parameters includes:

for the re-partition execution step, obtaining parameters N, L, RC, RNC, and WC as target parameters needed for the re-partition execution step;

for the broadcast execution step, obtaining parameters N i , N k , D, L, and RNC as target parameters needed for the broadcast execution step;

for the broadcast distribution execution step, obtaining parameters N, L, RC, RNC, and WC as target parameters needed for the broadcast distribution execution step; and

for the sort join execution step or the hash join execution step, obtaining NJ and n as target parameters needed for the sort join execution step or the hash join execution step, where

N represents a total number of data records;

L represents an average length of each of the data records;

RC represents a unit cost for the local read operation;

RNC represents a unit cost for the network read operation;

WC represents a unit cost for the local write operation;

Nk represents a number of data records contained in a main data table in the plurality of to-be-joined data tables, where k is any value of 1. . . n;

N i represents a number of data records contained in an i-th auxiliary data table in the plurality of to-be-joined data tables, where i=1. . . n and i≠k;

D represents a size of a data block supported by each of the plurality of storage nodes based on a file system of the distributed data warehouse;

N j represents a number of data records contained in a j-th data table in the plurality of to-be-joined data tables, where j=1. . . n; and

n represents a number of data tables in the plurality of to-be-joined data tables.

4. The method according to claim 3 , wherein estimating the execution cost includes:

for the re-partition execution step, estimating an execution cost for the local read operation as a data record number consumption of 0, a CPU consumption of 0, and an I/O consumption of N*L*RC (0, 0, N*L*RC), an execution cost for the network read operation as a data record number consumption of N, a CPU consumption of 0, and an I/O consumption of N*L*RNC (N, 0, N*L*RNC), an execution cost for the local sort operation as the data record number consumption of 0, the CPU consumption of N, and the I/O consumption of 0 (0, N, 0), and an execution cost for the local write operation as the data record number consumption of 0, a CPU consumption of 0, and the I/O consumption of N*L*WC (0, 0, N*L*WC) based on the parameters N, L, RC, RNC, and WC;

for the broadcast execution step, estimating the execution cost for the network read operation as the data record number consumption of ΣN i *M, the CPU consumption of 0, and the I/O consumption of ΣN i *M*L*RNC (ΣN i *M, 0, ΣN i *M*L*RNC) based on the parameters N i , N k , D, L, and RNC, where M=N k /D;

for the broadcast distribution execution step, estimating the execution cost for the local read operation as the data record number consumption of 0, the CPU consumption of 0, and the I/O consumption of N*L*RC (0, 0, N*L*RC), the execution cost for the network read operation as the data record number consumption of N, the CPU consumption of 0, and the I/O consumption of N*L*RNC (N, 0, N*L*RNC), and the execution cost for the local write operation as the data record number consumption of 0, the CPU consumption of 0, and the I/O consumption of N*L*WC (0, 0, N*L*WC) based on the parameters N, L, RC, RNC, and WC; and

for the sort join execution step or the hash join execution step, estimating an execution cost for the output operation as the data record number consumption of J, the CPU consumption of 0, and the I/O consumption of 0 (J, 0, 0) based on the parameters N j and n, where J=(ΠN j ) 1/n .

5. The method according to claim 4 , wherein estimating an execution cost further includes:

for the re-partition execution step, superposing the execution cost (0, 0, N*L*RC) of the local read operation, the execution cost (N, 0, N*L*RNC) of the network read operation, the execution cost (0, N, 0) of the local sort operation, and the execution cost (0, 0, N*L*WC) of the local write operation to obtain an execution cost (N, N, N*L*(RC+RNC+WC)) as an execution cost for the re-partition execution step;

for the broadcast execution step, using the execution cost (Σ i *M, 0, ΣN i *M*L*RNC) for the network read operation as an execution cost for the broadcast execution step;

for the broadcast distribution execution step, superposing the execution cost (0, 0, N*L*RC) for the local read operation, the execution cost (N, 0, N*L*RNC) for the network read operation, and the execution cost (0, 0, N*L*WC) for the local write operation to obtain an execution cost (N, 0, N*L*(RC+RNC+WC)) as an execution cost for the broadcast distribution execution step; and

for the sort join execution step or the hash join execution step, using the execution cost (J, 0, 0) for the output operation as an execution cost for the sort join execution step or the hash join execution step.

6. The method according to claim 5 , wherein prior to, for the re-partition execution step, superposing the execution cost (0, 0, N*L*RC) of the local read operation, the execution cost (N, 0, N*L*RNC) of the network read operation, the execution cost (0, N, 0) of the local sort operation, and the execution cost (0, 0, N*L*WC) of the local write operation to obtain an execution cost (N, N, N*L*(RC+RNC+WC)) as an execution cost for the re-partition execution step:

determining whether a skewed distribution occurs with respect to data records contained in the to-be-joined data tables; and

if a determination result is positive, then correcting the execution cost (N, 0, N*L*RNC) for the network read operation to (N, 0, P*N*L*p*RNC) and correcting the execution cost (0, 0, N*L*WC) for the local write operation to (0, 0, P*N*L*p*WC), where

p represents a distribution skewness; and

P represents a number of computing nodes for performing a join processing on the to-be-joined data tables.

7. The method according to claim 6 , wherein prior to, for the hash join execution step, using the execution cost (J, 0, 0) for the output operation as the execution cost for the hash join execution step:

determining whether a data table having a number of data records greater than the size D of the data block supported by each of the plurality of storage node exists in each auxiliary data table; and

if a determination result is positive, then correcting the execution cost (J, 0, 0) for the output operation to obtain a corrected execution cost (J, N k *ΣN l , N k *ΣN l *L*WC) as the execution cost for the hash join execution step,

where N l represent an l-st data table having a number of data records greater than the size D of the data block supported by each of the plurality of storage node; and l=1 . . . n and l≠k.

8. The method according to claim 1 , wherein setting the parameter list includes:

setting a number of data records, a total number of data records, and an average length of each data record comprised in each data table in the plurality of to-be-joined data tables;

setting, based on a file system of the distributed data warehouse, a size of a data block supported by each of the plurality of storage nodes; and

setting, based on hardware information of the distributed data warehouse, unit costs for operations needed for join calculations and a number of data records that each computing node can process.

9. The method according to claim 8 , wherein setting the unit costs for the operations needed for the join calculations includes:

determining, based on a storage medium used by the distributed data warehouse, a unit cost for a local read operation and a unit cost for a local write operation; and

determining, based on a network topology of the distributed data warehouse, a unit cost for a network read operation and a unit cost for a network write operation.

10. The method according to claim 1 , wherein the execution cost is represented by a data record number consumption, a CPU consumption, and an I/O consumption.

11. The method according to claim 1 , wherein estimating the execution cost for each of the execution steps is further based on:

determining at least one of:

a skewed distribution with respect to data records included in the plurality of to-be-joined data tables; and

existence of a data table including data records greater than a data block size supported by each of the plurality of storage nodes; and

correcting the computed execution cost of the one or more key operations based on the determination of at least one of the skewed distribution and the existence of a data table.

12. The method according to claim 1 , wherein one or more execution steps have the plurality of operations, and wherein the one or more key operations have higher execution costs than other operations of the plurality of operations.

13. A non-transitory computer-readable storage medium having embedded therein program instructions, which when executed by a processor causes the processor to execute a method of joining data tables, the method comprising:

determining a plurality of to-be-joined data tables in a distributed data warehouse that has a plurality of computing nodes and a plurality of storage nodes;

setting, based on environment of the distributed data warehouse where the plurality of to-be-joined data tables are located, a parameter list of cost estimation;

obtaining a plurality of table joining algorithms;

estimating a plurality of execution costs for the plurality of table joining algorithms such that each table joining algorithm has an estimated execution cost to join the plurality of to-be-joined data tables, wherein the estimating comprises;

determining, for each of the plurality of table joining algorithms, a number of execution steps and one or more key operations of a plurality of operations in each of the execution steps;

estimating an execution cost for each of the execution steps based on:

computing execution costs of the one or more key operations in each of the execution steps based on parameters selected from the parameter list for the one or more key operations; and

determining the execution cost for each of the execution steps based on at least one of: superposing the execution costs of the one or more key operations, and directly using an execution cost of a selected key operation out of the one or more key operations; and

determining the estimated execution cost for each of the plurality of table joining algorithms based on the execution cost for each of the execution steps;

selecting a target algorithm from the plurality of table joining algorithms based on the estimated execution cost of each of the plurality of table joining algorithms; and

joining the plurality of to-be-joined data tables with the target algorithm.

14. The non-transitory computer-readable storage medium according to claim 13 , wherein:

execution steps of a Partitioned Sort Join (PSJ) table joining algorithm include:

a re-partition execution step having a plurality of key operations, wherein the plurality of key operations in the re-partition execution step includes a local read operation, a network read operation, a local sort operation, and a local write operation; and

a sort join execution step having a key operation, wherein the key operation in the sort join execution step includes an output operation;

execution steps of a Broadcasted Hash Join (BHJ) table joining algorithm include:

a broadcast execution step having a key operation, wherein the key operation in the broadcast execution step includes the network read operation; and

a hash join execution step having a key operation, wherein the key operation in the hash join execution step includes the output operation; and

execution steps of a Blocked Hash Join (BKHJ) table joining algorithm include:

a broadcast distribution execution step having a plurality of key operations, wherein the plurality of key operations in the broadcast distribution execution step includes the local read operation, the network read operation, and the local write operation; and

a hash join execution step having a key operation, wherein the key operation in the hash join execution step includes the output operation.

15. The non-transitory computer-readable storage medium according to claim 14 , wherein the method further comprises obtaining target parameters needed for each of the execution steps, and wherein obtaining the target parameters includes:

for the re-partition execution step, obtaining parameters N, L, RC, RNC, and WC as target parameters needed for the re-partition execution step;

for the broadcast execution step, obtaining parameters N i , N k , D, L, and RNC as target parameters needed for the broadcast execution step;

for the broadcast distribution execution step, obtaining parameters N, L, RC, RNC, and WC as target parameters needed for the broadcast distribution execution step; and

for the sort join execution step or the hash join execution step, obtaining NJ and n as target parameters needed for the sort join execution step or the hash join execution step, where

N represents a total number of data records;

L represents an average length of each of the data records;

RC represents a unit cost for the local read operation;

RNC represents a unit cost for the network read operation;

WC represents a unit cost for the local write operation;

N k represents a number of data records contained in a main data table in the plurality of to-be-joined data tables, where k is any value of 1. . . n;

N i represents a number of data records contained in an i-th auxiliary data table in the plurality of to-be-joined data tables, where i=1...n and i≠k;

D represents a size of a data block supported by each of the plurality of storage nodes based on a file system of the distributed data warehouse;

N j represents a number of data records contained in a j-th data table in the plurality of to-be-joined data tables, where j=1. . . n; and

n represents a number of data tables in the plurality of to-be-joined data tables.

16. The non-transitory computer-readable storage medium according to claim 15 , wherein estimating the execution cost includes:

for the re-partition execution step, estimating an execution cost for the local read operation as a data record number consumption of 0, a CPU consumption of 0, and an I/O consumption of N*L*RC (0, 0, N*L*RC), an execution cost for the network read operation as a data record number consumption of N, a CPU consumption of 0, and an I/O consumption of N*L*RNC (N, 0, N*L*RNC), an execution cost for the local sort operation as the data record number consumption of 0, the CPU consumption of N, and the I/O consumption of 0 (0, N, 0), and an execution cost for the local write operation as the data record number consumption of 0, the CPU consumption of 0, and the I/O consumption of N*L*WC (0, 0, N*L*WC) based on the parameters N, L, RC, RNC, and WC;

for the broadcast execution step, estimating the execution cost for the network read operation as the data record number consumption of ΣN i *M, the CPU consumption of 0, and the I/O consumption of ΣN i *M*L*RNC (ΣN i *M, 0, ΣN i *M*L*RNC) based on the parameters N i , N k , D, L, and RNC, where M=N k /D;

for the broadcast distribution execution step, estimating the execution cost for the local read operation as the data record number consumption of 0, the CPU consumption of 0, and the I/O consumption of N*L*RC (0, 0, N*L*RC), the execution cost for the network read operation as the data record number consumption of N, the CPU consumption of 0, and the I/O consumption of N*L*RNC (N, 0, N*L*RNC), and the execution cost for the local write operation as the data record number consumption of 0, the CPU consumption of 0, and the I/O consumption of N*L*WC (0, 0, N*L*WC) based on the parameters N, L, RC, RNC, and WC; and

for the sort join execution step or the hash join execution step, estimating an execution cost for the output operation as the data record number consumption of J, the CPU consumption of 0, and the I/O consumption of 0 (J, 0, 0) based on the parameters N J and n, where J=(ΠN j ) 1/n .

17. The non-transitory computer-readable storage medium according to claim 15 , wherein estimating the execution cost further includes:

for the re-partition execution step, superposing the execution cost (0, 0, N*L*RC) of the local read operation, the execution cost (N, 0, N*L*RNC) of the network read operation, the execution cost (0, N, 0) of the local sort operation, and the execution cost (0, 0, N*L*WC) of the local write operation to obtain an execution cost (N, N, N*L*(RC+RNC+WC)) as an execution cost for the re-partition execution step;

for the broadcast execution step, using the execution cost (ΣN i *M, 0, ΣN i *M*L*RNC) for the network read operation as an execution cost for the broadcast execution step;

for the broadcast distribution execution step, superposing the execution cost (0, 0, N*L*RC) for the local read operation, the execution cost (N, 0, N*L*RNC) for the network read operation, and the execution cost (0, 0, N*L*WC) for the local write operation to obtain an execution cost (N, 0, N*L*(RC+RNC+WC)) as an execution cost for the broadcast distribution execution step; and

for the sort join execution step or the hash join execution step, using the execution cost (J, 0, 0) for the output operation as an execution cost for the sort join execution step or the hash join execution step.

18. The non-transitory computer-readable storage medium according to claim 17 , wherein prior to, for the repartition execution step, superposing the execution cost (0, 0, N*L*RC) of the local read operation, the execution cost (N, 0, N*L*RNC) of the network read operation, the execution cost (0, N, 0) of the local sort operation, and the execution cost (0, 0, N*L*WC) of the local write operation to obtain an execution cost (N, N, N*L*(RC+RNC+WC)) as an execution cost for the re-partition execution step:

determining whether a skewed distribution occurs with respect to data records contained in the plurality of to-be-joined data tables; and

if a determination result is positive, then correcting the execution cost (N, 0, N*L*RNC) for the network read operation to (N, 0, P*N*L*p*RNC) and correcting the execution cost (0, 0, N*L*WC) for the local write operation to (0, 0, P*N*L*p*WC), where

p represents a distribution skewness; and

P represents a number of computing nodes for performing a join processing on the to-be-joined data tables.

19. The non-transitory computer-readable storage medium according to claim 18 , wherein prior to, for the hash join execution step, using the execution cost (J, 0, 0) for the output operation as the execution cost for the hash join execution step:

determining whether a data table having a number of data records greater than the size D of the data block supported by each of the plurality of storage nodes exists in each auxiliary data table; and

if a determination result is positive, then correcting the execution cost (J, 0, 0) for the output operation to obtain a corrected execution cost (J, N k *ΣN l , N k *ΣN l *L*WC) as the execution cost for the hash join execution step,

where N l represent an l-st data table having a number of data records greater than the size D of the data block supported by each of the pluraity of storage node; and l=1 . . . n and l≠k.

20. A system for joining data tables, the system comprising:

a memory; and

a processor coupled to the memory, the processor to:

determine a plurality of to-be-joined data tables in a distributed data warehouse that has a plurality of computing nodes and a plurality of storage nodes;

set, based on an environment of the distributed data warehouse where the plurality of to-be-joined data tables are located, a parameter list for cost estimation

obtain a plurality of table joining algorithms;

estimate a plurality of execution costs for the plurality of table joining algorithms such that each of the plurality of table joining algorithm has an estimated execution cost to join the plurality of to-be-joined data tables, wherein the estimation comprises;

determine, for each of the plurality of table joining algorithms, a number of execution steps and one or more key operations of a plurality of operations in each of the execution steps; and

estimate an execution cost for each of the execution steps based on:

computing execution costs of the one or more key operations in each of the execution steps based on parameters selected from the parameter list for the one or more key operations; and

determining the execution cost for each of the execution steps based on at least one of: superposing the execution costs of the one or more key operations, and directly using an execution cost of a selected key operation out of the one or more key operations;

select a target algorithm from the plurality of table joining algorithms based on the estimated execution cost of each of the plurality of table joining algorithms; and

join the plurality of to-be-joined data tables with the target algorithm.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 21, 2026
From: ALIBABA GROUP HOLDING LIMITED
To: CLOUD INTELLIGENCE ASSETS HOLDING (SINGAPORE) PRIVATE LIMITED
Reel/Frame 075478/0225 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 15, 2018
From: XU, DONG; SUN, WEIGUANG; LIAN, JIEHONG; WANG, LONGZHONG
To: ALIBABA GROUP HOLDING LIMITED
Reel/Frame 047514/0392 →
Priority Claims (1)
CN 201610141198.4 · Mar 14, 2016 · national
Continuity (1)
Related Publication 20190171639A1 · Jun 6, 2019