IP Library Granted Patent US 12,067,058
Granted Patent B2
US 12,067,058 · App. 17/886,383 · Granted Aug 20, 2024

Executing database queries for joining tables using channel based flow control

Inventor: Adam Szymański (Warsaw, PL)
Assignee: Oxla sp. z o.o.
G06F16/90335G06F16/24537G06F16/24542G06F16/24544
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,067,058
App. No.
17/886,383
Granted
Aug 20, 2024
Kind
B2
Abstract

A database system generates an execution plan including multiple operators for processing a database query, for example, a join query or a group by query. The database system allocates a set of threads. Threads communicate with other threads via blocking channels. A blocking channel includes a buffer of a fixed capacity. The database system processes the database query by streaming data through operators of the execution plan. A thread sends data generated by an operator to another thread via the blocking channel if the buffer of the blocking channel has available capacity to store the data, or else the thread blocks until the buffer has capacity to store the data. Similarly, a thread receives data generated by an operator of another thread via the blocking channel if the buffer of the blocking channel has available data, or else the thread blocks until the buffer has data.

Claims (56)

1. A computer-implemented method for executing a database query for performing join of data stored in a plurality of input tables, the computer-implemented method comprising:

receiving, by a database system executing on a cluster of servers, the database query specifying a join of a scan table with a hash table;

generating an execution plan for executing the database query, the execution plan including instructions for a plurality of join strategies;

allocating a set of threads for executing operators of the execution plan, each thread processing one or more operators, the set of threads comprising at least a first thread and a second thread, wherein the first thread communicates with the second thread via a blocking channel comprising a buffer of fixed capacity for storing data; and

processing the database query by streaming data through the operators of the execution plan, the processing comprising:

determining a size of the hash table,

selecting a join strategy from the plurality of join strategies based on the size of the hash table,

executing the selected join strategy comprising, communicating data by the first thread via the blocking channel to the second thread to distribute one of the scan table or the hash table by performing a distributed shuffle, and

performing, by a server, a join of the scan table with a part of the hash table mapped to the server.

2. The computer-implemented method of claim 1 , wherein the blocking channel allows a thread to perform one or more of:

a push operation that pushes data to the buffer of the blocking channel, wherein the push operation causes the thread to block if the buffer of the blocking channel is full; and

a pull operation that pulls data from the buffer of the blocking channel, wherein the pull operation causes the thread to block if the buffer of the blocking channel is empty.

3. The computer-implemented method of claim 1 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table being below a threshold value, the join strategy distributing the hash table across servers of the cluster of servers and performing local scan at each server to join a portion of the scan table with the hash table.

4. The computer-implemented method of claim 1 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table exceeding a threshold value, the join strategy partitioning the hash table across servers of the cluster of servers and performing a distributed join at each server to join a portion of the scan table with the hash table.

5. The computer-implemented method of claim 1 , wherein the blocking channel is introduced between the first thread and the second thread, wherein the first thread executes an operator for broadcasting data of one of the hash table or the scan table to each server of the cluster of servers.

6. The computer-implemented method of claim 1 , wherein the blocking channel is introduced between the first thread and the second thread, wherein a second operator performed by the second thread performs an aggregation of data stored in one of the hash table or scan table.

7. The computer-implemented method of claim 1 , wherein processing the database query further comprises:

receiving, by the second thread, data stored in the buffer of the blocking channel if there is data available in the buffer of the blocking channel, and

blocking execution of the second thread if the buffer of the blocking channel is empty, wherein execution of the second thread is blocked until the buffer of the blocking channel has data available.

8. A non-transitory computer readable storage medium storing instructions that when executed by one or more computer processors, cause the one or more computer processors to perform steps comprising:

receiving, by a database system executing on a cluster of servers, a database query specifying a join of a scan table with a hash table;

generating an execution plan for executing the database query, the execution plan including instructions for a plurality of join strategies;

allocating a set of threads for executing operators of the execution plan, each thread processing one or more operators, the set of threads comprising at least a first thread and a second thread, wherein the first thread communicates with the second thread via a blocking channel comprising a buffer of fixed capacity for storing data; and

processing the database query by streaming data through the operators of the execution plan, the processing comprising:

determining a size of the hash table,

selecting a join strategy from the plurality of join strategies based on the size of the hash table,

executing the selected join strategy comprising, communicating data by the first thread via the blocking channel to the second thread to distribute one of the scan table or the hash table by performing a distributed shuffle, and

