> ## 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.

# distributed_plan_* session settings

> ClickHouse session settings in the distributed_plan_* generated group.

export const PrivatePreviewBadge = () => {
  return <div className="privatePreviewBadge">
            <div className="privatePreviewIcon">
            <svg width="16" height="16" viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
                <path d="M5.33301 6.66667V4.66667V4.66667C5.33301 3.194 6.52701 2 7.99967 2V2C9.47234 2 10.6663 3.194 10.6663 4.66667V4.66667V6.66667" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" />
                <path d="M8.00033 9.33337V11.3334" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" />
                <path fillRule="evenodd" clipRule="evenodd" d="M11.333 14H4.66634C3.92967 14 3.33301 13.4033 3.33301 12.6666V7.99996C3.33301 7.26329 3.92967 6.66663 4.66634 6.66663H11.333C12.0697 6.66663 12.6663 7.26329 12.6663 7.99996V12.6666C12.6663 13.4033 12.0697 14 11.333 14Z" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" />
            </svg>
        </div>
            {'Private preview'}
        </div>;
};

export const ExperimentalBadge = () => {
  return <a href="https://clickhouse.com/docs/reference/settings/beta-and-experimental-features#experimental-features" className="experimentalBadge">
            <div className="experimentalIcon">
            <svg width="16" height="16" viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
                <path strokeWidth="1.25" d="M5.5 2H10.5" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" />
                <path strokeWidth="1.25" d="M9.50015 2V6.19625L13.4283 12.7425C13.4738 12.8183 13.4985 12.9049 13.4996 12.9934C13.5008 13.0818 13.4785 13.169 13.435 13.246C13.3914 13.323 13.3283 13.3871 13.2519 13.4317C13.1755 13.4764 13.0886 13.4999 13.0002 13.5H3.00015C2.91164 13.5 2.8247 13.4766 2.74822 13.432C2.67174 13.3874 2.60847 13.3233 2.56487 13.2463C2.52126 13.1693 2.49889 13.082 2.50004 12.9935C2.50119 12.905 2.52582 12.8184 2.5714 12.7425L6.50015 6.19625V2" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" />
                <path strokeWidth="1.25" d="M4.47656 9.56754C5.30344 9.41254 6.47656 9.47942 7.99969 10.25C10.0153 11.2707 11.4216 11.0569 12.2184 10.7282" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" />
            </svg>
        </div>
            Experimental feature
        </a>;
};

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="distributed_plan_default_reader_bucket_count">
  distributed\_plan\_default\_reader\_bucket\_count
</h2>

<ExperimentalBadge />

<SettingsInfoBlock type="NonZeroUInt64" default_value="8" />

<VersionHistory rows={[{"id": "row-1","items": [{"label": "25.5"},{"label": "8"},{"label": "New experimental setting."}]}]} />

Default number of tasks for parallel reading in distributed query. Tasks are spread across between replicas.
Used by the rule-based distributed planner. The cost-based optimizer chooses the read fan-out by estimated cost and does not use this setting.

<h2 id="distributed_plan_default_shuffle_join_bucket_count">
  distributed\_plan\_default\_shuffle\_join\_bucket\_count
</h2>

<ExperimentalBadge />

<SettingsInfoBlock type="NonZeroUInt64" default_value="8" />

<VersionHistory rows={[{"id": "row-1","items": [{"label": "25.5"},{"label": "8"},{"label": "New experimental setting."}]}]} />

Default number of buckets for distributed shuffle-hash-join.
Used by the rule-based distributed planner. The cost-based optimizer chooses the fan-out by estimated cost and does not use this setting.

<h2 id="distributed_plan_execute_locally">
  distributed\_plan\_execute\_locally
</h2>

<ExperimentalBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "25.5"},{"label": "0"},{"label": "New experimental setting."}]}]} />

Run all tasks of a distributed query plan locally. Useful for testing and debugging.

<h2 id="distributed_plan_fallback_to_local_execution">
  distributed\_plan\_fallback\_to\_local\_execution
</h2>

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "26.9"},{"label": "1"},{"label": "New setting to fall back to local execution when a plan cannot be distributed (only takes effect under `make_distributed_plan`)."}]}]} />

When a query plan contains a step that does not support distributed execution, log the reason and execute the query on the initiator instead of throwing an exception. Disable to get an exception instead.

Only takes effect when `make_distributed_plan` (private preview) is enabled.

<h2 id="distributed_plan_force_exchange_kind">
  distributed\_plan\_force\_exchange\_kind
</h2>

<ExperimentalBadge />

