1 of 49

Data Engineering

Performance Tuning: Query Plan Selection

1

2 of 49

Given a SQL query Q

  • How should Q be executed by the data system?
  • What does the system know?
    • It knows Q, therefore the base relations/views, along with further intended processing of these relations
    • It knows any indexes on the relations, if the relation is clustered, statistics about the relations (sizes, # of columns and distinct values)…

  • Goal: a query execution plan, or a query plan for short

2

3 of 49

Query Execution Plans

Consider

SELECT id, age, zipcode

FROM Stops, Zips

WHERE Stops.location = Zips.location AND age > 35

Stops and Zips are both laid out as blocks

Q: What are ways in which we could execute this query?

4 of 49

Query Execution Plans: The Logical

SELECT id, age, zipcode

FROM Stops, Zips

WHERE Stops.location = Zips.location AND age > 35

 

 

 

 

 

 

 

 

 

 

 

 

 

 

 

5 of 49

Query Execution Plans: The Physical

SELECT id, age, zipcode

FROM Stops, Zips

WHERE Stops.location = Zips.location AND age > 35

 

 

 

 

 

 

 

 

 

 

 

 

 

 

 

What join algorithm do we use?

Do we use an index on age or not?

Do we wait until the results of the join are ready before we start doing projection? Or can we do it “on-the-fly”?

6 of 49

Given a SQL query Q

  • How should Q be executed by the data system?

  • Goal: select a query execution plan:
    • We think of plans at two levels
      • The logical level: the logical operators
      • The physical level: the operator implementations

6

 

7 of 49

Logical Operators vs. Physical Operators

  • Logical operators are extended relational algebra (RA) operators
    • Describe “what” is done
    • e.g., union, select, grouping, project

  • Physical operators describe implementations of these operators
    • Describe “how” to do it
    • e.g., for join
      • nested-loop, sort-merge, hash join, …
    • Physical operators also pertain to non-RA operators such as scanning a table

7

8 of 49

Side Note: Code-Centric vs. Query-Centric

  • In query-centric data systems, the system is responsible for both the logical and physical choices
  • In code-centric data systems, both options are possible:
    • The user supplies the logical query plan, but the system selects the physical implementation
      • E.g., in dataframes, users select the sequence of operators, but leave the physical implementations to the system
    • The user supplies both the logical query plan as well as some parts of the physical implementation
      • E.g., in Map-Reduce systems like Hadoop, as well as in certain NoSQL systems that support a MR interface (MongoDB)
  • Here, even more valuable to understand how to pick good query plans
  • Query plans also appear in many other specialized data systems …

9 of 49

A Query Plan By Any Other Name…

Streaming

NoSQL Pipelines

ML Computation Graphs

10 of 49

Recap: Why Users Should Worry About Performance Even in Query Centric Systems

Multiple reasons why:

Reason I: revisiting our specification

Reason 2: contract breakdown

Reason 3: many heavy-weight levers still under the control of users

The Declarative Contract

11 of 49

Getting started: Scanning a Relation R

  • Either all of R or parts of it that satisfy some condition

  • Simple Table Scan: Read all blocks on disk that contain tuples of R one-by-one

  • Index Scan: Read the index, and use it to read relevant blocks of R
    • More valuable if only a small fraction of R is “relevant”

11

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

1

2

3

4

5

12 of 49

Let’s Talk About Joins

  • Joins (theta or natural) are expensive
    • At a high level, need to match information across multiple relations
  • Many many different ways to do joins
  • We’ll talk about 3 different ways…. to show you how complex it is!

  • We won’t talk about how to pick between these ways
    • Query-centric systems will estimate “cost” for each way and pick one with lowest cost

  • Instead, goal is to give you a vocabulary and intuition for various algorithms

12

13 of 49

Join Approach 1: Nested Loop Joins

Simple Approach:

  • For every tuple of R
    • For every tuple of S
      • If they match, add to output!

  • Q: How would this work when you are accessing tuples at the granularity of blocks?
  • A: bring in k and k’ blocks each of R and S and join all tuples contained within!

13

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

1

2

3

4

5

6

7

8

9

10

11

12

R

S

14 of 49

Join Approach 1: Nested Loop Joins

  • For every k blocks of R
    • For every k’ blocks of S
      • ”Match” tuples across pairs of these blocks; if so, add to output

  • Variants:
    • Index-nested loop uses an index in the “inner” loop to look up only blocks of S that can match the k blocks of R
    • If one of the relations (say S) fits entirely in memory, we only need to cycle through the blocks of R in the “outer” loop

14

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

1

2

3

4

5

6

7

8

9

10

11

12

1

2

3

4

1

2

 

3

4

5

6

15

16

5

6

7

8

1

2

3

4

R

S

4

2

15 of 49

A Brief Recap: Merge Sort

  • Sorting an array by repeatedly merging two smaller, pre-sorted arrays
    • Walk down the pre-sorted sub-arrays during the merge
  • Can be extended to units of “blocks”

16 of 49

Join Approach 2: Sort-Merge Join

  • Phase 1: Sort
    • Sort portions of R on join attrib, write out sorted runs of blocks
      • [analogous to sub-arrays]
    • Sort portions of S on join attrib, write out sorted runs of blocks
  • Phase 2: Merge
    • Merge and match tuples across runs by walking down the runs in sorted order

  • Additional benefit: output is sorted

16

1

2

3

4

5

6

7

8

1

2

3

4

5

6

7

8

1

2

3

4

11

12

13

14

11

12

13

14

5

6

7

8

21

22

23

24

21

22

23

24

1

2

3

4

31

32

33

34

31

32

33

34

5

6

7

8

41

42

43

44

41

42

43

44

11

21

31

41

1

1

2

2

22

3

3

32

4

4

42

23

5

Sort

Merge

17 of 49

Join Approach 2: Sort-Merge Join

  • Phase 1: Sort
    • Sort portions of R on join attrib, write out sorted runs of blocks
      • [analogous to sub-arrays]
    • Sort portions of S on join attrib, write out sorted runs of blocks
  • Phase 2: Merge
    • Merge and match tuples across runs by walking down the runs in sorted order

17

  • Variants:
    • Even more convenient when one of the relations is sorted already
      • Can skip sorting that relation
      • An index (e.g., B+ tree) can be used to retrieve R or S in sorted order
    • Like nested-loop, can be easier if one of the relations fits in memory

18 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

18

1

2

3

4

5

6

7

8

1

2

3

4

5

6

11

12

13

14

1

Start of hashing R phase 1

S

R

19 of 49

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

19

1

2

3

4

5

6

7

8

1

2

3

4

5

6

11

12

13

14

1

21

Join Approach 3: Hash Join

20 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

20

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

14

2

21

11

21 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

21

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

14

3

21

11

22 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

22

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

14

4

21

11

23 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

23

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

14

4

21

11

23

24 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

24

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

8

21

11

22

14

22

End of hashing R phase 1

25 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

25

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

21

11

22

14

22

1

32

24

31

23

Start of hashing S phase 1

26 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

26

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

21

11

22

14

32

31

23

32

24

34

44

41

6

End of hashing S phase 1

27 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

27

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

21

11

22

14

32

31

23

32

24

34

44

41

21

11

31

41

1

Joining red buckets phase 2

28 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

28

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

21

11

22

14

32

31

23

32

24

34

44

41

21

11

31

41

1

2

3

4

5

Joining red buckets phase 2

29 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

29

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

21

11

22

14

32

31

23

32

24

34

44

41

1

2

3

4

5

13

23

6

Joining dark g. buckets phase 2

30 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

30

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

21

11

22

14

32

31

23

32

24

34

44

41

1

2

3

4

5

13

23

6

7

8

Joining dark g. buckets phase 2

31 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat

31

1

2

3

4

5

6

7

8

1

2

3

4

5

6

12

13

21

11

22

14

32

31

23

32

24

34

44

41

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

32 of 49

Join Approach 3: Hash Join

  • Phase 1:
    • Hash R into buckets b1, b2, … based on join attrib
    • Hash S into buckets b1, b2, … based on join attrib
  • Phase 2:
    • Read all tuples hashed to b1 from both R and S at a time, perform join, repeat
  • Variant:
    • If one of the relations (say S) fits entirely in memory, can create the hash table for S first, and then cycle through blocks of R, probe the hash table to find all the output tuples

33 of 49

Similar Variants Exist for Other Binary Operators

  • Typically Hashing, Sorting, and Index-Nested based variants

  • Exercise:
    • How would you use hashing to compute A - B?
    • How would you use sorting to compute A - B?
    • How would you use indexing to compute A - B?

33

34 of 49

Other Operators: Filters, Projects, Grouping

  • Filters/Projects can simply be “applied” to the results of an operation; no specific physical operator choices
    • Exception for index-scan, filtering done as part of scanning

  • For Aggregation/Grouping:
    • Goal is to ensure that “groups” of tuples are processed together
    • Similar to binary operators, hashing and sorting are key techniques
    • If an index is present, can use index to retrieve in sorted order (even cheaper!)

  • Q: How would we use hashing to perform an aggregation per group?

34

35 of 49

Overall: Physical Design of Operators is Hard!

  • Lots of variations, but most rely on a small set of techniques
    • Sorting, hashing, indexing, nested-looping
    • Some of these techniques produce outputs in sorted order, which may matter for subsequent ops.

  • We haven’t talked about what might work best given a certain setting
    • For that one needs a cost model: taking into account the sizes of the relations, the distributions of values, and buffer sizes

35

36 of 49

Recap: Logical Operators vs. Physical Operators

  • Logical operators are extended relational algebra (RA) operators
    • Describe “what” is done
    • e.g., union, select, grouping, project

  • Physical operators describe implementations of these operators
    • Describe “how” to do it
    • e.g., for join
      • nested-loop, sort-merge, hash join, …
    • Physical operators also pertain to non-RA operators such as scanning a table

36

37 of 49

OK…

  • So we talked about physical implementations. Now how do we use them?
  • Typical process of query optimization:
    • Step 1: convert the SQL query to a logical query plan
      • sometimes complex SQL queries get decomposed into multiple query plans, e.g., CTEs or subqueries
    • Step 2: apply rewriting to find other equivalent logical plans
      • needs rewriting rules
    • Step 3: use cost estimates pick among the logical plans, and the corr. physical plan
    • Step 4: feed the corresponding physical plan to the query processor

37

38 of 49

Step 1: Converting to a logical query plan

  • A logical query plan is simply an (extended) relational algebra expression
  • We already know how to do this

SELECT a1, a2, …, aggs

FROM R1, … Rk

WHERE C

GROUP BY b1, …, bm

HAVING H

38

 

Usually Joins

Extended RA operator for grouping and aggregation

39 of 49

Some subqueries can be rewritten!

SELECT DISTINCT Stops.location FROM Stops

WHERE Stops.location IN (SELECT Zips.location FROM Zips)

Q: Can we rewrite without using a subquery?

SELECT DISTINCT Stops.location FROM Stops, Zips

WHERE Stops.location = Zips.location

Q: What if we dropped the DISTINCT keyword?

39

40 of 49

Step 2: Rewriting the Logical Plan

SELECT id, age, zipcode

FROM Stops, Zips

WHERE Stops.location = Zips.location AND age > 35

 

 

 

 

 

 

 

 

 

 

 

 

 

 

 

Need to be able to tell that these three plans are equivalent

41 of 49

Step 2: Rewriting the Logical Plan

  • We need algebraic laws that allow us to manipulate relational algebra expressions

  • Commutative, associative, and distributive laws, like:

41

 

 

 

42 of 49

Laws Involving Selection: Examples

42

 

 

 

 

 

43 of 49

Laws Involving Selection: Examples

  •  

43

 

 

44 of 49

Laws Involving Projection

  •  

44

 

 

 

 

45 of 49

How to use these rules

Product (maker, price, pname, category)

Company (name, city, owner, marketcap)

Query plans (RA exps) also depicted as trees

Q: What does this query evaluate to?

Q: Can we push the predicates down?

45

Product

Company

 

 

 

maker = name

price>100 &

city = “berkeley”

pname

46 of 49

How to use these rules

  • Product (maker, price, pname, category)
  • Company (name, city, owner, marketcap)

46

Product

Company

 

 

maker = name

pname

 

price>100

 

city = “berkeley”

Product

Company

 

 

 

maker = name

price>100 &

city = “berkeley”

pname

47 of 49

How to use these rules

  • Product (maker, price, pname, category)
  • Company (name, city, owner, marketcap)

47

Product

Company

 

 

maker = name

pname

 

price>100

 

city = “berkeley”

Q: can we push

projections down?

48 of 49

How to use these rules

  • Product (maker, price, pname, category)
  • Company (name, city, owner, marketcap)

48

Product

Company

 

 

maker = name

pname

 

price>100

 

city = “berkeley”

 

pname,price,

maker

 

name,city

Product

Company

 

 

maker = name

pname

 

price>100

 

city = “berkeley”

49 of 49

OK…

  • So we talked about physical implementations. Now how do we use them?
  • Let’s talk about the heart of the database engine
  • Step 1: convert the SQL query to a logical query plan
  • Step 2: apply rewriting to find other equivalent logical plans
  • Step 3: use “optimization” pick among the logical plans, and pick the corresponding physical plan
  • Step 4: feed the corresponding plan to the query processor

49