performing, by a server, a join of the scan table with a part of the hash table mapped to the server.

9. The non-transitory computer readable storage medium of claim 8 , wherein the blocking channel allows a thread to perform one or more of:

a push operation that pushes data to the buffer of the blocking channel, wherein the push operation causes the thread to block if the buffer of the blocking channel is full; and

a pull operation that pulls data from the buffer of the blocking channel, wherein the pull operation causes the thread to block if the buffer of the blocking channel is empty.

10. The non-transitory computer readable storage medium of claim 8 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table being below a threshold value, the join strategy distributing the hash table across servers of the cluster of servers and performing local scan at each server to join a portion of the scan table with the hash table.

11. The non-transitory computer readable storage medium of claim 8 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table exceeding a threshold value, the join strategy partitioning the hash table across servers of the cluster of servers and performing a distributed join at each server to join a portion of the scan table with the hash table.

12. The non-transitory computer readable storage medium of claim 8 , wherein the blocking channel is introduced between the first thread and the second thread, wherein the first thread executes an operator for broadcasting data of one of the hash table or the scan table to each server of the cluster of servers.

13. The non-transitory computer readable storage medium of claim 8 , wherein the blocking channel is introduced between the first thread and the second thread, wherein a second operator performed by the second thread performs an aggregation of data stored in one of the hash table or scan table.

14. The non-transitory computer readable storage medium of claim 8 , wherein the instructions for processing the database query further cause the one or more computer processors to perform steps comprising:

receiving, by the second thread, data stored in the buffer of the blocking channel if there is data available in the buffer of the blocking channel, and

blocking execution of the second thread if the buffer of the blocking channel is empty, wherein execution of the second thread is blocked until the buffer of the blocking channel has data available.

15. A computer system comprising:

one or more computer processors; and

a non-transitory computer readable storage medium storing instructions that when executed by the one or more computer processors, cause the one or more computer processors to perform steps comprising:

receiving, by a database system executing on a cluster of servers, a database query specifying a join of a scan table with a hash table;

generating an execution plan for executing the database query, the execution plan including instructions for a plurality of join strategies;

allocating a set of threads for executing operators of the execution plan, each thread processing one or more operators, the set of threads comprising at least a first thread and a second thread, wherein the first thread communicates with the second thread via a blocking channel comprising a buffer of fixed capacity for storing data; and

processing the database query by streaming data through the operators of the execution plan, the processing comprising:

determining a size of the hash table,

selecting a join strategy from the plurality of join strategies based on the size of the hash table,

executing the selected join strategy comprising, communicating data by the first thread via the blocking channel to the second thread to distribute one of the scan table or the hash table according to the selected join strategy by performing a distributed shuffle, and

performing, by a server, a join of the scan table with a part of the hash table mapped to the server.

16. The computer system of claim 15 , wherein the blocking channel allows a thread to perform one or more of:

a push operation that pushes data to the buffer of the blocking channel, wherein the push operation causes the thread to block if the buffer of the blocking channel is full; and

a pull operation that pulls data from the buffer of the blocking channel, wherein the pull operation causes the thread to block if the buffer of the blocking channel is empty.

17. The computer system of claim 15 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table being below a threshold value, the join strategy distributing the hash table across servers of the cluster of servers and performing local scan at each server to join a portion of the scan table with the hash table.

18. The computer system of claim 15 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table exceeding a threshold value, the join strategy partitioning the hash table across servers of the cluster of servers and performing a distributed join at each server to join a portion of the scan table with the hash table.

19. The computer system of claim 15 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table being below a threshold value, the join strategy distributing the hash table across servers of the cluster of servers and performing local scan at each server to join a portion of the scan table with the hash table.

20. The computer system of claim 15 , wherein a join strategy is selected from the plurality of join strategies responsive to the size of the hash table exceeding a threshold value, the join strategy partitioning the hash table across servers of the cluster of servers and performing a distributed join at each server to join a portion of the scan table with the hash table.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 15, 2025
From: OXLA SP. Z O.O
To: REDPANDA DATA, INC.
Reel/Frame 073213/0277 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 12, 2022
From: SZYMANSKI, ADAM
To: OXLA SP. Z O.O.
Reel/Frame 062053/0017 →
Priority Claims (1)
PL 441869 · Jul 28, 2022 · national
Continuity (1)
Related Publication 20240037099A1 · Feb 1, 2024