Sublink pushdown

Updated at:

PolarDB for PostgreSQL support sublink pushdown for query rewrite, which improves the execution efficiency of SQL statements that contain IN or ANY clauses.

Background

In PostgreSQL, ANY-type sublinks (IN and ANY clauses) are pulled up as semi joins (semi join). However, if the joined relation is a subquery that cannot be pulled up, PostgreSQL cannot generate a parameterized path for it. The subquery runs as an independent whole, which may significantly degrade the overall SQL execution efficiency when the subquery processes a large amount of data.

For example, the following SQL statement contains a GROUP BY clause in the subquery and therefore cannot be pulled up. Its execution time is dominated by scanning and sorting the t_big table, and increases as the t_big table grows.

EXPLAIN ANALYZE SELECT * FROM (SELECT a, sum(b) b FROM t_big GROUP BY a)v WHERE a IN (SELECT a FROM t_small);
                                                                   QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------------------
 Merge Semi Join  (cost=0.55..59523.15 rows=10000 width=12) (actual time=0.064..1237.621 rows=2 loops=1)
   Merge Cond: (t_big.a = t_small.a)
   ->  GroupAggregate  (cost=0.42..46910.99 rows=1000000 width=12) (actual time=0.033..1113.615 rows=1000000 loops=1)
         Group Key: t_big.a
         ->  Index Scan using t_big_a_idx on t_big  (cost=0.42..31910.99 rows=1000000 width=8) (actual time=0.024..420.575 rows=1000000 loops=1)
   ->  Index Only Scan using t_small_a_idx on t_small  (cost=0.13..12.16 rows=2 width=4) (actual time=0.028..0.030 rows=2 loops=1)
         Heap Fetches: 2
 Planning Time: 0.256 ms
 Execution Time: 1237.700 ms
(9 rows)

If the ANY-type sublink can be pushed down into the subquery, the index on the subquery can be used to improve execution efficiency.

EXPLAIN ANALYZE SELECT * FROM (SELECT a, sum(b) b FROM t_big WHERE a IN (SELECT a FROM t_small) GROUP BY a)v;
                                                             QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------
 GroupAggregate  (cost=17.96..17.99 rows=2 width=12) (actual time=0.061..0.064 rows=2 loops=1)
   Group Key: t_big.a
   ->  Sort  (cost=17.96..17.96 rows=2 width=8) (actual time=0.054..0.056 rows=2 loops=1)
         Sort Key: t_big.a
         Sort Method: quicksort  Memory: 25kB
         ->  Nested Loop  (cost=1.46..17.95 rows=2 width=8) (actual time=0.031..0.045 rows=2 loops=1)
               ->  Unique  (cost=1.03..1.04 rows=2 width=4) (actual time=0.014..0.017 rows=2 loops=1)
                     ->  Sort  (cost=1.03..1.03 rows=2 width=4) (actual time=0.013..0.014 rows=2 loops=1)
                           Sort Key: t_small.a
                           Sort Method: quicksort  Memory: 25kB
                           ->  Seq Scan on t_small  (cost=0.00..1.02 rows=2 width=4) (actual time=0.005..0.006 rows=2 loops=1)
               ->  Index Scan using t_big_a_idx on t_big  (cost=0.42..8.44 rows=1 width=8) (actual time=0.010..0.011 rows=1 loops=2)
                     Index Cond: (a = t_small.a)
 Planning Time: 0.527 ms
 Execution Time: 0.143 ms
(15 rows)

Prerequisites

The following PolarDB for PostgreSQL versions are supported:

  • PostgreSQL 14 (minor engine version 2.0.14.13.28.0 or later)

  • PostgreSQL 11 (minor engine version 2.0.11.15.44.0 or later)

Note

You can view the minor engine version in the console, or run the SHOW polardb_version; statement to check it. If the minor engine version does not meet the requirement, upgrade the minor engine version.

Use cases

Sublink pushdown applies to subqueries that are referenced by an IN or ANY clause and contain a GROUP BY clause, especially when the subquery involves a large table. Pushing the IN or ANY clause into the subquery enables use of the index on the large table and reduces access to its data.

Limitations

Sublink pushdown has the following limitations:

  • The IN or ANY clause must reference a subquery that contains a GROUP BY clause. Otherwise, open source PostgreSQL generates a parameterized path directly and sublink pushdown is not used.

  • The columns in the IN or ANY clause must be included in the GROUP BY columns. Otherwise, the rewritten SQL is not equivalent to the original.

  • The current query block must not contain outer joins. Otherwise, the rewritten SQL is not equivalent to the original.

  • Only single-column scenarios are supported. Examples: a in (select a from t) and a = any(select a from t).

  • Only SELECT and CREATE TABLE AS statements are supported.

Usage notes

Sublink pushdown is controlled by parameters. The related parameters and their descriptions are as follows:

Parameter

Description

polar_cbqt_pushdown_sublink

Controls whether sublink pushdown is enabled. Valid values:

  • OFF (default): disables sublink pushdown.

  • ON: enables sublink pushdown. The Cost-based query transformation (CBQT) framework decides whether to apply pushdown based on cost.

  • FORCE: forcibly enables sublink pushdown and bypasses the CBQT framework. This value is typically used in hints and is not recommended as a global setting.

Examples

Data preparation

CREATE TABLE t_small(a int);
CREATE TABLE t_big(a int, b int, c int);

CREATE INDEX ON t_big(a);

INSERT INTO t_big SELECT i, i, i FROM generate_series(1, 1000000)i;
INSERT INTO t_small VALUES(1), (1000000);

