What does ClickHouse's max_threads setting control, and when does raising it not help?
answer
- it is a per-query cap, not a server-wide one
- its default follows the hardware
- threads need something to divide up
- one hash table per thread during aggregation
- some stages funnel down to a single stream
basics
~20 smax_threads caps how many threads ClickHouse uses to process one query on one server, defaulting to the machine's physical core count. Raising it cannot help when there are too few mark ranges to split, when the query is I/O-bound, or when the expensive stage is single-stream.
solid answer
~50 s`max_threads` is a per-query, per-server setting limiting the number of query-processing threads, and it defaults to the number of physical CPU cores. ClickHouse splits a MergeTree read into ranges of granules and hands them to that many streams, which is why you see `× N` next to processors in `EXPLAIN PIPELINE`. Raising it does nothing in three common situations: there are fewer mark ranges than threads, so extra threads have nothing to read; the bottleneck is disk or network rather than CPU; or the dominant stage is inherently single-stream, such as the final merge of a sort or a global `LIMIT`. It is not free either — most parallel aggregations build one hash table per thread, so memory grows roughly with the thread count, and on a busy server a high `max_threads` on one query starves the rest. Set it per query in a `SETTINGS` clause rather than globally.
code
sql · 6 lines-- cap threads for one heavy report so it cannot monopolise the box
SELECT user_id, count(), uniqExact(session_id)
FROM events
WHERE event_date >= today() - 30
GROUP BY user_id
SETTINGS max_threads = 4, max_memory_usage = 8000000000;go deeper
Know that max_threads limits the threads used for one query, that its default follows the CPU core count, and that you can set it in a SETTINGS clause on the statement itself.
Explain how mark ranges are split across streams, why the multiplier in EXPLAIN PIPELINE tracks the setting, and why per-thread aggregation hash tables make memory grow with the thread count.
Demonstrate judgment about when it is the wrong lever: too few ranges, an I/O-bound read, or a single-stream final stage. Be ready to argue for lowering it to protect memory or to isolate a noisy reporting workload.
Own the thread budget across the fleet: per-user profiles and query limits are a workload-isolation policy, not a per-query tweak, and the setting interacts with concurrency, memory ceilings and the cost of the cluster.
## What the setting is `max_threads` limits the number of threads that process **one query** on **one server**. Its default is `auto`, which resolves to the number of physical CPU cores on that machine. It is a normal query-level setting, so it can be set in a `SETTINGS` clause, in a session, in a user profile, or in the server configuration — and the narrowest scope you can get away with is the right one. ```sql SELECT domain, count() FROM hits WHERE event_date >= today() - 7 GROUP BY domain SETTINGS max_threads = 4; ``` ## How the parallelism is actually created A MergeTree table is stored as parts, and each part's data is addressed by **marks**, one per granule of rows. When a query reads, the engine works out which parts and which ranges of marks survive the primary key and partition conditions, and then distributes those ranges across up to `max_threads` reading streams. Each stream pulls blocks of columns and pushes them into the rest of the pipeline, which is likewise instantiated once per stream. `EXPLAIN PIPELINE` shows this directly as the `× N` multiplier next to each processor. That mechanism is why the setting is an *upper bound* and not a promise. If the surviving work is a single small range of marks, there is nothing to split, and you will see `× 1` no matter how high you set the number. ## When raising it does not help **Too little work to divide.** A selective query against a single small part, or a table where a strict filter leaves one range of granules, cannot be split further. More threads simply idle. **The bottleneck is not CPU.** If the query is waiting on disk reads, on decompression from cold storage, or on the network for a distributed leg, additional CPU threads change nothing. Check `read_bytes` and the timing counters in `system.query_log` before assuming a CPU bound. **The expensive stage is single-stream.** The final merge of a global sort, a top-level `LIMIT`, and the last merge of aggregation states all funnel to one stream. `EXPLAIN PIPELINE` shows this as a `Resize N → 1`; anything above that point runs on a single thread regardless of `max_threads`. **The server is already saturated.** Threads are drawn from a shared pool. If dozens of concurrent queries each want the full core count, they contend, and the aggregate throughput can fall while individual latencies rise. ## What raising it costs Parallel aggregation builds a **separate hash table per thread** and merges them at the end, so peak memory for a high-cardinality `GROUP BY` grows roughly in proportion to the thread count. A query that fits comfortably at `max_threads = 4` can exceed the per-query memory limit at 32 with no change in the data. The merge of those partial states is itself extra work, so on small results the parallel version can even be slower. ## When lowering it is the right move Lowering `max_threads` is an underrated tool: - To **cut memory** on a `GROUP BY` that is close to its limit, without changing the query. - To **isolate workloads**: give a background reporting user a profile with a small `max_threads` so their scans cannot monopolise the cores that interactive dashboards need. - To make a benchmark **reproducible**, since thread scheduling introduces variance. ## Related settings, and what max_threads is not `max_threads` governs query processing. It does not govern insert parallelism (`max_insert_threads` covers the writing side), background merges (which have their own pool sized by the server configuration), or the number of concurrent queries the server accepts. Nor does it control across-shard parallelism: on a distributed query, `max_threads` applies **on each server independently**, so a five-shard query with `max_threads = 8` may be using forty threads across the cluster. ## Diagnosing it properly The honest routine is: run the query, read `query_duration_ms`, `read_rows` and `memory_usage` from `system.query_log`, then look at `EXPLAIN PIPELINE` for the multipliers. If the multiplier already equals `max_threads` and the query is still slow, the answer is not more threads — it is reading less data. If the multiplier is stuck at one, find out why the read could not be split before touching the setting at all.
- Why can raising max_threads make a GROUP BY fail with a memory error?Parallel aggregation gives each thread its own hash table and merges them at the end, so peak memory scales roughly with the thread count. On a high-cardinality GROUP BY, going from four to thirty-two threads can multiply the aggregation state enough to breach `max_memory_usage`. Lowering `max_threads` is often the quickest way to bring such a query back under its limit.
- On a five-shard cluster, does max_threads = 8 mean eight threads in total?No. The setting is per server. The initiating node forwards the query to each shard, and each shard applies `max_threads` locally, so the cluster may be running up to forty query-processing threads plus the initiator's own merge work. Sizing thread limits for a cluster means thinking per node, not per query.
- EXPLAIN PIPELINE shows × 1 on a big table even with max_threads = 16. What do you check?How much of the table survived pruning. If the key conditions left a single range of marks — one small part, or one granule range — there is nothing to split across streams. Check the parts and granules with `EXPLAIN PLAN indexes = 1` and the `SelectedParts` and `SelectedMarks` counters in `system.query_log`. Very small tables and heavily filtered reads legitimately run single-stream.
saying these in an interview costs you the question
- Treats max_threads as a server-wide concurrency limit
- Assumes doubling threads always halves the runtime
- Ignores that per-thread hash tables multiply aggregation memory
- Sets it globally in a profile instead of per query
- Thinks it also controls background merge or insert threads