IP Library › Granted Patent US 12,271,379
Granted Patent B2
US 12,271,379 · App. 18/348,400 · Granted Apr 8, 2025

Cross-database join query

Inventors: Hai Jun Shen (Tianjin, CN); Chang Sheng Liu (Haidian, CN); Jun Hui Liu (Xi'an, CN); Ying Qi Pan (Beijing, CN); Liam Loucks (Calgary, CA)
Assignee: International Business Machines Corporation
G06F16/24544G06F11/3409G06F16/24561G06F16/256
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 12,271,379
App. No.
18/348,400
Filed
Jul 7, 2023
Granted
Apr 8, 2025
Kind
B2
Art Unit
2154
USPC
707/718
Abstract

A data virtualization layer (DV) receives a join query request related to a plurality of tables respectively stored in a plurality of distributed database servers. A plurality of candidate query plans for the join query request is generated where each of the plurality of candidate query plans indicates an order for transmitting the tables respectively stored in the database servers to the DV. For each of the plurality of candidate query plans, a query cost for the candidate query plan is calculated based on a data amount of the tables to be transmitted according to the candidate query plan. From the plurality of candidate query plans, a query plan is determined for the join query request which has a lowest query cost.

Claims (42)

1. A computer-implemented method comprising:

receiving at a data virtualization layer (DV), by one or more processors, a join query request related to a plurality of tables {T 1 , T 2 , . . . , T n } respectively stored in a plurality of distributed database servers {S 1 , S 2 , . . . , S n }, wherein n is an integer larger than 1;

implementing a catalog agent for the data virtualization layer;

mapping tables, using the catalog agent, in the next database server into virtual tables;

generating, by the one or more processors, a plurality of candidate query plans for the join query request, each of the plurality of candidate query plans indicating an order for transmitting the tables {T 1 , T 2 , . . . , T n } respectively stored in the database servers {S 1 , S 2 , . . . , S n } to the DV, wherein for one candidate query plan P={S i →S i → . . . S n →, . . . , →DV}, 1≤i, j≤n, table T i stored in database server S i is transmitted to and stored in database server S j before being transmitted together with table T j which is stored in database server S j to a next database server, and wherein in at least one of the database servers {S 1 , S 2 , . . . , S n }, a join work is performed on the stored tables based on the join query request before being transmitted to the next database server, the join work being performed by the catalog agent;

for each of the plurality of candidate query plans, calculating, by the one or more processors, a query cost for the candidate query plan based on a data amount of the tables to be transmitted according to the candidate query plan; and

determining, by the one or more processors, from the plurality of candidate query plans, a query plan for the join query request which has a lowest query cost.

2. The computer-implemented method of claim 1 , wherein the data amount of the tables to be transmitted is determined based on statistic information of the tables {T 1 , T 2 , . . . , T n }, the statistic information of the tables {T 1 , T 2 , . . . , T n } being respectively obtained from the database servers {S 1 , S 2 , . . . , S n }.

3. The computer-implemented method of claim 1 , wherein the query cost for the candidate query plan is further calculated based on a processing rate of the database servers {S 1 , S 2 , . . . , S n }.

4. The computer-implemented method of claim 3 , wherein the processing rate is determined by performance information of the at least one of the database servers {S 1 , S 2 , . . . , S n } in which the join work is performed, the performance information being obtained from catalog parameters stored in respective database servers {S 1 , S 2 , . . . , S n }.

5. The computer-implemented method of claim 4 , wherein the performance information is indicated by at least one of a CPU status, a memory status, a disk utilization of the at least one of the database servers {S 1 , S 2 , . . . , S n }.

6. The computer-implemented method of claim 1 , wherein the query cost for the candidate query plan is further calculated based on network latencies for transmitting the tables {T 1 , T 2 , . . . , T n } respectively stored in the database servers {S 1 , S 2 , . . . , S n } to the DV based on the candidate query plan.

7. The computer-implemented method of claim 6 , wherein for one database server included in the database servers {S 1 , S 2 , . . . , S n }, the network latencies with other database servers and DV are obtained from catalog parameters stored in the one database server.

8. The computer-implemented method of claim 1 , wherein the method further comprises:

fetching, by the one or more processors, data related to the join query request based on the determined query plan.

9. A system comprising:

one or more processors;

a memory coupled to at least one of the one or more processors;

a set of computer program instructions stored in the memory and executed by at least one of the one or more processors in order to perform actions of:

receiving, at a data virtualization layer (DV), a join query request related to a plurality of tables {T 1 , T 2 , . . . , T n } respectively stored in a plurality of distributed database servers {S 1 , S 2 , . . . , S n }, wherein n is an integer larger than 1;

implementing a catalog agent for the data virtualization layer;

mapping tables, using the catalog agent, in the next database server into virtual tables;

generating a plurality of candidate query plans for the join query request, each of the plurality of candidate query plans indicating an order for transmitting the tables {T 1 , T 2 , . . . , T n } respectively stored in the database servers {S 1 , S 2 , . . . , S n } to the DV, wherein for one candidate query plan P={S i →S j →, . . . , S n →, . . . , →DV}, 1≤i, j≤n, table T i stored in database server S i is transmitted to and stored in database server S j before being transmitted together with table T j which is stored in database server S j to a next database server, and wherein in at least one of the database servers {S 1 , S 2 , . . . , S n }, a join work is performed on the stored tables based on the join query request before being transmitted to the next database server, the join work being performed by the catalog agent;

for each of the plurality of candidate query plans, calculating a query cost for the candidate query plan based on a data amount of the tables to be transmitted according to the candidate query plan; and

determining from the plurality of candidate query plans, a query plan for the join query request which has a lowest query cost.

10. The system of claim 9 , wherein the data amount of the tables to be transmitted is determined based on statistic information of the tables {T 1 , T 2 , . . . , T n }, the statistic information of the tables {T 1 , T 2 , . . . , T n } being respectively obtained from the database servers {S 1 , S 2 , . . . , S n }.

11. The system of claim 9 , wherein the query cost for the candidate query plan is further calculated based on a processing rate of the database servers {S 1 , S 2 , . . . , S n }.

12. The system of claim 11 , wherein the processing rate is determined by performance information of the at least one of the database servers {S 1 , S 2 , . . . , S n } in which the join work is performed, the performance information being obtained from catalog parameters stored in respective database servers {S 1 , S 2 , . . . , S n }.

13. The system of claim 12 , wherein the performance information is indicated by at least one of a CPU status, a memory status, a disk utilization of the at least one of the database servers {S 1 , S 2 , . . . , S n }.

14. The system of claim 9 , wherein the query cost for the candidate query plan is further calculated based on network latencies for transmitting the tables {T 1 , T 2 , . . . , T n } respectively stored in the database servers {S 1 , S 2 , . . . , S n } to the DV based on the candidate query plan.

15. A computer program product comprising a computer readable storage medium having program instructions embodied therewith, wherein the program instructions being executable by a device to perform a method comprising:

receiving, at a data virtualization layer (DV), a join query request related to a plurality of tables {T 1 , T 2 , . . . , T n } respectively stored in a plurality of distributed database servers {S 1 , S 2 , . . . , S n }, wherein n is an integer larger than 1;

implementing a catalog agent for the data virtualization layer;

mapping tables, using the catalog agent, in the next database server into virtual tables;

generating a plurality of candidate query plans for the join query request, each of the plurality of candidate query plans indicating an order for transmitting the tables {T 1 , T 2 , . . . , T n } respectively stored in the database servers {S 1 , S 2 , . . . , S n } to the DV, wherein for one candidate query plan P={S i →S j → . . . S n →, . . . , →DV}, 1≤i, j≤n, table T i stored in database server S i is transmitted to and stored in database server S j before being transmitted together with table T j which is stored in database server S j to a next database server, and wherein in at least one of the database servers {S 1 , S 2 , . . . , S n }, a join work is performed on the stored tables based on the join query request before being transmitted to the next database server, the join work being performed by the catalog agent;

for each of the plurality of candidate query plans, calculating a query cost for the candidate query plan based on a data amount of the tables to be transmitted according to the candidate query plan; and

determining from the plurality of candidate query plans, a query plan for the join query request which has a lowest query cost.

16. The computer program product of claim 15 , wherein the data amount of the tables to be transmitted is determined based on statistic information of the tables {T 1 , T 2 , . . . , T n }, the statistic information of the tables {T 1 , T 2 , . . . , T n } being respectively obtained from the database servers {S 1 , S 2 , . . . , S n }.

17. The computer program product of claim 15 , wherein the query cost for the candidate query plan is further calculated based on a processing rate of the database servers {S 1 , S 2 , . . . , S n }.

18. The computer program product of claim 17 , wherein the processing rate is determined by performance information of the at least one of the database servers {S 1 , S 2 , . . . , S n } in which the join work is performed, the performance information being obtained from catalog parameters stored in respective database servers {S 1 , S 2 , . . . , S n }.

19. The computer program product of claim 18 , wherein the performance information is indicated by at least one of a CPU status, a memory status, a disk utilization of the at least one of the database servers {S 1 , S 2 , . . . , S n }.

20. The computer program product of claim 15 , wherein the query cost for the candidate query plan is further calculated based on network latencies for transmitting the tables {T 1 , T 2 , . . . , T n } respectively stored in the database servers {S 1 , S 2 , . . . , S n } to the DV based on the candidate query plan.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jul 7, 2023
From: SHEN, HAI JUN; LIU, CHANG SHENG; LIU, JUN HUI; PAN, YING QI; LOUCKS, LIAM
To: INTERNATIONAL BUSINESS MACHINES CORPORATION
Reel/Frame 064178/0472 →
Continuity (1)
Related Publication 20250013643A1 · Jan 9, 2025
References Cited (16)
US 7454462B2 · Belfiore · 2008 [cited by applicant]
US 9582221B2 · Du · 2017 [cited by applicant]
US 10592562B2 · Pal · 2020 [cited by applicant]
US 11354312B2 · Liu · 2022 [cited by applicant]
US 20110258179A1 · Weissman · 2011 [cited by applicant]
US 20140188841A1 · Sun · 2014 [cited by applicant]
US 20190303475A1 · Jindal · 2019 [cited by examiner]
US 20200073986A1 · Purcell · 2020 [cited by examiner]
US 20210224275A1 · Maheshwari · 2021 [cited by examiner]
US 20210342351A1 · Mathew · 2021 [cited by examiner]
US 20220300491A1 · Gruszecki · 2022 [cited by examiner]
US 20220374404A1 · Johnson · 2022 [cited by examiner]
Jiang, “Optimizing join query in distributed database”, A Capstone Project (or Thesis) Submitted to the University of North Carolina Wilmington in Partial Fulfillment of the Requirements for the Degree of Master of Scie… [cited by applicant]
Kang, “Global query management in heterogeneous distributed database systems”, Microprocessing and Microprogramming, vol. 38, (1993), pp. 377-384. [cited by applicant]
Karpathiotakis et al., “Fast queries over heterogeneous data through engine customization”, Proceedings of the VLDB Endowment, vol. 9, No. 12, pp. 972-983. [cited by applicant]
Ran Tan et al., “Enabling Query Processing across Heterogeneous Data Models: A Survey”, 2017 IEEE International Conference on Big Data (Big Data), Published Date: Dec. 14, 2017, 10 pages. [cited by applicant]