IP Library Granted Patent US 12,189,655
Granted Patent B2
US 12,189,655 · App. 18/539,079 · Granted Jan 7, 2025

Adaptive distribution method for hash operations

Inventors: Benoit Dageville (San Carlos, CA); Thierry Cruanes (San Mateo, CA); Marcin Zukowski (San Mateo, CA); Allison Waingold Lee (San Carlos, CA); Philipp Thomas Unterbrunner (Belmont, CA)
Assignee: Snowflake Inc.
G06F16/273A61F5/566G06F9/4881G06F9/5016G06F9/5044G06F9/5083G06F9/5088G06F16/148G06F16/1827G06F16/211G06F16/221G06F16/2365G06F16/24532G06F16/24545G06F16/24552G06F16/2456G06F16/2471G06F16/254G06F16/27G06F16/283G06F16/951G06F16/9535G06F16/9538H04L67/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 12,189,655
App. No.
18/539,079
Granted
Jan 7, 2025
Kind
B2
Abstract

A method, apparatus, and system for join operations of a plurality of relations that are distributed over a plurality of storage locations over a network of computing components.

Claims (64)

1. A method, comprising:

receiving a relational join query comprising a join operation and an indication of a first relation and a second relation to be joined, wherein the first relation and the second relation are partitioned over processing nodes of a cluster;

re-partitioning the first relation to a plurality of build operator instances;

determining, by a processing device, whether to replicate the first relation to each of a plurality of probe operator instances or to re-partition the second relation between the plurality of probe operator instances, wherein the determining is based at least in part on an actual size of the first relation and an estimated size of the second relation;

based on the determining, distributing the first relation or the second relation over a plurality of communication links of a data communication network to the processing nodes of the cluster associated with the probe operation; and

performing, at the processing nodes of the cluster associated with the probe operation, the relational join query using at least one of a hash join, a sort-merge join, or a nested-loop join to generate a third relation.

2. The method of claim 1 , wherein the determining is further based on a cost metric.

3. The method of claim 1 , wherein the actual size of the first relation is determined during execution of the join operation.

4. The method of claim 1 , further comprising building a hash index for the first relation, wherein the actual size of the first relation is determined based on the building of the hash index.

5. The method of claim 1 , wherein the first relation comprises tuples, the method further comprising:

forwarding the tuples to a plurality of build operator instances in accordance with a partition move; and

determining a total number of tuples processed for the partition move;

wherein the actual size of the first relation is based on the total number of tuples processed for the partition move.

6. The method of claim 1 , further comprising, upon determining to replicate the first relation to each of the plurality of probe operator instances:

setting links between the plurality of build operator instances and the plurality of probe operator instances to broadcast links;

setting links between the second relation and the plurality of probe operator instances to synchronous links; and

sending the first relation through the broadcast links so that each partition of the first relation is broadcasted to every one of the plurality of probe operator instances.

7. The method of claim 1 , further comprising, upon determining to re-partition the second relation between the plurality of probe operator instances:

setting links between the plurality of build operator instances and the plurality of probe operator instances to synchronous links;

setting links between the second relation and the plurality of probe operator instances to partition links; and

sending the second relation through the partition links to the plurality of probe operator instances.

8. A system, comprising:

a memory to store a plurality of relations; and

a processing device operatively coupled with the memory, the processing device to:

receive a relational join query comprising a join operation and an indication of a first relation and a second relation to be joined, wherein the first relation and the second relation are partitioned over processing nodes of a cluster;

re-partition the first relation to a plurality of build operator instances;

determine, based at least in part on an actual size of the first relation and an estimated size of the second relation, whether to replicate the first relation to each of a plurality of probe operator instances or to re-partition the second relation between the plurality of probe operator instances;

based on the determination, distribute the first relation or the second relation over a plurality of communication links of a data communication network to the processing nodes of the cluster associated with the probe operation; and

perform, at the processing nodes of the cluster associated with the probe operation, the relational join query via at least one of a hash join, a sort-merge join, or a nested-loop join to generate a third relation.

9. The system of claim 8 , wherein the determination is further based on a cost metric.

10. The system of claim 8 , wherein the actual size of the first relation is determined during execution of the join operation.

11. The system of claim 8 , wherein the processing device is further to build a hash index for the first relation, wherein the actual size of the first relation is determined based on the build of the hash index.

12. The system of claim 8 , wherein the first relation comprises tuples, and the processing device is further to:

forward the tuples to the plurality of build operator instances in accordance with a partition move; and

determine a total number of tuples processed for the partition move;

wherein the actual size of the first relation is based on the total number of tuples processed for the partition move.

13. The system of claim 8 , wherein if the processing device determines to replicate the first relation to each of the plurality of probe operator instances, the processing device is further to:

set links between the plurality of build operator instances and the plurality of probe operator instances to broadcast links;

