IP Library Granted Patent US 11,734,303
Granted Patent B2
US 11,734,303 · App. 16/805,632 · Granted Aug 22, 2023

Query processing distribution

Inventors: Thierry Cruanes (San Mateo, CA); Benoit Dageville (Foster City, CA); Marcin Zukowski (San Mateo, CA)
Assignee: Snowflake Inc.
G06F16/27A61F5/566G06F9/4881G06F9/5016G06F9/5044G06F9/5083G06F9/5088G06F16/148G06F16/1827G06F16/211G06F16/221G06F16/2365G06F16/2456G06F16/2471G06F16/24532G06F16/24545G06F16/24552G06F16/254G06F16/283G06F16/951G06F16/9535H04L67/1095H04L67/1097H04L67/568
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,734,303
App. No.
16/805,632
Granted
Aug 22, 2023
Kind
B2
Abstract

Example caching systems and methods are described. In one implementation, a method identifies multiple files used to process a query and distributes each of the multiple files to a particular execution node to execute the query. Each execution node determines whether the distributed file is stored in the execution node's cache. If the execution node determines that the file is stored in the cache, it processes the query using the cached file. If the file is not stored in the cache, the execution node retrieves the file from a remote storage device, stores the file in the execution node's cache, and processes the query using the file.

Claims (56)

1. A non-transitory computer-readable medium storing instructions which, when executed by one or more second processors of a computing device, cause the one or more second processors to:

receive a plurality of queries for information stored in one or more databases to be processed by a plurality of virtual warehouses that includes a first plurality of processors, wherein each of the plurality of virtual warehouses includes multiple ones of the first plurality of processors and corresponding local storage;

identify a plurality of query tasks associated with the plurality of queries;

distribute, with the one or more second processors, the plurality of query tasks to a plurality of virtual warehouses using at least a set of statistics associated with the first plurality of processors, wherein the set of statistics is accumulated from at least one of the previously executed query tasks by at least one of the first plurality of processors, the set of statistics is stored separately from a local storage of each of the first plurality of processors and each of first plurality of processors is stateless with respect to the set of statistics;

change, during the processing of the plurality of the query tasks by the plurality of virtual warehouses, a total number of virtual warehouses in the plurality of virtual warehouses is using a load of the plurality of virtual warehouses due to the processing of the plurality of query tasks by the plurality of virtual warehouses, wherein the change is at least one of creating a new virtual warehouse or deleting an existing virtual warehouse and the change in the total number of virtual warehouses are made in a unit of the multiple ones of the first plurality of processors and the corresponding local storage, wherein cache resources associated with the new virtual warehouse are populated with one or more data files associated with the processing of the plurality of query tasks at the time the virtual warehouse is created and wherein the cache resources vary among the plurality of processors, wherein a first subset of the plurality of processors comprises minimal cache resources and a second subset of the plurality of processors comprises cache resources providing faster input-output operations;

execute the plurality of query tasks with the first plurality of processors;

as a result of the change in the total number of virtual warehouses, redistribute the one or more query tasks of the plurality of query tasks to the plurality of virtual warehouses, the one or more query tasks to be executed by the one or more associated processors of a second plurality of processors; and

in response to executing the plurality of query tasks, update at least a portion of the set of statistics based on at least one of the executed query tasks with the plurality of processors.

2. The non-transitory machine-readable medium of claim 1 , wherein the executing of the plurality of query tasks uses at least a plurality of database tables.

3. The non-transitory machine-readable medium of claim 2 , wherein at least some of the plurality of database tables are encrypted and are subsequently decrypted before the executing of the plurality of query tasks.

4. The non-transitory machine-readable medium of claim 2 , wherein at least some of the plurality of database tables are compressed and are subsequently decompressed before the executing of the plurality of query tasks.

5. The non-transitory machine-readable medium of claim 2 , wherein each of the first plurality of processors processes a corresponding one of the plurality of database tables, and wherein data from the plurality of database tables is stored in a cache associated with that processor.

6. The non-transitory machine-readable medium of claim 1 , wherein the query is received from a client, and wherein the instructions further cause the one or more second processors to:

generate a result from the execution of the plurality of query tasks, and data from the result to the client associated with the query.

7. The non-transitory machine-readable medium of claim 1 , wherein the instructions further cause the computing device to:

optimize the plurality of queries.

8. The non-transitory computer-readable medium of claim 1 , wherein the set of data is stored in a relational database.

9. The non-transitory computer-readable medium of claim 8 , wherein the relational database is a structured query language database.

10. The non-transitory computer-readable medium of claim 1 , wherein the database system is a multi-tenant database that isolates computing resources and data between different customers.

11. The non-transitory computer-readable medium of claim 1 , wherein at least one of the databases is external to a system that includes the first plurality of processors.

12. The non-transitory computer-readable medium of claim 1 , wherein the statistics comprise metadata related to the one or more databases.

13. The non-transitory computer-readable medium of claim 1 , wherein the accumulation of the set of statistics from the at least one of the previously executed query tasks occurs automatically.

14. The non-transitory computer-readable medium of claim 1 , wherein the updating of the at least the portion of the set of statistics occurs automatically.

15. A method comprising:

receiving a plurality of queries with a database system to be processed by a plurality of virtual warehouses that includes a first plurality of processors, wherein each of the plurality of virtual warehouses includes multiple ones of the first plurality of processors and corresponding local storage;

identifying a plurality of query tasks associated with the plurality of queries;

distributing, with the one or more second processors, the plurality of query tasks to a plurality of virtual warehouses using at least a set of statistics associated with the first plurality of processors, wherein the set of statistics is accumulated from at least one of the previously executed query tasks by at least one of the first plurality of processors, the set of statistics is stored separately from a local storage of each of the first plurality of processors and each of first plurality of processors is stateless with respect to the set of statistics;

