note

Nobody trusts the query planner

Spark re-plans mid-query. Postgres 19 lets you pin the plan instead. The difference comes from what each one can do while a query is running.

· postgres / spark / clickhouse

ANALYZE samples your table, it doesn’t read it. When the sample misleads — skew, correlated columns, a value that was rare last week — the planner picks a different plan than it did yesterday and runs it with total confidence. Same query, same schema, same data volume, ten times the runtime.

Every engine has this problem. They’ve dealt with it in completely different ways.

Spark changes its mind

AQE has been on by default since Spark 3.2. It doesn’t commit to the plan up front: when a shuffle finishes it looks at how much data actually came out and re-plans from there. The join it costed as sort-merge turns out to be 30MB, so it broadcasts instead.

The catch is better than the feature. A statically hinted broadcast is often still faster than letting AQE work it out, because AQE only learns the size by shuffling both sides first. You pay for the shuffle, then find out you didn’t need it. That’s why /*+ BROADCAST(t) */ hasn’t left anyone’s codebase.

Postgres can’t

Spark materialises at shuffle boundaries — a stage finishes, the data sits somewhere, nothing is in flight. That’s the moment that makes re-planning possible. Postgres streams tuples through a pipeline of executor nodes and never has one. By the time the plan is wrong you’re thirty seconds into running it.

So PG19 went the other way and lets you pin the plan up front.

pg_plan_advice

Run EXPLAIN (COSTS OFF, PLAN_ADVICE) on a query whose plan you like and it hands back JOIN_ORDER(f d), HASH_JOIN(d), SEQ_SCAN(f d). You apply it through a setting instead of wedging it into the SQL.

It works by ruling options out, so you can only ever get a plan the planner already considered viable. Worst case you constrain it to something slower.

Every piece of advice comes back marked matched, partially matched, inapplicable, conflicting or failed.

That last part is what I’d pay for. A Spark hint dies quietly — the table creeps past autoBroadcastJoinThreshold, the broadcast stops happening, nothing anywhere says so, and you find out when a four minute job takes forty.

ClickHouse has less to argue with

The sorting key decides so much that there are fewer plan decisions available to go wrong, so you tune the table instead of the query. Get the ORDER BY wrong and no runtime cleverness saves you.

Three engines, three moments at which the decision gets made: mid-query, at plan time, at table design time. That’s most of why your tuning instincts don’t survive moving between them.

← all notes