> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-parallel-read-in-order-multi-part.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

# join_* session settings

> ClickHouse session settings in the join_* generated group.

export const VersionHistory = ({rows = []}) => {
  if (rows.length === 0) {
    return null;
  }
  const headers = ["Version", "Default value", "Comment"];
  const border = "1px solid rgba(128, 128, 128, 0.3)";
  const cell = {
    border,
    padding: "0.25rem 0.5rem",
    textAlign: "start",
    verticalAlign: "top"
  };
  return <details className="not-prose" style={{
    border,
    borderRadius: "0.5rem",
    margin: "0.5rem 0",
    padding: "0.5rem 0.75rem",
    fontSize: "0.8125rem",
    lineHeight: "1.125rem"
  }}>
      <summary style={{
    cursor: "pointer",
    fontWeight: 600,
    opacity: 0.72
  }}>
        Version history
      </summary>
      <table style={{
    borderCollapse: "collapse",
    width: "100%",
    margin: "0.5rem 0 0"
  }}>
        <thead>
          <tr>
            {headers.map(header => <th key={header} style={{
    ...cell,
    fontWeight: 600,
    opacity: 0.72
  }}>
                {header}
              </th>)}
          </tr>
        </thead>
        <tbody>
          {rows.map((row, row_index) => <tr key={row.id ?? row_index}>
              {(row.items ?? []).map((item, item_index) => <td key={item_index} style={{
    ...cell,
    overflowWrap: "anywhere"
  }}>
                  {item?.label}
                </td>)}
            </tr>)}
        </tbody>
      </table>
    </details>;
};

export const SettingsInfoBlock = ({type, default_value, changeable_without_restart}) => {
  return <div className="not-prose" style={{
    display: "flex",
    flexWrap: "wrap",
    alignItems: "baseline",
    columnGap: "0.5rem",
    rowGap: "0.125rem",
    margin: "0.375rem 0",
    fontSize: "0.8125rem",
    lineHeight: "1.125rem"
  }}>
      <div style={{
    fontWeight: 600,
    opacity: 0.72
  }}>Type</div>
      <div style={{
    overflowWrap: "anywhere"
  }}>{type}</div>
      <div style={{
    fontWeight: 600,
    opacity: 0.72,
    marginInlineStart: "0.5rem"
  }}>Default</div>
      <div style={{
    overflowWrap: "anywhere"
  }}>{default_value}</div>
      {changeable_without_restart && <div style={{
    fontWeight: 600,
    opacity: 0.72,
    marginInlineStart: "0.5rem"
  }}>
          Changeable without restart
        </div>}
      {changeable_without_restart && <div style={{
    overflowWrap: "anywhere"
  }}>
          {changeable_without_restart}
        </div>}
    </div>;
};

