How to make distributed join of three or more tables as local join?

Viewed 1092

I am mainly confused about the execution plan of three tables, this is the execution plan of query as below(note: t1d,t2d,t3d are distributed table):

select xxx
from t1d t1d
left join 
    (select * from t2d where xxx group by xxx) t2d 
using A
left join
    (select * from t3d where xxx group by xxx) t3d
using A
where t1d.xxx
group by t1d.xxx
SETTINGS distributed_product_mode='local'
┌─explain───────────────────────────────────────────────────────────────────────────────────────────┐
│ Expression (Projection)                                                                           │
│   CreatingSets (Create sets before main query execution)                                          │
│     Expression (Before ORDER BY)                                                                  │
│       Join (JOIN)                                                                                 │
│         Expression ((Before JOIN + Projection))                                                   │
│           SettingQuotaAndLimits (Set limits and quota after reading from storage)                 │
│             Union                                                                                 │
│               Expression ((Convert block structure for query from local replica + ))              │
│                 CreatingSets (Create sets before main query execution)                            │
│                   Expression (Before ORDER BY)                                                    │
│                     AddingDelayedSource (Add non-joined rows after JOIN)                          │
│                       Join (JOIN)                                                                 │
│                         Expression (Before JOIN)                                                  │
│                           SettingQuotaAndLimits (Set limits and quota after reading from storage) │
│                             ReadFromPreparedSource (Read from NullSource)                         │
│                   CreatingSet (Create set for JOIN)                                               │
│                     Expression ((Projection + Before ORDER BY))                                   │
│                       Aggregating                                                                 │
│                         Expression (Before GROUP BY)                                              │
│                           SettingQuotaAndLimits (Set limits and quota after reading from storage) │
│                             ReadFromStorage (MergeTree)                                           │
│               ReadFromPreparedSource (Read from remote replica)                                   │
│     CreatingSet (Create set for JOIN)                                                             │
│       Expression (Projection)                                                                     │
│         SettingQuotaAndLimits (Set limits and quota after reading from storage)                   │
│           Union                                                                                   │
│             Expression ((Convert block structure for query from local replica + Before ORDER BY)) │
│               SettingQuotaAndLimits (Set limits and quota after reading from storage)             │
│                 ReadFromStorage (MergeTree)                                                       │
│             ReadFromPreparedSource (Read from remote replica)                                     │

From my understanding, I think the step as below:

  1. t1_local and t2_local do local join on each shard as your reply, and I use explain syntax to find that t2d is written to t2_local, it is true, I am clear about this.

  2. Initiator host combines the results from all shard of local join

  3. each shard do query "select * from t3_local where xxx group by xxx" and combines on the initiator(maybe this is synchronized with step1)

  4. Initiator do join between result of step2 and result of step3.

and my question is that: I hope on each shard can do local join like

select xxx
from t1_local 
left join (select xxx from t2_local)
using A
left join (select xxx from t3_local)
using A

and then the initiator combines results from all shards. I think this is faster than above. But actually the execution plan can't show it.

3 Answers

Seems like this query should work as you expected, but I prefer to accomplish this without the distributed_product_mode setting.

I can assume that you are joining 3 Distributed tables: t1d, t2d, t3d. Clickhouse will work as you expected: it will execute your request on each shard locally and then combine results at initiator.

That means that you can use join of the Distributed table with local tables to achieve expected result:

SELECT xxx
FROM t1d t1d
LEFT JOIN 
    (SELECT * FROM t2d_local WHERE xxx GROUP BY xxx) t2d 
USING A
LEFT JOIN
    (SELECT * FROM t3d_local WHERE xxx GROUP BY xxx) t3d
USING A
WHERE t1d.xxx
GROUP BY t1d.xxx

Change t2d_local and t3d_local with the corresponding local tables

How does JOINs work at Clickhouse?

I'll try to explain with an example of joining 2 tables.

ALGORITHM

Let's assume each shard has a local_table and distributed wrapper over it.

When you join distributed_table with other table (e.g. subquery):

SELECT <fields> FROM distributed_table JOIN (SELECT * FROM some_other_table) USING field
  1. Initiator host sends query to each shard with left table replaced by the corresponding local table:
SELECT <fields> FROM local_table JOIN (SELECT * FROM some_other_table) USING field
  1. Every shard executes the subquery:
SELECT * FROM some_other_table
  1. local_table is joined with previous result of each shard
  2. Results are sent to the initiator host from all the shards
  3. Initiator host combines the results

Let's take a look at an example and play around with the distributed_product_mode setting and local/distributed tables

EXAMPLE

Say we have a cluster cluster_name of two shards: host1 and host2. Let's create tables there:

CREATE TABLE source_local ON CLUSTER cluster_name
(
    `key` Int32,
    `value` String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/default.source_local', '{replica}')
ORDER BY key;

CREATE TABLE source ON CLUSTER cluster_name AS source_local
ENGINE = Distributed(cluster_name, default, source_local)
CREATE TABLE to_join_local ON CLUSTER cluster_name
(
    `key` Int32,
    `value` String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/default.to_join_local', '{replica}')
ORDER BY key;

CREATE TABLE to_join ON CLUSTER cluster_name AS to_join_local
ENGINE = Distributed(cluster_name, default, to_join_local)

Then, perform several insert at host1:

INSERT INTO source_local VALUES (1, 'source1'), (2, 'source2');
INSERT INTO to_join_local VALUES (1, 'tojoin1');

And one insert at host2:

INSERT INTO to_join_local VALUES (2, 'tojoin2');

For better understanding let's visualize local tables:

Shard 1 and Shard 2 local tables

Let's start with the basic configuration ofdistributed_product_mode setting, setting it just to allow. It will not modify the algorithm but also will not throw unnecessary exceptions:

SET distributed_product_mode = 'allow'

JOINING WITH THE DISTRIBUTED TABLE

SELECT
    source.key,
    source.value,
    to_join.value
FROM source AS source
INNER JOIN
(
    SELECT *
    FROM to_join AS tj
) AS to_join USING (key)

Short explanation: Every host perfoms join of left local table with right subquery and then results are combined at the initiator host.

Let's do this step by step according to the algorithm

  1. The following query is sent to the shards:
SELECT
    source.key,
    source.value,
    to_join.value
FROM source_local AS source
INNER JOIN
(
    SELECT *
    FROM to_join AS tj
) AS to_join USING (key)

(note: source table is replaced by source_local table)

  1. All the shards executes the same subquery:
SELECT * FROM to_join AS tj

Since it is a Distributed table, the result will be the same on both of the shards:

┏━━━━━┳━━━━━━━━━┓
┃ key ┃ value   ┃
┡━━━━━╇━━━━━━━━━┩
│   1 │ tojoin1 │
├─────┴─────────┤
│   2 │ tojoin2 │
└─────┴─────────┘
  1. source_local table joins this result on every shard.

Host1: source_local contains both rows with keys 1 and 2, exactly as a result of the subquery. So the result of the join on host1 will contain 2 rows.

Host2: since source_local contains nothing on host2, result of the join will be empty.

Visualization:

Local joins 4-5. These result transferred to the initiator and combined there. The final result:

┏━━━━━┳━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ key ┃ value   ┃ to_join.value ┃
┡━━━━━╇━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│   1 │ source1 │ tojoin1       │
├─────┼─────────┼───────────────┤
│   2 │ source2 │ tojoin2       │
└─────┴─────────┴───────────────┘

JOINING WITH THE LOCAL TABLE

SELECT
    source.key,
    source.value,
    to_join.value
FROM source AS source
INNER JOIN
(
    SELECT *
    FROM to_join_local AS tj
) AS to_join USING (key)

Short explanation: Each shard performs join of two local tables and then results are combined on the initiator

  1. The following query is sent to the shards:
SELECT
    source.key,
    source.value,
    to_join.value
FROM source_local AS source
INNER JOIN
(
    SELECT *
    FROM to_join_local AS tj
) AS to_join USING (key)

2-3. enter image description here 4-5. The final result therefore differs from the previous one:

┏━━━━━┳━━━━━━━━━┳━━━━━━━━━━━━━━━┓
┃ key ┃ value   ┃ to_join.value ┃
┡━━━━━╇━━━━━━━━━╇━━━━━━━━━━━━━━━┩
│   1 │ source1 │ tojoin1       │
└─────┴─────────┴───────────────┘

I think this request should solve your issue

JOIN WITH DISTRIBUTED TABLE , distributed_product_mode = 'local'

SELECT
    source.key,
    source.value,
    to_join.value
FROM source AS source
INNER JOIN
(
    SELECT *
    FROM to_join AS tj
) AS to_join USING (key)
SETTINGS distributed_product_mode = 'local'

The result:

┌─key─┬─value───┬─to_join.value─┐
│   1 │ source1 │ tojoin1       │
└─────┴─────────┴───────────────┘

We perfomed join with the Distributed table, but got the same result as for joining with local table. The reason is that distributed_product_mode = 'local' Clickhouse implicitly does the same as we did when joining with local table. Here are the docs

This join also performs as you want

JOIN WITH DISTRIBUTED TABLE , distributed_product_mode = 'global'

The same query:

SELECT
    source.key,
    source.value,
    to_join.value
FROM source AS source
INNER JOIN
(
    SELECT *
    FROM to_join AS tj
) AS to_join USING (key)
SETTINGS distributed_product_mode = 'global'

results in:

┌─key─┬─value───┬─to_join.value─┐
│   1 │ source1 │ tojoin1       │
│   2 │ source2 │ tojoin2       │
└─────┴─────────┴───────────────┘

We've the same result as is in the first case with Distributed table. Are there any difference? Yeah, there is a difference in a way Clickhouse performs the query:

At the point 2 right subquery will be executed only at one shard and then it will be spreaded across other shards. It can be used for query optimizations, but do no affect the result. You can achive the same result by using GLOBAL JOIN instead of JOIN.

I do this with a subquery.

For example, table t1d, t2d, t3d are distributed table, and i have a query like this:

SELECT t1d.a,
       t2d.b,
       t3d.c
FROM t1d
INNER JOIN t2d ON t1d.id = t2d.id
INNER JOIN t3d ON t2d.id = t3d.id;

so change the sql to this:

set distributed_product_mode = 'local';

SELECT t1d.a,
       t2d.b,
       t3d.c
FROM t1d
INNER JOIN
  (SELECT t2d.*, t3d.*
   FROM t2d
   INNER JOIN t3d ON t2d.id = t3d.id)bc ON t1d.id = bc.id;

then it can go with local join on each shard.

NOTICE: join key and sharding_key must be the same column.

Related