set links between the second relation and the plurality of probe operator instances to synchronous links; and

send the first relation through the broadcast links so that each partition of the first relation is broadcasted to every one of the plurality of probe operator instances.

14. The system of claim 8 , wherein if the processing device determines to re-partition the second relation between the plurality of probe operator instances, the processing device is further to:

set links between the plurality of build operator instances and the plurality of probe operator instances to synchronous links;

set links between the second relation and the plurality of probe operator instances to partition links; and

send the second relation through the partition links to the plurality of probe operator instances.

15. A non-transitory computer readable medium having instructions stored thereon that, when executed by a processing device, cause the processing device to:

receive a relational join query comprising a join operation and an indication of a first relation and a second relation to be joined, wherein the first relation and the second relation are partitioned over processing nodes of a cluster;

re-partition the first relation to a plurality of build operator instances;

determine, based at least in part on an actual size of the first relation and an estimated size of the second relation, whether to replicate the first relation to each of a plurality of probe operator instances or to re-partition the second relation between the plurality of probe operator instances;

based on the determination, distribute the first relation or the second relation over a plurality of communication links of a data communication network to the processing nodes of the cluster associated with the probe operation; and

perform, at the processing nodes of the cluster associated with the probe operation, the relational join query via at least one of a hash join, a sort-merge join, or a nested-loop join to generate a third relation.

16. The non-transitory computer readable medium of claim 15 , wherein the determination is further based on a cost metric.

17. The non-transitory computer readable medium of claim 15 , wherein the processing device is further to build a hash index for the first relation, wherein the actual size of the first relation is determined based on the build of the hash index.

18. The non-transitory computer readable medium of claim 15 , wherein the first relation comprises tuples, and the processing device is further to:

forward the tuples to the plurality of build operator instances in accordance with a partition move; and

determine a total number of tuples processed for the partition move;

wherein the actual size of the first relation is based on the total number of tuples processed for the partition move.

19. The non-transitory computer readable medium of claim 15 , wherein if the processing device determines to replicate the first relation to each of the plurality of probe operator instances, the processing device is further to:

set links between the plurality of build operator instances and the plurality of probe operator instances to broadcast links;

set links between the second relation and the plurality of probe operator instances to synchronous links; and

send the first relation through the broadcast links so that each partition of the first relation is broadcasted to every one of the plurality of probe operator instances.

20. The non-transitory computer readable medium of claim 15 , wherein if the processing device determines to re-partition the second relation between the plurality of probe operator instances, the processing device is further to:

set links between the plurality of build operator instances and the plurality of probe operator instances to synchronous links;

set links between the second relation and the plurality of probe operator instances to partition links; and

send the second relation through the partition links to the plurality of probe operator instances.

