How does ClickHouse execute a query against a Distributed table across shards?
answer
- one node receives it, many nodes do the work
- the distributed table stores nothing itself
- shards send partial state, not finished rows
- one node has to merge everything at the end
- a subquery over a distributed table sees only a slice
basics
~20 sThe node receiving the query rewrites it against each shard's local table, sends it to one replica per shard, and merges what comes back. Shards return partially aggregated state rather than raw rows, so the initiator does the final merge — which is why it can become the bottleneck.
solid answer
~50 sA `Distributed` table is a routing layer over local tables on each shard. The node that receives the query — the initiator — rewrites it to reference the local table, forwards it to one replica per shard, and combines the results. For aggregations the shards do not return rows; they return **partial aggregation states**, and the initiator merges them, so the data crossing the network is proportional to the number of groups rather than the number of rows. Two consequences dominate in practice. First, the initiator does the final merge, sort and limit alone, so a query with millions of groups or a huge global `ORDER BY` concentrates memory and CPU on one node. Second, a subquery on the right of `IN` or `JOIN` that references a distributed table is executed *per shard against that shard's local data*, which silently returns wrong results — `GLOBAL IN` and `GLOBAL JOIN` fix this by evaluating once on the initiator and shipping a temporary table to every shard.
code
sql · 9 lines-- WRONG on a cluster: each shard evaluates the subquery against its own slice
SELECT count()
FROM events
WHERE user_id IN (SELECT user_id FROM vip_users);
-- RIGHT: evaluated once on the initiator, broadcast to every shard
SELECT count()
FROM events
WHERE user_id GLOBAL IN (SELECT user_id FROM vip_users);go deeper
Know that a ClickHouse Distributed table holds no data — it forwards a query to the local tables on each shard and combines the answers, and that inserts are routed by a sharding key.
Explain the rewrite to the local table, the fan-out to one replica per shard, and that aggregations return partial states merged by the initiator so network cost tracks group count, not row count.
Show the failure modes you have hit: the initiator as the merge bottleneck on high-cardinality grouping, silent wrong answers from IN over a distributed subquery, and shard skew diagnosed by joining per-shard log rows on initial_query_id.
Own the cluster's shape: choosing a sharding key that keeps groups and joins shard-local, deciding whether to forbid distributed subqueries by policy, and sizing the coordinating node for the merge work the workload actually generates.
## The two table layers A ClickHouse cluster separates storage from routing. Each shard holds a **local** MergeTree table — say `events_local`. On top sits a `Distributed` table, `events`, declared with the cluster name, the database, the local table's name and a sharding key. The distributed table stores no data of its own; it knows how to fan a query out and how to route inserts. ## The execution flow 1. A client sends a query to some node; that node becomes the **initiator**. 2. The initiator analyses the query and **rewrites** it so the distributed table is replaced by the local table name. 3. It sends that rewritten query to one replica of each shard (choosing among replicas by a load-balancing policy; `prefer_localhost_replica` makes it use its own copy when it hosts one). 4. Each shard executes locally, in full parallel, applying its own `max_threads` — the setting is per server, so a five-shard cluster with `max_threads = 8` may run forty query threads plus the initiator's. 5. Shards stream results back to the initiator, which performs the final combination and returns to the client. In a query plan on the initiator, the remote legs appear as a read-from-remote step; EXPLAIN there does not expand into each shard's own pipeline. To see a shard's plan you run EXPLAIN against the local table on that shard. ## Partial aggregation states are the key idea For a `GROUP BY`, the shards do not ship rows. Each shard aggregates its own data and returns **intermediate aggregation state** — for `sum` that is a running total per group, for a distinct-count sketch it is the sketch itself — and the initiator merges the states per group. This is what makes MPP aggregation cheap: network traffic scales with distinct groups, not with rows scanned. It also explains the failure mode. If the grouping key has tens of millions of distinct values, every shard sends a large state set and the **initiator** must hold and merge all of them on one machine. The memory limit that fires will be on the initiator, not the shards, and the fix is to reduce cardinality, not to add shards. Several settings shape this. `distributed_group_by_no_merge` skips the merge entirely — correct only when the data is already sharded such that no group spans shards. `optimize_distributed_group_by_sharding_key` lets the engine apply that reasoning automatically when the grouping key is compatible with the sharding key. `distributed_aggregation_memory_efficient` merges incoming states in a streaming fashion rather than accumulating everything first. ## The GLOBAL trap This is the correctness question interviewers actually ask. Consider: ```sql SELECT count() FROM events WHERE user_id IN (SELECT user_id FROM vip_users); ``` If `vip_users` is itself distributed, the rewritten query sent to each shard evaluates the subquery **against that shard's local slice** of `vip_users`. Each shard therefore filters against a different, partial set, and the total is wrong — quietly, with no error. `GLOBAL IN` changes the protocol: the initiator evaluates the subquery once, materialises the result as a temporary table, and ships it to every shard, so all shards filter against the same complete set. `GLOBAL JOIN` does the same for the right-hand side of a join. The cost is the broadcast: the temporary set is sent to every shard, so it must be small enough that shipping it is cheaper than the alternative. `distributed_product_mode` governs how the server treats such distributed subqueries — from allowing them through to rejecting them outright — and setting it to reject is a reasonable guardrail on a cluster where this mistake keeps happening. ## Sorting, limits and skew A global `ORDER BY ... LIMIT n` pushes a local sort and limit to each shard, then merges the pre-limited streams on the initiator — efficient. A global `ORDER BY` with no limit, or one over a large result, concentrates the sort on the initiator. Data skew is the other operational hazard: if the sharding key sends a disproportionate share of rows to one shard, that shard becomes the straggler and the query is as slow as it is, no matter how idle the rest are. Compare per-shard rows in `system.parts` on each node, and compare per-shard timings by joining the shard rows in `system.query_log` on `initial_query_id`. ## Observing it Every participating server writes its own `system.query_log` rows. On the initiator the user's query has `is_initial_query = 1`; each shard's sub-query has the same `initial_query_id` with `is_initial_query = 0`. That join is how you find out which shard was slow, which one read the most, and whether the imbalance is data or hardware.
- Why can a distributed GROUP BY exhaust memory on the coordinating node while every shard is comfortable?Each shard aggregates only its own slice and sends partial states, but the initiator receives every shard's states and merges them in one place. With a high-cardinality grouping key the merged state is far larger than any single shard's. Adding shards makes this worse, not better. Reduce the cardinality, pre-aggregate, or enable memory-efficient distributed aggregation so states are merged as they stream in.
- When is distributed_group_by_no_merge safe to enable?Only when no group can span shards — that is, when the sharding key guarantees all rows for a given grouping key live on one shard. Then each shard's aggregation is already final and the initiator can simply concatenate. Enable it on a query whose grouping key is compatible with the sharding key; enable it blindly and you get duplicate group rows, one per shard, with partial values.
- How do you tell whether one shard is the straggler in a distributed ClickHouse query?Every server logs its own row. Take the `initial_query_id` from the initiator's row and look up the matching rows across the cluster, comparing `query_duration_ms`, `read_rows` and `memory_usage` per host. A shard reading far more rows points at data skew from the sharding key; similar rows but a longer duration points at hardware, merge backlog or local contention.
The initiator is a survey office: it posts the same questionnaire to every regional branch, each branch tallies its own returns, and only the tallies come back to be added up centrally.
saying these in an interview costs you the question
- Thinks the Distributed table stores its own copy of the data
- Believes shards return raw rows for the initiator to aggregate
- Uses IN with a distributed subquery and trusts the result
- Assumes max_threads is a cluster-wide total
- Expects more shards to fix a memory error on the coordinator