<VersionHistory rows={[{"id": "row-1","items": [{"label": "25.5"},{"label": ""},{"label": "New experimental setting."}]}]} />

Force specified kind of Exchange operators between distributed query stages.

Possible values:

* '' - do not force any kind of Exchange operators, let the optimizer choose,
* 'Persisted' - use temporary files in object storage,
* 'Streaming' - stream exchange data over network.

<h2 id="distributed_plan_force_shuffle_aggregation">
  distributed\_plan\_force\_shuffle\_aggregation
</h2>

<ExperimentalBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "25.7"},{"label": "0"},{"label": "New experimental setting"}]}]} />

Use Shuffle aggregation strategy instead of PartialAggregation + Merge in distributed query plan.
Ignored where the Shuffle strategy cannot produce a correct result, for example for `GROUPING SETS` or when the aggregation must produce results in bucket order.

<h2 id="distributed_plan_max_buffered_log_rows">
  distributed\_plan\_max\_buffered\_log\_rows
</h2>

<ExperimentalBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "26.10"},{"label": "100000"},{"label": "New setting bounding how many log lines a stateless-worker task buffers for forwarding to the coordinator between status polls; excess lines are dropped and counted. New feature, so the previous value equals the default."}]}]} />

When `send_logs_level` forwards stateless-worker task logs to the coordinator, each worker task buffers at most this many log lines between status polls. Lines beyond the bound are dropped and their count is reported to the client. `0` means unbounded (never drops, but a stalled status poll can grow the buffer without limit).

<h2 id="distributed_plan_max_rows_to_broadcast">
  distributed\_plan\_max\_rows\_to\_broadcast
</h2>

<ExperimentalBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "25.7"},{"label": "20000"},{"label": "New experimental setting."}]}]} />

Maximum rows to use broadcast join instead of shuffle join in distributed query plan.
A heuristic for the rule-based distributed planner. When the cost-based optimizer is enabled, the broadcast-vs-shuffle choice is made by estimated cost and this setting has no effect.

<h2 id="distributed_plan_optimize_exchanges">
  distributed\_plan\_optimize\_exchanges
</h2>

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "25.5"},{"label": "1"},{"label": "New experimental setting."}]}]} />

Removes unnecessary exchanges in distributed query plan. Disable it for debugging.

<h2 id="distributed_plan_prefer_replicas_over_workers">
  distributed\_plan\_prefer\_replicas\_over\_workers
</h2>

<PrivatePreviewBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "26.4"},{"label": "0"},{"label": "New setting to serialize distributed plan for replicas"}]}]} />

Serialize the distributed query plan for execution at replicas.

<h2 id="distributed_plan_read_in_order">
  distributed\_plan\_read\_in\_order
</h2>

<ExperimentalBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "26.9"},{"label": "0"},{"label": "New setting to allow the read-in-order optimization for `ORDER BY` in a distributed query plan, so a sorted read of the table's sorting key can skip the sort and stop early. Off by default: only shapes where no exchange survives between the read and the sort are safe today."}]}]} />

Allow the read-in-order optimization for `ORDER BY` in a distributed query plan, so a sorted read of the
table's sorting key can skip the sort and stop early instead of scanning and sorting.

Off by default: the rewrite that distributes a sort assumes the sort it wraps does not depend on its input
already being ordered, and a sort that does can be fed rows through an exchange that does not preserve
order. Only shapes where no exchange survives between the read and the sort are safe today.

<h2 id="distributed_plan_workers_num">
  distributed\_plan\_workers\_num
</h2>

<PrivatePreviewBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "26.6"},{"label": "0"},{"label": "New experimental setting for how many stateless workers to lease for a distributed query. Zero disables leasing."}]}]} />

How many stateless workers will be used to execute this query. Zero disables stateless-worker leasing for distributed plans.

<h2 id="distributed_plan_workers_provisioning_timeout_ms">
  distributed\_plan\_workers\_provisioning\_timeout\_ms
</h2>

<PrivatePreviewBadge />

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

<VersionHistory rows={[{"id": "row-1","items": [{"label": "26.8"},{"label": "10000"},{"label": "New setting bounding how long a query waits for leased stateless workers to become reachable before execution; `compatibility` below 26.8 restores the previous no-wait behavior."}]}]} />

Total wall-clock time, in milliseconds, a query may spend provisioning stateless workers before execution: leasing them from the discovery service and verifying they are reachable. The query blocks up to this budget for the leased workers to become ready; when it elapses the query proceeds with the workers verified so far, or fails if none became available. Zero waits only for the initial lease-and-verify pass (no retries).