changing, during the processing of the plurality of the query tasks by the plurality of virtual warehouses, a total number of virtual warehouses in the plurality of virtual warehouses is using a load of the plurality of virtual warehouses due to the processing of the plurality of query tasks by the plurality of virtual warehouses, wherein the change is creating a new virtual warehouse and the change in the total number of virtual warehouses are made in a unit of the multiple ones of the first plurality of processors and the corresponding local storage, wherein cache resources associated with the new virtual warehouse are populated with data files associated with the processing of the plurality of query tasks at the time the virtual warehouse is created, based at least in part on the expected tasks to be performed by the new virtual warehouse and wherein the cache resources vary among the plurality of processors, wherein a first subset of the plurality of processors comprises minimal cache resources and a second subset of the plurality of processors comprises cache resources providing faster input-output operations;

executing the plurality of query tasks with the first plurality of processors;

as a result of the change in the total number of virtual warehouses, redistributing the one or more query tasks of the plurality of query tasks to the plurality of virtual warehouses, the one or more query tasks to be executed by the one or more associated processors of a second plurality of processors; and

in response to executing the plurality of query tasks, updating at least a portion of the set of statistics based on at least one of the executed query tasks with the first plurality of processors.

16. The method of claim 15 , wherein the executing of the plurality of query tasks uses at least a plurality of database tables.

17. The method of claim 16 , wherein at least some of the plurality of database tables are encrypted and are subsequently decrypted before the executing of the plurality of query tasks.

18. The method of claim 16 , wherein at least some of the plurality of database tables are compressed and are subsequently decompressed before the executing of the plurality of query tasks.

19. The method of claim 16 , wherein each of the first plurality of processors processes a corresponding one of the plurality of database tables, and wherein data from the plurality of database tables is stored in a cache associated with those processors.

20. The method of claim 15 , wherein the query is received from a client, and further comprising:

returning a result to a query coordinator.

21. The method of claim 15 , further comprising:

optimizing the plurality of queries.

22. The method of claim 15 , wherein the set of data is stored in a relational database.

23. The method of claim 15 , wherein the relational database is a structured query language database.

24. The method of claim 15 , wherein the database system is a multi-tenant database that isolates computing resources and data between different customers.

25. The method of claim 15 , wherein at least one of the databases is external to a system that includes the first plurality of processors.

26. The method of claim 15 , wherein the statistics comprise metadata related to the one or more databases.

27. The method of claim 15 , wherein the accumulation of the set of statistics from the at least one of the previously executed query tasks occurs automatically.

28. The method of claim 15 , wherein the updating of the at least the portion of the set of statistics occurs automatically.

29. A system comprising:

one or more first processors programmed to execute a query coordinator process, the query coordinator process to:

receive a plurality of queries for information stored in one or more databases to be processed by a plurality of virtual warehouses that includes a first plurality of processors, wherein each of the plurality of virtual warehouses includes multiple ones of the first plurality of processors and corresponding local storage,

identify a plurality of query tasks associated with the plurality of queries,

distribute the plurality of query tasks to a plurality of virtual warehouses using at least a set of statistics associated with the first plurality of processors, wherein the set of statistics is accumulated from at least one of the previously executed query tasks by at least one of the first plurality of processors, the set of statistics is stored separately from a local storage of each of the first plurality of processors and each of first plurality of processors is stateless with respect to the set of statistics, and

change, during the processing of the plurality of the query tasks by the plurality of virtual warehouses, a total number of virtual warehouses in the plurality of virtual warehouses is using a load of the plurality of virtual warehouses due to the processing of the plurality of query tasks by the plurality of virtual warehouses, wherein the change is at least one of creating a new virtual warehouse or deleting an existing virtual warehouse and the change in the total number of virtual warehouses are made in a unit of the multiple ones of the first plurality of processors and the corresponding local storage, wherein cache resources associated with the new virtual warehouse are populated with data files associated with the processing of the plurality of query tasks at the time the virtual warehouse is created, based at least in part on the expected tasks to be performed by the new virtual warehouse and wherein the cache resources vary among the plurality of processors, wherein a first subset of the plurality of processors comprises minimal cache resources and a second subset of the plurality of processors comprises cache resources providing faster input-output operations; and

the plurality of second processors programmed to:

execute the plurality of query tasks, wherein the query coordinator process is further programmed to, in response to executing the plurality of query tasks, update at least a portion of the set of statistics based on at least one of the executed query tasks with the plurality of processors; and

as a result of the change in the total number of virtual warehouses, redistribute the one or more query tasks of the plurality of query tasks to the plurality of virtual warehouses, the one or more query tasks to be executed by the one or more associated processors of a second plurality of processors.

30. The system of claim 29 , wherein the executing of the plurality of query tasks uses at least a plurality of database tables.

Assignments (2)
CHANGE OF NAME Recorded May 25, 2021
From: SNOWFLAKE COMPUTING, INC.
To: SNOWFLAKE INC.
Reel/Frame 056367/0145 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 5, 2020
From: DAGEVILLE, BERNOIT; CRUANES, THIERRY; ZUKOWSKI, MARCIN
To: SNOWFLAKE COMPUTING, INC.
Reel/Frame 052033/0014 →
Continuity (3)
Continuation 14518971 · Oct 20, 2014
Provisional Application 61941986 · Feb 19, 2014
Related Publication 20200201880A1 · Jun 25, 2020
Cited By (12)
US 12,189,623 US 12,235,962 US 12,244,626 US 12,259,967 US 12,363,151 US 12,418,565 US 12,432,253 US 12,450,351 US 12,452,273 US 12,468,810 US 12,579,268 US 12,664,258