These settings are available in [system.settings](/reference/system-tables/settings) and are autogenerated from [source](https://github.com/ClickHouse/ClickHouse/blob/master/src/Core/Settings.cpp).

<h2 id="join_algorithm">
  join\_algorithm
</h2>

<SettingsInfoBlock type="JoinAlgorithm" default_value="direct,parallel_hash,hash,ie_join" />

<VersionHistory rows={[{"id": "row-1","items": [{"label": "26.8"},{"label": "direct,parallel_hash,hash,ie_join"},{"label": "Appended `ie_join` to the default list, so a join whose `ON` section has only inequality conditions is executed with IEJoin instead of a `CROSS JOIN` with a filter. Being last, it is used only when the other algorithms do not apply."}]}, {"id": "row-2","items": [{"label": "24.12"},{"label": "direct,parallel_hash,hash"},{"label": "'default' was deprecated in favor of explicitly specified join algorithms, also parallel_hash is now preferred over hash"}]}]} />

Specifies which [JOIN](/reference/statements/select/join) algorithm is used.

Several algorithms can be specified, and an available one would be chosen for a particular query based on kind/strictness and table engine.

Spilling for hash-based algorithms is configured separately from `join_algorithm`. All hash-based algorithms use [`max_bytes_before_external_join`](/reference/settings/session-settings/max-bytes#max_bytes_before_external_join) and [`max_bytes_ratio_before_external_join`](/reference/settings/session-settings/max-bytes#max_bytes_ratio_before_external_join) as spill thresholds.

Most algorithms affect a query only when they are the one selected for it. Some, however, change planning merely by being listed — even as a lower-priority fallback that is not ultimately selected — because the decision is made before the algorithm is picked. There are two such effects:

* Join-key type inference becomes stricter (a merge join cannot join keys of different types, for example `String` and `Nullable(String)`). This can change the result types of `USING` columns, and can make a join into a `Join`-engine table fail with `TYPE_MISMATCH`. Triggered by `full_sorting_merge` and `parallel_full_sorting_merge`.
* `ORDER BY ... LIMIT` on the preserved side of a join gets an explicit sort instead of a read in primary-key order, because the join is assumed to break the ordered read (a merge join inserts its own pre-join sort; a partial merge join re-sorts the left blocks; a join that can produce delayed blocks does not propagate the ordered read either). The result is the same, the plan is less efficient. Triggered by `full_sorting_merge`, `parallel_full_sorting_merge`, `partial_merge`, `prefer_partial_merge`, `grace_hash` and `auto`, and also by a non-zero `max_bytes_before_external_join` / `max_bytes_ratio_before_external_join`.

Both apply even when the query eventually runs with `hash` or another algorithm. If this is undesirable, do not list the algorithms above in `join_algorithm` for the affected queries.

Possible values:

* grace\_hash

[Grace hash join](https://en.wikipedia.org/wiki/Hash_join#Grace_hash_join) is used.  Grace hash provides an algorithm option that provides performant complex joins while limiting memory use.
Selecting `grace_hash` explicitly is intended primarily for diagnostic use. To enable join spilling, set [`max_bytes_before_external_join`](/reference/settings/session-settings/max-bytes#max_bytes_before_external_join) or [`max_bytes_ratio_before_external_join`](/reference/settings/session-settings/max-bytes#max_bytes_ratio_before_external_join) instead.

`grace_hash` is external from the first block: the right table is partitioned straight away, where `hash` and `parallel_hash` collect it in memory first and partition it only once it crosses the spill threshold. Pick it when you already know the right side will not fit in memory and want to skip the in-memory phase. The spill threshold itself is the same one every hash algorithm uses, [`max_bytes_before_external_join`](/reference/settings/session-settings/max-bytes#max_bytes_before_external_join) / [`max_bytes_ratio_before_external_join`](/reference/settings/session-settings/max-bytes#max_bytes_ratio_before_external_join), and one of the two has to be non-zero unless `legacy_join_size_limits_trigger_spilling` is on. Without a threshold `grace_hash` is passed over for the next algorithm in the list, and rejected if it is the only one.

The first phase of a grace join reads the right table and splits it into N buckets depending on the hash value of key columns (initially, N is `grace_hash_join_initial_buckets`). This is done in a way to ensure that each bucket can be processed independently. Rows from the first bucket are added to an in-memory hash table while the others are saved to disk. If the hash table grows beyond the spill threshold, the number of buckets is increased along with the assigned bucket for each row. Any rows which don't belong to the current bucket are flushed and reassigned.

Supports `INNER/LEFT/RIGHT/FULL ALL/ANY JOIN`.

* hash

[Hash join algorithm](https://en.wikipedia.org/wiki/Hash_join) is used. The most generic implementation that supports all combinations of kind and strictness and multiple join keys that are combined with `OR` in the `JOIN ON` section.

When using the `hash` algorithm, the right part of `JOIN` is uploaded into RAM.

* parallel\_hash

A variation of `hash` join that splits the data into buckets and builds several hashtables instead of one concurrently to speed up this process.

When using the `parallel_hash` algorithm, the right part of `JOIN` is uploaded into RAM.

* partial\_merge

A variation of the [sort-merge algorithm](https://en.wikipedia.org/wiki/Sort-merge_join), where only the right table is fully sorted.

The `RIGHT JOIN` and `FULL JOIN` are supported only with `ALL` strictness (`SEMI`, `ANTI`, `ANY`, and `ASOF` are not supported).

When using the `partial_merge` algorithm, ClickHouse sorts the data and dumps it to the disk. The `partial_merge` algorithm in ClickHouse differs slightly from the classic realization. First, ClickHouse sorts the right table by joining keys in blocks and creates a min-max index for sorted blocks. Then it sorts parts of the left table by the `join key` and joins them over the right table. The min-max index is also used to skip unneeded right table blocks.

* direct

The `direct` (also known as nested loop) algorithm performs a lookup in the right table using rows from the left table as keys.
It's supported by special storages such as [Dictionary](/reference/engines/table-engines/special/dictionary), [EmbeddedRocksDB](/reference/engines/table-engines/integrations/embedded-rocksdb), and [MergeTree](/reference/engines/table-engines/mergetree-family/mergetree) tables.

For MergeTree tables, the algorithm pushes join key filters directly to the storage layer. This can be more efficient when the key can use the table's primary key index for lookups, otherwise it performs full scans of the right table for each left table block.

Supports `INNER` and `LEFT` joins and only single-column equality join keys without other conditions.

* auto

When set to `auto`, `hash` join is tried first, and the algorithm is switched on the fly to another algorithm if the memory limit is violated.

* full\_sorting\_merge

[Sort-merge algorithm](https://en.wikipedia.org/wiki/Sort-merge_join) with full sorting of joined tables before joining.

* ie\_join

The sort-based [IEJoin](https://vldb.org/pvldb/vol8/p2074-khayyat.pdf) algorithm for a `JOIN` whose `ON` section has two inequality comparisons (`<`, `<=`, `>`, `>=`) between expressions of the joined tables. Supports `ALL INNER/LEFT/RIGHT/FULL JOIN` and `SEMI`/`ANTI` `LEFT/RIGHT JOIN`.

The position in the list sets the priority. Listed after other algorithms, as in the default value, `ie_join` is used only when the `ON` section has no equality conditions. Listed first, it is used whenever the `ON` section has two inequality conditions, and any remaining conditions are applied as a filter over the join result. When there are more than two inequality conditions, the two used by the algorithm are chosen by their estimated selectivity from [column statistics](/reference/engines/table-engines/mergetree-family/mergetree#column-statistics); without statistics, the first two in syntax order are used. Without `ie_join` in the list, an `INNER JOIN` with only inequality conditions is executed as a `CROSS JOIN` with a filter, and the other join kinds as a [block nested loop join](/reference/statements/select/join#join-with-an-arbitrary-on-condition).

Both inputs are accumulated in memory before joining: [`max_rows_in_join`](/reference/settings/session-settings/max-rows#max_rows_in_join) and [`max_bytes_in_join`](/reference/settings/session-settings/max-bytes#max_bytes_in_join) limit the accumulated input of both sides together (not just the right side), with the action on overflow set by [`join_overflow_mode`](/reference/settings/session-settings/join#join_overflow_mode); the sort indexes the operator builds on top of the accumulated input are not counted against the limit. The join operator itself runs in a single thread; only the pre-join sorts of the inputs are parallelized.

* parallel\_full\_sorting\_merge

Same as `full_sorting_merge`, but hash-compatible equality joins are sharded by the hash of the join keys into independent per-shard merge joins that run in parallel (up to `max_threads`), instead of a single merge join. This keeps the low, streaming memory usage of a merge join while using all threads, and the result is not ordered.

Hash-sharding by the join keys is applied only to plain equality joins on key types whose hash agrees with the merge-join comparison, and only when neither side is already sorted. It is skipped in these cases:

* `ASOF` joins, and floating-point / `JSON` / `Object` / `Dynamic` key types: their hashes are not consistent with the merge-join comparison, so equal keys could land in different shards.
* Sides that are already sorted (a MergeTree read in order, or any pre-sorted input): an order-preserving scatter into the per-shard merges can deadlock the pipeline. The in-order read and its `read_in_order_use_virtual_row` optimization are kept instead.
* While the initiator builds a distributed plan (`make_distributed_plan`), because the scattered sort is not serializable for remote execution. The local single-fragment plan and the per-worker fragments re-optimize with that setting disabled, so they can still be sharded.

Skipping it disables only this rewrite, not parallelism in general: the join runs as a single `full_sorting_merge`, and MergeTree sides read in order can still be sharded at the source by primary-key ranges (which order by the same comparison the join uses, so equal keys stay together) when `query_plan_join_shard_by_pk_ranges` is enabled.

* prefer\_partial\_merge

ClickHouse always tries to use `partial_merge` join if possible, otherwise, it uses `hash`. *Deprecated*, same as `partial_merge,hash`.

* default (deprecated)

Legacy value, please don't use anymore.
Same as `direct,hash`, i.e. try to use direct join and hash join (in this order).

<h2 id="join_any_take_last_row">
  join\_any\_take\_last\_row
</h2>

<SettingsInfoBlock type="Bool" default_value="0" />

Changes the behaviour of join operations with `ANY` strictness when the right table has more than one matching row for a key.

<Note>
  This setting applies to [`Join`](/reference/engines/table-engines/special/join) engine tables and hash-based join algorithms. It has no effect on a [block nested loop join](/reference/statements/select/join#join-with-an-arbitrary-on-condition): with no join key there is no group of matching rows to take the last one of.

  If a join is built in parallel, the order of rows can be non-deterministic. This means that `join_any_take_last_row = 1` can return a non-deterministic row for `ANY JOIN` queries.
</Note>

Possible values:

* 0 — If the right table has more than one matching row, only the first one found is joined.
* 1 — If the right table has more than one matching row, only the last one found is joined.

See also:

* [JOIN clause](/reference/statements/select/join)
* [Join table engine](/reference/engines/table-engines/special/join)
* [join\_default\_strictness](/reference/settings/session-settings/join#join_default_strictness)

<h2 id="join_default_strictness">
  join\_default\_strictness
</h2>

<SettingsInfoBlock type="JoinStrictness" default_value="ALL" />

Sets default strictness for [JOIN clauses](/reference/statements/select/join).

Possible values:

* `ALL` — If the right table has several matching rows, ClickHouse creates a [Cartesian product](https://en.wikipedia.org/wiki/Cartesian_product) from matching rows. This is the normal `JOIN` behaviour from standard SQL.
* `ANY` — If the right table has several matching rows, only the first one found is joined. If the right table has only one matching row, the results of `ANY` and `ALL` are the same.
* `ASOF` — For joining sequences with an uncertain match.
* `Empty string` — If `ALL` or `ANY` is not specified in the query, ClickHouse throws an exception.

<h2 id="join_on_disk_max_files_to_merge">
  join\_on\_disk\_max\_files\_to\_merge
</h2>

<SettingsInfoBlock type="UInt64" default_value="64" />

Limits the number of files allowed for parallel sorting in MergeJoin operations when they are executed on disk.

The bigger the value of the setting, the more RAM is used and the less disk I/O is needed.

Possible values:

* Any positive integer, starting from 2.

<h2 id="join_output_by_rowlist_perkey_rows_threshold">
  join\_output\_by\_rowlist\_perkey\_rows\_threshold
</h2>

<SettingsInfoBlock type="UInt64" default_value="5" />

<VersionHistory rows={[{"id": "row-1","items": [{"label": "24.9"},{"label": "5"},{"label": "The lower limit of per-key average rows in the right table to determine whether to output by row list in hash join."}]}]} />

The lower limit of per-key average rows in the right table to determine whether to output by row list in hash join.

<h2 id="join_overflow_mode">
  join\_overflow\_mode
</h2>

<SettingsInfoBlock type="OverflowMode" default_value="throw" />

Defines what action ClickHouse performs when a join reaches any of the following limits:

* [max\_bytes\_in\_join](/reference/settings/session-settings/max-bytes#max_bytes_in_join)
* [max\_rows\_in\_join](/reference/settings/session-settings/max-rows#max_rows_in_join)

Every hash-based [`join_algorithm`](/reference/settings/session-settings/join#join_algorithm)
value honors this setting, including the ones that spill to disk: reaching the
limit stops the query rather than triggering a spill. The exception is
`legacy_join_size_limits_trigger_spilling`: with it on, the part of a join that
already runs on disk spills further instead of acting on this setting.
`ie_join` honors it as well, on the input it accumulates from both sides. `partial_merge` still
handles the limits by switching strategy — see
[`join_algorithm`](/reference/settings/session-settings/join#join_algorithm).

Possible values:

* `THROW` — ClickHouse throws an exception and stops the query.
* `BREAK` — ClickHouse stops the query and does not throw an exception.

Default value: `THROW`.

**See Also**

* [JOIN clause](/reference/statements/select/join)
* [Join table engine](/reference/engines/table-engines/special/join)

<h2 id="join_use_nulls">
  join\_use\_nulls
</h2>

<SettingsInfoBlock type="Bool" default_value="0" />

Sets the type of [JOIN](/reference/statements/select/join) behaviour. When merging tables, empty cells may appear. ClickHouse fills them differently based on this setting.

Possible values:

* 0 — The empty cells are filled with the default value of the corresponding field type.
* 1 — `JOIN` behaves the same way as in standard SQL. The type of the corresponding field is converted to [Nullable](/reference/data-types/nullable), and empty cells are filled with [NULL](/reference/syntax).