Assignments (3)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 9, 2024
From: SNOWFLAKE COMPUTING, INC.
To: SNOWFLAKE INC.
Reel/Frame 069546/0861 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Nov 14, 2024
From: DAGEVILLE, BENOIT; CRUANES, THIERRY; ZUKOWKSI, MARCIN; LEE, ALLISON WAINGOLD; UNTERBRUNNER, PHILIPP THOMAS
To: SNOWFLAKE COMPUTING INC.
Reel/Frame 069380/0717 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 14, 2023
From: DAGEVILLE, BENOIT; CRUANES, THIERRY; LEE, ALLISON WAINGOLD; UNTERBRUNNER, PHILLIPP THOMAS; ZUKOWSKI, MARCIN
To: SNOWFLAKE INC.
Reel/Frame 066013/0404 →
Continuity (9)
Continuation 18118595 · Mar 7, 2023
Continuation 17655491 · Mar 18, 2022
Continuation 17358988 · Jun 25, 2021
Continuation 17080219 · Oct 26, 2020
Continuation 16858518 · Apr 24, 2020
Continuation 16039710 · Jul 19, 2018
Continuation 14626836 · Feb 19, 2015
Provisional Application 61941986 · Feb 19, 2014
Related Publication 20240111787A1 · Apr 4, 2024
References Cited (133)
US 4967341A · Yamamoto et al. · 1990 [cited by applicant]
US 5325509A · Lautzenheiser · 1994 [cited by applicant]
US 5546571A · Shan · 1996 [cited by examiner]
US 5787466A · Berliner · 1998 [cited by applicant]
US 5812840A · Shwartz · 1998 [cited by applicant]
US 5864842A · Pederson · 1999 [cited by applicant]
US 5873074A · Kashyap et al. · 1999 [cited by applicant]
US 5898176A · Mori · 1999 [cited by examiner]
US 6112198A · Lohman et al. · 2000 [cited by applicant]
US 6185557B1 · Liu · 2001 [cited by applicant]
US 6226639B1 · Lindsay · 2001 [cited by examiner]
US 6301796B1 · Cresson · 2001 [cited by applicant]
US 6374235B1 · Chen et al. · 2002 [cited by applicant]
US 6507835B1 · Amundsen · 2003 [cited by examiner]
US 6618720B1 · On Au · 2003 [cited by applicant]
US 6865567B1 · Oommen · 2005 [cited by examiner]
US 7085769B1 · Luo et al. · 2006 [cited by applicant]
US 7092951B1 · Luo et al. · 2006 [cited by applicant]
US 7092954B1 · Ramesh · 2006 [cited by applicant]
US 7149737B1 · Luo et al. · 2006 [cited by applicant]
US 7478080B2 · Pirahesh et al. · 2009 [cited by applicant]
US 7577667B2 · Hinshaw et al. · 2009 [cited by applicant]
US 7617179B2 · Nica · 2009 [cited by applicant]
US 7634477B2 · Hinshaw · 2009 [cited by applicant]
US 7644083B1 · Sirek · 2010 [cited by applicant]
US 7702610B2 · Zane et al. · 2010 [cited by applicant]
US 7761477B1 · Luo et al. · 2010 [cited by applicant]
US 7882100B2 · Andrei · 2011 [cited by applicant]
US 7925656B2 · Liu · 2011 [cited by applicant]
US 8001109B2 · Lohman et al. · 2011 [cited by applicant]
US 8001110B2 · Hattori · 2011 [cited by applicant]
US 8015180B2 · Hu et al. · 2011 [cited by applicant]
US 8055651B2 · Barsness et al. · 2011 [cited by applicant]
US 8122008B2 · Li · 2012 [cited by examiner]
US 8126870B2 · Chowdhuri et al. · 2012 [cited by applicant]
US 8255388B1 · Luo et al. · 2012 [cited by applicant]
US 8386473B2 · Rugg et al. · 2013 [cited by applicant]
US 8386532B2 · Annapragada · 2013 [cited by applicant]
US 8621145B1 · Kimmel et al. · 2013 [cited by applicant]
US 8640137B1 · Bostic et al. · 2014 [cited by applicant]
US 8732118B1 · Cole et al. · 2014 [cited by applicant]
US 8805818B2 · Zane et al. · 2014 [cited by applicant]
US 8825678B2 · Potapov et al. · 2014 [cited by applicant]
US 8886631B2 · Adabi et al. · 2014 [cited by applicant]
US 8943103B2 · Annapragada · 2015 [cited by applicant]
US 8949834B2 · Olston · 2015 [cited by applicant]
US 9183256B2 · Zane et al. · 2015 [cited by applicant]
US 9256631B2 · Sen et al. · 2016 [cited by applicant]
US 9275110B2 · Pradhan et al. · 2016 [cited by applicant]
US 9292558B2 · Sen et al. · 2016 [cited by applicant]
US 9569493B2 · Gaza et al. · 2017 [cited by applicant]
US 9569494B2 · Gaza et al. · 2017 [cited by applicant]
US 9659046B2 · Sen · 2017 [cited by applicant]
US 9679011B2 · Dixit et al. · 2017 [cited by applicant]
US 9720967B2 · Lee · 2017 [cited by examiner]
US 9779123B2 · Sen et al. · 2017 [cited by applicant]
US 9792328B2 · Cheng et al. · 2017 [cited by applicant]
US 9836505B2 · Cheng · 2017 [cited by examiner]
US 9875280B2 · Amdt et al. · 2018 [cited by applicant]
US 10019481B2 · Jagtap et al. · 2018 [cited by applicant]
US 10120901B2 · Larriba-Pey · 2018 [cited by applicant]
US 10380112B2 · Bodziony et al. · 2019 [cited by applicant]
US 10397317B2 · Balkesen et al. · 2019 [cited by applicant]
US 20030065688A1 · Dageville et al. · 2003 [cited by applicant]
US 20040167904A1 · Wen · 2004 [cited by applicant]
US 20040172400A1 · Zarom et al. · 2004 [cited by applicant]
US 20040220904A1 · Finlay · 2004 [cited by applicant]
US 20040220923A1 · Nica · 2004 [cited by applicant]
US 20050210049A1 · Foster · 2005 [cited by applicant]
US 20060074872A1 · Gordon · 2006 [cited by applicant]
US 20060074901A1 · Pirahesh et al. · 2006 [cited by applicant]
US 20060080285A1 · Chowdhuri · 2006 [cited by applicant]
US 20060117036A1 · Cruanes · 2006 [cited by applicant]
US 20060136354A1 · Bell et al. · 2006 [cited by applicant]
US 20060218123A1 · Chowdhuri et al. · 2006 [cited by applicant]
US 20070073643A1 · Ghosh et al. · 2007 [cited by applicant]
US 20070276861A1 · Pryce et al. · 2007 [cited by applicant]
US 20080027965A1 · Garrett et al. · 2008 [cited by applicant]
US 20080222346A1 · Raciborski et al. · 2008 [cited by applicant]
US 20080288473A1 · Hu et al. · 2008 [cited by applicant]
US 20090228514A1 · Liu · 2009 [cited by applicant]
US 20090300043A1 · MacLennan · 2009 [cited by applicant]
US 20100005054A1 · Smith et al. · 2010 [cited by applicant]
US 20100082599A1 · Graefe · 2010 [cited by applicant]
US 20100082671A1 · Li · 2010 [cited by examiner]
US 20100131490A1 · Lamb et al. · 2010 [cited by applicant]
US 20100205170A1 · Barsness et al. · 2010 [cited by applicant]
US 20110252427A1 · Olston · 2011 [cited by applicant]
US 20110302151A1 · Abadi et al. · 2011 [cited by applicant]
US 20110313999A1 · Bruno · 2011 [cited by examiner]
US 20120047125A1 · Day · 2012 [cited by applicant]
US 20120101860A1 · Ezzat · 2012 [cited by applicant]
US 20120109888A1 · Zhang et al. · 2012 [cited by applicant]
US 20120191699A1 · George et al. · 2012 [cited by applicant]
US 20120254154A1 · Rugg et al. · 2012 [cited by applicant]
US 20120296883A1 · Ganesh et al. · 2012 [cited by applicant]
US 20120310916A1 · Adabi et al. · 2012 [cited by applicant]
US 20120311065A1 · Ananthanarayanan et al. · 2012 [cited by applicant]
US 20120317094A1 · Bear et al. · 2012 [cited by applicant]
US 20120323971A1 · Pasupuleti · 2012 [cited by applicant]
US 20130013585A1 · Graefe · 2013 [cited by applicant]
US 20130110778A1 · Taylor et al. · 2013 [cited by applicant]
US 20130117255A1 · Liu et al. · 2013 [cited by applicant]
US 20130124545A1 · Holmberg · 2013 [cited by applicant]
US 20130124565A1 · Annapragada · 2013 [cited by applicant]
US 20130132967A1 · Soundararajan et al. · 2013 [cited by applicant]
US 20130205092A1 · Roy et al. · 2013 [cited by applicant]
US 20130226902A1 · Lariba-Pey · 2013 [cited by applicant]
US 20130232133A1 · Al-omari et al. · 2013 [cited by applicant]
US 20130318123A1 · Annapragada · 2013 [cited by applicant]
US 20140006380A1 · Arndt et al. · 2014 [cited by applicant]
US 20140250142A1 · Pradhan et al. · 2014 [cited by applicant]
US 20140280023A1 · Jagtap et al. · 2014 [cited by applicant]
US 20140372365A1 · Weyerhaeuser et al. · 2014 [cited by applicant]
US 20150039627A1 · Sen et al. · 2015 [cited by applicant]
US 20150039852A1 · Sen et al. · 2015 [cited by applicant]
US 20150039853A1 · Sen et al. · 2015 [cited by applicant]
US 20150120555A1 · Jung et al. · 2015 [cited by applicant]
US 20150186465A1 · Gaza et al. · 2015 [cited by applicant]
US 20150220600A1 · Bellamkonda · 2015 [cited by applicant]
US 20150261820A1 · Cheng · 2015 [cited by examiner]
US 20150278306A1 · Cheng · 2015 [cited by applicant]
US 20160098451A1 · Dickie · 2016 [cited by applicant]
US 20170185648A1 · Kavulya et al. · 2017 [cited by applicant]
US 20190034486A1 · Bodziony et al. · 2019 [cited by applicant]
US 20190104175A1 · Balkesen et al. · 2019 [cited by applicant]
US 20190303370A1 · Bodziony et al. · 2019 [cited by applicant]
JP 2009015534A · 2009 [cited by applicant]
“Oracle9i Database New Features, Release 2 (9.2)” Mar. 2002, p. 166. [cited by applicant]
Sergey Melnik et al.: “Dremel: Interactive Analysis of WebScale Datasets”, Proceedings of the VLDB Endowment, vol. 3, 2010, Jan. 1, 2010, pp. 330-339. [cited by applicant]
David a Maluf et al: “NASA Technology Transfer System”, Space Mission Challenges for Information Technology (SMC-IT), 2011 IEEE Fourth International Conference on, IEEE, Aug. 2, 2011 (Aug. 2, 2011), pp. 111-117. [cited by applicant]
Hollman J et al: “Empirical observations regarding predictability in user access-behavior in a distributed digital library system”, Parallel and Distributed Processing Symposium, Proceedings International, IPDPS 2002, A… [cited by applicant]
CNIPA—Second Office Action dated Nov. 8, 2023, Chinese Application No. 2019105957089, 10 pages. [cited by applicant]