ANALYZE t_small, t_big;

Original query

In the original query plan, the join condition t_big.a = t_small.a cannot be pushed down as a parameterized path, which forces a full table scan on t_big and leads to low execution efficiency.

EXPLAIN ANALYZE SELECT * FROM (SELECT a, sum(b) b FROM t_big GROUP BY a)v WHERE a IN (SELECT a FROM t_small);

The following result is returned:

                                                                   QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------------------
 Merge Semi Join  (cost=1.46..59510.27 rows=10000 width=12) (actual time=0.049..1239.128 rows=2 loops=1)
   Merge Cond: (t_big.a = t_small.a)
   ->  GroupAggregate  (cost=0.42..46909.23 rows=1000000 width=12) (actual time=0.034..1113.324 rows=1000000 loops=1)
         Group Key: t_big.a
         ->  Index Scan using t_big_a_idx on t_big  (cost=0.42..31909.23 rows=1000000 width=8) (actual time=0.025..412.650 rows=1000000 loops=1)
   ->  Sort  (cost=1.03..1.03 rows=2 width=4) (actual time=0.012..0.013 rows=2 loops=1)
         Sort Key: t_small.a
         Sort Method: quicksort  Memory: 25kB
         ->  Seq Scan on t_small  (cost=0.00..1.02 rows=2 width=4) (actual time=0.005..0.006 rows=2 loops=1)
 Planning Time: 0.219 ms
 Execution Time: 1239.208 ms
(11 rows)

Enable sublink pushdown through CBQT

After CBQT and sublink pushdown are enabled, the a in (select a from t_small) clause is pushed down into the subquery. This generates a parameterized path for t_big based on the join condition, significantly reducing the amount of data scanned and the execution time.

Note

The sublink pushdown rewrite is applied only when the cost of the original plan exceeds the value of the polar_cbqt_cost_threshold parameter.

-- Enable CBQT
SET polar_enable_cbqt to on;

-- Enable sublink pushdown
SET polar_cbqt_pushdown_sublink to on;

EXPLAIN ANALYZE SELECT * FROM (SELECT a, sum(b) b FROM t_big GROUP BY a)v WHERE a IN (SELECT a FROM t_small);

The following result is returned:

                                                             QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------
 GroupAggregate  (cost=17.96..17.99 rows=2 width=12) (actual time=0.056..0.059 rows=2 loops=1)
   Group Key: t_big.a
   ->  Sort  (cost=17.96..17.96 rows=2 width=8) (actual time=0.051..0.052 rows=2 loops=1)
         Sort Key: t_big.a
         Sort Method: quicksort  Memory: 25kB
         ->  Nested Loop  (cost=1.46..17.95 rows=2 width=8) (actual time=0.032..0.045 rows=2 loops=1)
               ->  Unique  (cost=1.03..1.04 rows=2 width=4) (actual time=0.014..0.018 rows=2 loops=1)
                     ->  Sort  (cost=1.03..1.03 rows=2 width=4) (actual time=0.014..0.015 rows=2 loops=1)
                           Sort Key: t_small.a
                           Sort Method: quicksort  Memory: 25kB
                           ->  Seq Scan on t_small  (cost=0.00..1.02 rows=2 width=4) (actual time=0.007..0.008 rows=2 loops=1)
               ->  Index Scan using t_big_a_idx on t_big  (cost=0.42..8.44 rows=1 width=8) (actual time=0.010..0.010 rows=1 loops=2)
                     Index Cond: (a = t_small.a)
 Planning Time: 0.518 ms
 Execution Time: 0.141 ms
(15 rows)

Enable sublink pushdown through a hint

You can forcibly enable sublink pushdown at the SQL level by using a hint. The a in (select a from t_small) clause is also pushed down into the subquery. This generates a parameterized path for t_big based on the join condition, significantly reducing the amount of data scanned and the execution time.

-- polar_cbqt_pushdown_sublink is disabled by default
SET polar_cbqt_pushdown_sublink to off;

EXPLAIN ANALYZE /*+ Set(polar_cbqt_pushdown_sublink force) */ SELECT * FROM (SELECT a, sum(b) b FROM t_big GROUP BY a)v WHERE a IN (SELECT a FROM t_small);

The following result is returned:

                                                             QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------
 GroupAggregate  (cost=17.96..17.99 rows=2 width=12) (actual time=0.073..0.076 rows=2 loops=1)
   Group Key: t_big.a
   ->  Sort  (cost=17.96..17.96 rows=2 width=8) (actual time=0.067..0.069 rows=2 loops=1)
         Sort Key: t_big.a
         Sort Method: quicksort  Memory: 25kB
         ->  Nested Loop  (cost=1.46..17.95 rows=2 width=8) (actual time=0.026..0.040 rows=2 loops=1)
               ->  Unique  (cost=1.03..1.04 rows=2 width=4) (actual time=0.011..0.015 rows=2 loops=1)
                     ->  Sort  (cost=1.03..1.03 rows=2 width=4) (actual time=0.010..0.011 rows=2 loops=1)
                           Sort Key: t_small.a
                           Sort Method: quicksort  Memory: 25kB
                           ->  Seq Scan on t_small  (cost=0.00..1.02 rows=2 width=4) (actual time=0.005..0.006 rows=2 loops=1)
               ->  Index Scan using t_big_a_idx on t_big  (cost=0.42..8.44 rows=1 width=8) (actual time=0.009..0.009 rows=1 loops=2)
                     Index Cond: (a = t_small.a)
 Planning Time: 0.788 ms
 Execution Time: 0.156 ms
(15 rows)