JOIN clause produces a new table by combining columns from one or multiple tables by using values common to each. It is a common operation in databases with SQL support, which corresponds to relational algebra join. The special case of one table join is often referred to as a “self-join”.
Syntax
ON clause and columns from the USING clause are called “join keys”. Unless otherwise stated, a JOIN produces a Cartesian product from rows with matching “join keys”, which might produce results with many more rows than the source tables.
Supported types of JOIN
All standard SQL JOIN types are supported:JOINwithout a type specified impliesINNER.- The keyword
OUTERcan be safely omitted. - An alternative syntax for
CROSS JOINis specifying multiple tables in theFROMclause separated by commas. - If there are no matching columns for a
NATURAL JOIN, it functions like aCROSS JOIN.
When using the analyzer, disabling the
semi_join_include_columns_from_both_sides or anti_join_include_columns_from_both_sides setting makes the corresponding join expose only its preserved side to expressions resolved after the join result is formed.LEFT SEMI JOINandLEFT ANTI JOINexpose only left-side columns.RIGHT SEMI JOINandRIGHT ANTI JOINexpose only right-side columns.- This affects clauses such as
SELECT,PREWHERE,WHERE,GROUP BY,HAVING,QUALIFY,ORDER BY, andLIMIT BY, including qualified wildcards liket1.*. - The
ONexpression of the sameJOINcan still reference both sides.
SELECT * expands columns from both tables.When join_algorithm is set to
partial_merge, RIGHT JOIN and FULL JOIN are supported only with ALL strictness (SEMI, ANTI, ANY, and ASOF are not supported).LATERAL JOIN
JOIN LATERAL lets the subquery on the right side of a join reference columns of the table
expressions on its left side; the subquery is evaluated for each distinct combination of the left-side
column values it references, and its result is joined to every left row with that combination:
ON true predicate is mandatory, as for any other INNER or LEFT JOIN; omitting it is a syntax error.
It is experimental and disabled by default; enable it with the
allow_experimental_lateral_join setting.
Only the following subset is supported so far; anything else is rejected with an error:
INNER JOIN LATERALandLEFT JOIN LATERALonly;RIGHT,FULL,PASTEandNATURALjoins are not supported, andLATERALcannot be combined with aCROSSor comma join at all.- The default
ALLstrictness only;ANY,SEMI,ANTIandASOFare not supported. - No join predicate other than
ON true(ON 1is also accepted);USINGis not supported, and the predicate cannot be omitted. Put the filters that relate the two sides into theWHEREclause of the lateral subquery. - The
GLOBALandLOCALjoin modifiers are not supported. - The lateral subquery must reference at least one column of the left side. Use a regular join for a non-correlated subquery.
- The lateral subquery is evaluated once per distinct value of the left-side columns it references, not
once per left row, so it must not contain functions that are non-deterministic within a query, such as
randorgenerateUUIDv4, or table functions that generate random rows, such asgenerateRandom. Functions that are constant within a query, such asnow, are allowed. - Only a subquery is supported as the lateral table expression. The PostgreSQL table-source forms
LATERAL unnest(...)andCROSS JOIN UNNEST(...)are not supported - use theARRAY JOINclause instead. - The
GROUP BYandORDER BYof the lateral subquery run once over all evaluations together, so themax_rows_to_group_by,max_rows_to_sortandmax_bytes_to_sortlimits count the rows of all evaluations, not of one. They are only supported with thethrowoverflow mode;anyandbreakare rejected. - The rows of all evaluations are matched to the left rows by a single join. As for any hash join,
max_rows_in_joinandmax_bytes_in_joinlimit the side of this join that is kept in memory, and the planner chooses that side: with the default settings (correlated_subqueries_use_in_memory_buffer = 1) it is the left side ofJOIN LATERAL, because the left rows must be fully read before the lateral subquery is evaluated; otherwise it can be the results of all evaluations together. The limits are always enforced as ifjoin_overflow_modewerethrow: withbreak, the join would silently drop unrelated left rows.
Settings
The default join type can be overridden usingjoin_default_strictness setting.
The behavior of the ClickHouse server for ANY JOIN operations depends on the any_join_distinct_right_table_keys setting.
See also
join_algorithmjoin_any_take_last_rowjoin_use_nullspartial_merge_join_rows_in_right_blocksjoin_on_disk_max_files_to_mergeany_join_distinct_right_table_keys
cross_to_inner_join_rewrite setting to define the behavior when ClickHouse fails to rewrite a CROSS JOIN as an INNER JOIN. The default value is 1, which allows the join to continue but it will be slower. Set cross_to_inner_join_rewrite to 0 if you want an error to be thrown, and set it to 2 to not run the cross joins but instead force a rewrite of all comma/cross joins. If the rewriting fails when the value is 2, you will receive an error message stating “Please, try to simplify WHERE section”.
ON section conditions
AnON section can contain several conditions combined using the AND and OR operators. Conditions specifying join keys must:
- reference both left and right tables
- use the equality operator
JOIN type. Note that if the same conditions are placed in a WHERE section and they are not met, then rows are always filtered out from the result.
The OR operator inside the ON clause works using the hash join algorithm — for each OR argument with join keys for JOIN, a separate hash table is created, so memory consumption and query execution time grow linearly with an increase in the number of expressions OR of the ON clause.
Only the equality operator (
=) makes a condition a join key. Other operators between columns of different tables are supported, but without any equality the join has to examine every pair of rows and is much slower; see JOIN with an arbitrary ON condition.table_1 and table_2:
table_2:
Query
C and the empty text column. It is included into the result because an OUTER type of a join is used.
Response
INNER type of a join and multiple conditions:
Query
Response
INNER type of a join and condition with OR:
Query
Response
INNER type of a join and conditions with OR and AND:
A non-equal condition that uses columns from a single table, such as
t1.a = t2.key AND t1.b > 0 AND t2.b > t2.c, is applied to that table alone: t1.b > 0 uses columns only from t1 and t2.b > t2.c uses columns only from t2.
A non-equal condition that compares columns of different tables, such as t1.a = t2.key AND t1.b > t2.key, is also supported; check out the sections below for more details.Query
Response
JOIN with inequality conditions for columns from different tables
ClickHouse supportsALL/ANY/SEMI/ANTI INNER/LEFT/RIGHT/FULL JOIN with inequality conditions in addition to equality conditions.
When the ON section contains an equality between the two tables next to the inequality, the equality is the join key and the inequality is checked on the rows it matched. Such a mixed condition is executed only by the hash, parallel_hash and grace_hash join algorithms.
When the ON section contains no equality between the two tables, there is no join key to match on. Such a join is executed by ie_join when the condition is a pair of inequalities, or otherwise as a block nested loop join.
A condition evaluated during the join may not contain arrayJoin, because it must preserve the number of rows; use ARRAY JOIN in a subquery instead.
Example
Table t1:
t2
JOIN with an arbitrary ON condition
TheON section may be an arbitrary boolean expression over the columns of both tables, such as a range check, an arithmetic comparison, or a call to a scalar function. When it contains no equality between the two tables, there is no join key to match rows on. A pair of inequality conditions is executed by ie_join when that algorithm is listed in join_algorithm, as it is by default. Any other condition is executed as a block nested loop join: the right table is materialized and the condition is evaluated on every pair of rows. An ALL INNER JOIN takes the equivalent form of a CROSS JOIN with the condition as a filter.
The block nested loop join supports every join type and strictness except ASOF JOIN, PASTE JOIN and ANY FULL JOIN. It is not one of the join_algorithm values: it is the last resort, used only when no other algorithm can execute the condition. In EXPLAIN output it appears as a BlockNestedLoopJoin step. When allow_block_nested_loop_join is disabled, a query that would need it is rejected with INVALID_JOIN_ON_EXPRESSION.
Example
Query
Response
join_use_nulls, as in any other join: with a default value above, and with NULL when the setting is enabled.
Strictness
ANY and SEMI normally keep one row per group of rows that share a join key. There is no such group here, so they keep one row per row of the table that drives the join: LEFT ANY and LEFT SEMI emit each left row at most once, RIGHT ANY and RIGHT SEMI each right row at most once. ANY INNER limits both sides at once — each row of either table is used at most once, so the result has at most as many rows as the smaller table. Which pairs make up such a result is arbitrary, exactly as ANY implies, and it may differ between two runs of the same query. For ANY INNER this extends to the number of rows: pairing each row of both sides greedily, in whatever order the rows are examined, may leave a different number of rows unpaired, so the result of count() over an ANY INNER join with no join key is not reproducible and depends on the number of threads and on the physical order of the right table.
Performance
Every pair of rows is examined, so the work grows with the product of the two tables’ row counts rather than with their sum. A block nested loop join is therefore orders of magnitude more expensive than a hash join on the same data, and the gap widens as the tables grow. If a query can be written with at least one equality in its ON section, write it that way.
LEFT ANY, ANY INNER, LEFT SEMI and LEFT ANTI joins stop scanning the right table at a left row’s first matching row, which usually makes them cheaper than the corresponding ALL join. Their right-driven counterparts (RIGHT ANY, RIGHT SEMI, RIGHT ANTI) examine every pair, because the result depends on which right rows matched. join_any_take_last_row has no effect here: with no join key there is no group of matching rows to take the last one of.
Memory is bounded by the materialized right table and does not grow with the size of the result. The right table is subject to the same settings as in other join algorithms, listed under Memory limitations: max_rows_in_join, max_bytes_in_join and join_overflow_mode limit it, and max_bytes_before_external_join makes it spill to disk.
NULL and NaN values in JOIN keys
NULL is not equal to any value, including itself. This means that if a JOIN key has a NULL value in one table, it won’t match a NULL value in the other table.
Example
Table A:
B:
Charlie from table A and the row with score 88 from table B are not in the result because of the NULL value in the JOIN key.
In case you want to match NULL values, use the isNotDistinctFrom function to compare the JOIN keys.
NaN values in float JOIN keys do not follow the NULL rule above.
A scalar comparison of two NaN values (NaN = NaN) is 0, however, JOIN keys are not compared using scalar semantics - NaN keys may match.
Whether they actually do is an implementation detail and depends on the join algorithm, the key type, and the session settings.
Do not rely on a specific behavior.
If you require that NaN rows do not match, map them to NULL values: ON if(isNaN(A.id), NULL, A.id) = B.id.
In an ASOF JOIN, the closest-match column is compared by ordering, which NaN does not support — filter such rows out on both sides.
ASOF JOIN usage
ASOF JOIN is useful when you need to join records that have no exact match.
This JOIN algorithm requires a special column in tables. This column:
- Must contain an ordered sequence.
- Can be one of the following types: Int, UInt, Float, Date, DateTime, Decimal.
- For the
hashjoin algorithm it can’t be the only column in theJOINclause.
ASOF JOIN ... ON:
SELECT count() FROM table_1 ASOF LEFT JOIN table_2 ON table_1.a == table_2.b AND table_2.t <= table_1.t.
Conditions supported for the closest match: >, >=, <, <=.
Syntax ASOF JOIN ... USING:
ASOF JOIN uses equi_columnX for joining on equality and asof_column for joining on the closest match with the table_1.asof_column >= table_2.asof_column condition. The asof_column column is always the last one in the USING clause.
For example, consider the following tables:
ASOF JOIN can take the timestamp of a user event from table_1 and find an event in table_2 where the timestamp is closest to the timestamp of the event from table_1 corresponding to the closest match condition. Equal timestamp values are the closest if available. Here, the user_id column can be used for joining on equality and the ev_time column can be used for joining on the closest match. In our example, event_1_1 can be joined with event_2_1 and event_1_2 can be joined with event_2_3, but event_2_2 can’t be joined.
ASOF JOIN is supported only by hash and full_sorting_merge join algorithms.
It’s not supported in the Join table engine.PASTE JOIN usage
The result ofPASTE JOIN is a table that contains all columns from left subquery followed by all columns from the right subquery.
The rows are matched based on their positions in the original tables (the order of rows should be defined).
If the subqueries return a different number of rows, extra rows will be cut.
Example:
Distributed JOIN
There are two ways to execute a JOIN involving distributed tables:- When using a normal
JOIN, the query is sent to remote servers. Subqueries are run on each of them in order to make the right table, and the join is performed with this table. In other words, the right table is formed on each server separately. - When using
GLOBAL ... JOIN, first the requestor server runs a subquery to calculate one side of the join and collects the result into a temporary table. This temporary table is then passed to each remote server, and queries are run on them using the temporary data that was transmitted. ForLEFTandINNERjoins, the right table is calculated as the subquery. ForRIGHTjoins, the left table is calculated instead, since the right table is the one being preserved and should be read from shards.
GLOBAL. For more information, see the Distributed subqueries section.
Implicit type conversion
INNER JOIN, LEFT JOIN, RIGHT JOIN, and FULL JOIN queries support the implicit type conversion for “join keys”. However the query can not be executed, if join keys from the left and the right tables cannot be converted to a single type (for example, there is no data type that can hold all values from both UInt64 and Int64, or String and Int32).
Example
Consider the table t_1:
t_2:
Usage recommendations
Processing of empty or NULL cells
While joining tables, the empty cells may appear. The setting join_use_nulls define how ClickHouse fills these cells. If theJOIN keys are Nullable fields, the rows where at least one of the keys has the value NULL are not joined.
Syntax
The columns specified inUSING must have the same names in both subqueries, and the other columns must be named differently. You can use aliases to change the names of columns in subqueries.
The USING clause specifies one or more columns to join, which establishes the equality of these columns. The list of columns is set without brackets. More complex join conditions are not supported.
Syntax Limitations
For multipleJOIN clauses in a single SELECT query:
- Taking all the columns via
*is available only if tables are joined, not subqueries. - The
PREWHEREclause is not available. - The
USINGclause is not available.
ON, WHERE, and GROUP BY clauses:
- Arbitrary expressions cannot be used in
ON,WHERE, andGROUP BYclauses, but you can define an expression in aSELECTclause and then use it in these clauses via an alias.
Performance
When running aJOIN, there is no optimization of the order of execution in relation to other stages of the query. The join (a search in the right table) is run before filtering in WHERE and before aggregation.
Each time a query is run with the same JOIN, the subquery is run again because the result is not cached. To avoid this, use the special Join table engine, which is a prepared array for joining that is always in RAM.
In some cases, it is more efficient to use IN instead of JOIN.
If you need a JOIN for joining with dimension tables (these are relatively small tables that contain dimension properties, such as names for advertising campaigns), a JOIN might not be very convenient due to the fact that the right table is re-accessed for every query. For such cases, there is a “dictionaries” feature that you should use instead of JOIN. For more information, see the Dictionaries section.
Memory limitations
By default, ClickHouse uses the hash join algorithm. ClickHouse takes the right_table and creates a hash table for it in RAM. Ifjoin_algorithm = 'auto' is enabled, then after some threshold of memory consumption, ClickHouse falls back to merge join algorithm. For JOIN algorithms description see the join_algorithm setting.
If you need to restrict JOIN operation memory consumption use the following settings:
- max_rows_in_join — Limits number of rows in the hash table.
- max_bytes_in_join — Limits size of the hash table.
enable_adaptive_memory_spill_scheduler can still spill a join that the threshold below made
spill-capable, and
legacy_join_size_limits_trigger_spilling turns the two caps back into spill triggers on disk.
To let a join keep running by spilling the right side to disk instead of failing, use:
- max_bytes_before_external_join — Absolute spill threshold.
- max_bytes_ratio_before_external_join — Spill threshold as a ratio of available memory.
grace_hash; under memory
pressure enable_adaptive_memory_spill_scheduler can spill earlier than they ask for — but only once one of
them is non-zero, since a join with no threshold at all never spills. The join_algorithm you pick
decides how a join spills — grace_hash partitions the right table from the first block, hash and parallel_hash collect
it in memory and switch over when the threshold is crossed — not whether these settings apply. The one exception is
legacy_join_size_limits_trigger_spilling: with it on, standalone grace_hash ignores both thresholds and spills on the two hard caps instead.
Examples
Example:Related content
- Blog: ClickHouse: A Blazingly Fast DBMS with Full SQL Join Support - Part 1
- Blog: ClickHouse: A Blazingly Fast DBMS with Full SQL Join Support - Under the Hood - Part 2
- Blog: ClickHouse: A Blazingly Fast DBMS with Full SQL Join Support - Under the Hood - Part 3
- Blog: ClickHouse: A Blazingly Fast DBMS with Full SQL Join Support - Under the Hood - Part 4