A ClickHouse GROUP BY aborts with MEMORY_LIMIT_EXCEEDED — how do you diagnose and fix it?
answer
- find out which ceiling fired first
- the query log keeps the failed run
- which side of a JOIN gets built into memory?
- thread count multiplies the aggregation state
- spilling to disk trades speed for survival
basics
~20 sRead the exception to see which limit fired — per query, per user or server-wide — then find the real consumer: usually a high-cardinality GROUP BY or a huge JOIN right side. Fix by shrinking the state, enabling spill to disk, or lowering threads, before raising max_memory_usage.
solid answer
~50 sFirst identify **which** ceiling was hit: the exception text names it, and `system.query_log` gives the failed query's `memory_usage`. `max_memory_usage` is per query per server, `max_memory_usage_for_user` covers everything one user runs, and the server-wide limit protects the process. Then find the consumer. A `GROUP BY` on a high-cardinality key builds a hash table per thread; a `JOIN` builds a hash table from the **right-hand** table, so a query joining a large table on the right is the classic accident. The fixes in order of preference: reduce the state — pre-aggregate, narrow the grouping key, use `LowCardinality` columns, or an approximate function such as `uniqCombined` instead of an exact distinct count; put the smaller table on the right, or pick a different `join_algorithm` such as `grace_hash` or `partial_merge`; enable spilling with `max_bytes_before_external_group_by`; lower `max_threads` so fewer hash tables exist at once. Raising the limit comes last, because it just moves the failure to the whole server.
code
sql · 7 lines-- was it the query, the user, or the server? and how big did it get?
SELECT event_time, memory_usage, exception, query
FROM system.query_log
WHERE event_date = today()
AND type = 'ExceptionWhileProcessing'
ORDER BY event_time DESC
LIMIT 5;go deeper
Know that ClickHouse enforces a per-query memory limit and that GROUP BY and JOIN are what usually consume it. Recognise the MEMORY_LIMIT_EXCEEDED error and know where the failed query is logged.
Explain where the memory goes: one aggregation hash table per thread, the right-hand table of a hash join, exact distinct structures. Know the spill settings and the effect of lowering max_threads.
Diagnose before tuning: identify which ceiling fired, pull the failed query from system.query_log, and pick the fix that shrinks the state rather than the one that raises the ceiling. Be ready to defend approximate aggregation as an engineering trade.
Set memory as policy: per-user profiles with matched query, execution-time and rows-read limits, sized against realistic concurrency and the server ceiling, so that one bad query fails alone instead of taking the node with it.
## Read the exception before you touch anything ClickHouse raises `MEMORY_LIMIT_EXCEEDED` from several different ceilings and the message says which one: - `max_memory_usage` — the limit for a **single query on a single server**. - `max_memory_usage_for_user` — the total across everything one user is running concurrently. - The server-wide limit (`max_server_memory_usage`, normally derived from a configured ratio of physical RAM) — the process protecting itself. These lead to completely different conclusions. A per-query breach is a query problem. A per-user or server-wide breach with modest individual queries is a **concurrency** problem: no single query is greedy, but twenty of them together are. Fixing the second by raising the first makes things worse. ## Find the real consumer Pull the failed run from the log: ```sql SELECT event_time, query_duration_ms, memory_usage, exception, query FROM system.query_log WHERE event_date = today() AND type = 'ExceptionWhileProcessing' ORDER BY event_time DESC LIMIT 5; ``` There are only a handful of things in an analytical query that hold large state: **High-cardinality GROUP BY.** The aggregation hash table holds one entry per distinct key, and **each thread builds its own** before the final merge. Grouping by something near-unique — a raw URL, a session id, a user id on a billion-row table — is the most common cause, and it scales with `max_threads`. **JOIN build side.** ClickHouse's default hash join builds its hash table from the **right-hand** table and streams the left through it. `big JOIN small` is fine; `small JOIN big` loads the big one into memory. This ordering is explicit in ClickHouse and is not silently corrected for you. **Exact distinct counts.** `uniqExact` keeps every distinct value; `count(DISTINCT ...)` maps to a similar exact structure by default. **Global ORDER BY.** Sorting a large intermediate result holds blocks in memory before merging. **Large IN sets and array functions** that materialise big intermediates per row. ## The fixes, in the order a good engineer tries them **1. Shrink the state.** Filter earlier so fewer rows reach the aggregation. Reduce the grouping key — group by a truncated timestamp instead of an exact one, by domain instead of full URL. Store repetitive string columns as `LowCardinality`, which lets aggregation work over dictionary codes rather than strings. Swap exact distinct counting for an approximate function such as `uniqCombined`, which holds bounded state and is usually accurate enough for analytics. Where the query is run repeatedly, pre-aggregating into a rollup table removes the problem at the source rather than the query. **2. Fix the JOIN shape.** Put the smaller relation on the right so it is the build side. Push filters into the right-hand subquery so less of it is materialised. If it genuinely does not fit, change `join_algorithm` — `grace_hash` partitions the build side to disk in buckets, `partial_merge` sorts and merges instead of hashing, and `full_sorting_merge` is another non-hash option. Each trades memory for CPU and I/O. **3. Let it spill.** Setting `max_bytes_before_external_group_by` makes the aggregation flush partial state to disk when it grows past that threshold, and merge from disk at the end. The conventional guidance is to set it to roughly half of `max_memory_usage`, because the merge phase itself needs headroom. `max_bytes_before_external_sort` does the analogous thing for `ORDER BY`. Spilling turns a failure into a slower success — a good trade for a nightly job, a poor one for an interactive dashboard. **4. Aggregate in order.** If the `GROUP BY` keys form a prefix of the table's sorting key, `optimize_aggregation_in_order` lets the engine emit finished groups as it goes instead of holding every group, which can collapse memory dramatically. This only applies when the data really is ordered by those columns. **5. Lower `max_threads`.** Fewer threads means fewer simultaneous hash tables. It is the cheapest one-line mitigation and costs only wall-clock time. **6. Only then raise the limit** — and if you do, check it against the server-wide ceiling and the expected concurrency. A per-query limit that several concurrent queries can each reach is a server outage waiting to happen. ## Preventing the repeat Set the limits in **user profiles**, not per query, so that ad-hoc analysts get a lower ceiling than the ETL user, and pair them with `max_execution_time` and `max_rows_to_read` so runaway queries fail fast and cheaply. Watch `memory_usage` percentiles in `system.query_log` rather than waiting for the next exception; queries that trend toward the limit announce themselves well before they break.
- Which side of a ClickHouse JOIN is loaded into memory, and why does that matter here?The right-hand table is the build side for the default hash join: it is read fully into a hash table while the left side streams through. So `large JOIN small` is safe and `small JOIN large` is not, and ClickHouse will not reorder them for you. When the right side genuinely cannot fit, switch `join_algorithm` to `grace_hash` or `partial_merge`, which trade memory for extra CPU and disk I/O.
- When does optimize_aggregation_in_order actually help a GROUP BY's memory?Only when the grouping keys are a prefix of the table's sorting key. Then rows arrive already clustered by group, so the engine can finish and emit a group as soon as it sees the next one, instead of keeping every group's state until the end. Memory becomes roughly proportional to the concurrent groups in flight rather than the total distinct count.
- Why is raising max_memory_usage a last resort rather than a first fix?Because it is per query per server. Twenty concurrent queries each entitled to the new ceiling can collectively exceed the server-wide limit and take the process down, converting one failed query into an outage. Raise it only after checking the server ceiling and the realistic concurrency, and prefer setting it in a user profile so the entitlement matches the workload.
saying these in an interview costs you the question
- Immediately raises max_memory_usage without diagnosing the query
- Assumes ClickHouse reorders a JOIN to build the smaller side
- Thinks spilling to disk is free rather than slower
- Ignores that per-thread hash tables scale with max_threads
- Confuses a per-user breach with a single greedy query