In a columnar engine, why can a filter written as a row-wise scalar UDF run an order of magnitude slower than the same logic in built-in operators?
answer
- the engine cannot see inside the function
- one call per row, not one per batch
- the optimizer loses cost and selectivity too
- the biggest loss is data that no longer gets skipped
- transform the constant, not the column
basics
~20 sA row-wise UDF is a black box invoked once per value, so it collapses batch-at-a-time execution back to row-at-a-time: no SIMD, no branch-free evaluation, often a cross-language boundary per call, and the optimizer can neither push it to storage nor reorder it confidently.
solid answer
~50 sVectorized speed comes from typed, branch-free kernels applied to a whole batch. A scalar UDF has none of those properties: the engine must unpack each value from the vector, call opaque code, and box the result back, once per row. If the function runs in another language or a sandbox, each call also crosses a boundary with serialization and marshalling, which can cost far more than the arithmetic. Just as damaging is what the planner loses — an opaque predicate cannot be pushed into the scan for pruning, cannot be evaluated against zone-map statistics, and has an unknown cost and selectivity, so it may be ordered badly relative to cheap predicates. The fixes, in order: express the logic in built-in expressions so it becomes a kernel; if a UDF is unavoidable, use the engine's batched or columnar UDF interface so one call handles a whole vector; and make sure cheap, selective predicates run first so the UDF only sees survivors.
code
sql · 10 lines-- slow: opaque per-row call, column wrapped, no block pruning
SELECT count(*)
FROM events
WHERE my_month_key(event_ts) = '2026-01';
-- fast: built-in kernel over the vector, and a prunable range on the stored column
SELECT count(*)
FROM events
WHERE event_ts >= TIMESTAMP '2026-01-01 00:00:00'
AND event_ts < TIMESTAMP '2026-02-01 00:00:00';go deeper
Know that calling a custom function on every row is much more expensive than using the built-in operators, and prefer built-in string and date functions when they exist.
Explain the mechanism: an opaque per-row call breaks batch-at-a-time execution, prevents SIMD, and may cross a language boundary with marshalling on every value.
Lead with the optimizer loss, not just the loop cost: no pushdown, no statistics-based skipping, unknown selectivity. Show the sargable rewrite and how you would confirm the diagnosis from scanned bytes.
Treat the extension surface as a platform decision. Offering only scalar UDFs invites this failure at scale; batched interfaces, a curated built-in library and load-time derived columns are the levers that keep a shared platform's cost predictable.
## Why the built-in path is fast A built-in predicate compiles down to a typed kernel over a column vector: fixed-width values, no per-value dispatch, no branches, SIMD-friendly, with everything staying in cache between operators. The engine also *understands* the expression — it knows its cost, its selectivity, whether it is deterministic, whether it can be evaluated against per-block min/max statistics, and whether it can be pushed down into the scan so that whole blocks or files are never read. ## What a scalar UDF removes A row-wise scalar UDF is, by definition, a function from one value to one value. Every property above is lost: **Execution-level losses.** - **One call per row.** The engine unpacks a value from the vector, invokes the function, and writes the result back — thousands of times per batch. Call overhead that vectorization amortized away is back at full price. - **No SIMD, no branch-free evaluation.** The engine cannot fuse an opaque callee into a packed loop, so lanes go unused. - **Boxing and type marshalling.** Values often must be converted from the engine's internal columnar representation into whatever the function's runtime expects, and the result converted back. - **Cross-runtime boundaries.** If the UDF runs in a separate language runtime or a sandboxed process, each invocation may involve serialization, a context switch, or interpreter overhead. This is usually the single largest term, and it is why a UDF doing trivial arithmetic can still be orders of magnitude slower than the built-in equivalent. - **Cache pollution.** Whatever the UDF's runtime touches evicts the vectors the pipeline was relying on staying hot. **Planner-level losses, which are often worse.** - **No pushdown, no pruning.** `WHERE my_udf(ts) = '2026-01'` is opaque, so the engine cannot translate it into a range on `ts` and cannot use per-block statistics to skip data. The equivalent built-in range predicate might have eliminated 99% of blocks before any CPU work happened at all. This is the reason a UDF filter can be slower by far more than the per-call factor suggests. - **Unknown cost and selectivity.** The optimizer has no statistics for the function's output, so predicate ordering and join ordering downstream of it may be badly chosen. A UDF might be evaluated on every row before a cheap, highly selective predicate that would have removed most of them. - **Blocked rewrites.** Common subexpression elimination, constant folding and simplification generally cannot cross a function the engine cannot reason about, and a function not declared deterministic cannot be hoisted out of a loop or cached. ## Diagnosing it The signature is a query that is CPU-bound with modest bytes scanned, where the profile attributes most time to the filter or projection operator rather than to the scan, and where the amount of data read did not drop the way you expected from the predicate. Comparing the same query with the UDF removed — even with a deliberately wrong but cheap predicate of similar selectivity — isolates the effect quickly. Another tell: query time scales with *rows examined* rather than with *rows returned*. ## The fixes, in order of preference 1. **Rewrite in built-ins.** The great majority of scalar UDFs in real warehouses are string manipulation, date arithmetic or `CASE` logic that the SQL dialect already expresses natively. This is the fix that also restores pushdown and pruning. 2. **Rewrite the predicate so the *stored* column is compared directly.** Instead of wrapping the column in a function, transform the constant: compare `ts >= '2026-01-01' AND ts < '2026-02-01'` rather than `month_of(ts) = '2026-01'`. This is the classic sargability rule, and in a columnar engine its payoff is block and file pruning, not index usage. 3. **Use a batched or columnar UDF interface.** Many engines offer an API where the function receives a whole batch — a vector or an array of values — and returns a vector. One call per batch instead of per row removes most of the boundary cost even if the body is still interpreted. 4. **Precompute at load time.** If the derived value is used repeatedly, materialize it as a column when the data is written, so queries filter a real, prunable, encodable column. 5. **Order the predicates.** Where the UDF genuinely cannot be avoided, ensure the cheap selective predicates are evaluated first so the function only runs over the survivors. Vectorized engines apply later predicates only to the current selection, so this can reduce UDF invocations by orders of magnitude. ## What the interviewer is testing Whether you can separate two distinct costs — per-row invocation overhead, and lost optimizer knowledge — and whether you reach for the rewrite that restores pruning rather than only for the one that speeds up the loop. Candidates who mention only "UDFs are slow" miss that the biggest term is usually the data that no longer gets skipped.
- Why is losing pushdown often more expensive than the per-row call overhead itself?Per-row overhead multiplies the cost of processing the rows you read. Losing pushdown changes how many rows you read at all: an opaque predicate cannot be checked against per-block min/max statistics, so blocks that a plain range predicate would have skipped entirely are now fetched, decoded and evaluated. Going from 1% of blocks scanned to 100% dwarfs any constant factor in the loop.
- If a UDF is genuinely unavoidable, what reduces its cost the most?Use a batched interface if the engine offers one, so a single call processes a whole vector instead of a single value — that removes most of the boundary and marshalling cost. Then make sure cheap, selective predicates run first, since later predicates are applied only to surviving positions, and consider precomputing the derived value as a stored column at load time.
- How would you confirm a UDF is the cause rather than guessing?Compare the query against one with the UDF predicate replaced by a cheap built-in of similar selectivity, and look at bytes or blocks scanned in both. If the scanned volume drops sharply, the problem is lost pruning; if scanned volume is the same but CPU time collapses, it is invocation overhead. The profile attributing most time to the filter operator rather than the scan points the same way.
saying these in an interview costs you the question
- Says UDFs are slow only because the language is interpreted
- Misses that an opaque predicate blocks block and file pruning
- Assumes the optimizer can estimate a UDF's selectivity
- Wraps the column in a function instead of transforming the constant
- Thinks adding compute fixes a per-row invocation bottleneck