When choosing an analytical engine, how much weight should per-core execution efficiency get versus simply adding more compute?
answer
- weight it by the CPU-bound fraction
- Amdahl: improving 20% of the time buys little
- more nodes fix throughput, not query latency
- efficiency compounds with every concurrent user
- fix pruning first, it is the bigger lever
basics
~20 sWeight it by how much of your workload is actually CPU-bound after pruning. Efficiency buys lower cost per query, better tail latency and more concurrency per node, but scale-out fixes throughput and cannot fix single-query latency or a bad data layout.
solid answer
~50 sStart by measuring where time goes on *your* queries. If most of the wall clock is fetching from object storage, waiting on shuffle, or scanning data that better clustering would have skipped, execution efficiency is a small term and you should spend your effort on layout and pruning. If queries are genuinely CPU-bound in filter, join and aggregation — which is typical of interactive dashboards over already-pruned data — then per-core efficiency is close to a linear discount on cost, because elastic compute is billed in CPU-time units. Two things scale-out cannot buy: single-query latency below what one pass of the pipeline costs, and concurrency density, since an engine burning three times the CPU per row supports a third of the concurrent users per node. Against that, weigh the factors that usually decide platform choice anyway — ecosystem, operability, elasticity, security model, correctness and migration cost. Benchmark with your own data, your own queries and your real concurrency, never a published number.
go deeper
You will not be asked to make this call. Know only that reading less data usually matters more than how fast an engine processes what it reads.
Be able to say that adding compute helps throughput but not the fixed costs of a single query, and that the size of any execution win depends on how much time is actually spent on CPU.
Show the measurement discipline: profile where time goes, separate bytes scanned from CPU time, and test at realistic concurrency before attributing a difference to execution efficiency.
Own the whole tradeoff. Express it as cost per query at a latency target, weigh concurrency density and idle cost, and be explicit that operability, ecosystem and data-format coupling usually outrank raw efficiency in a platform decision.
## Frame the question as a cost model, not a speed contest In an elastic analytical platform you pay, in one currency or another, for compute-time. Whether the meter reads credits, slots, node-hours or CPU-seconds, the quantity is roughly *cores × seconds*. Per-core execution efficiency therefore acts as a direct multiplier on the bill for the CPU-bound fraction of the workload — an engine that does the same work with half the instructions and half the memory stalls costs about half as much for that fraction. That is the strongest argument for taking it seriously, and it is a better argument than "it is faster", because faster can always be bought with more nodes while cheaper cannot. ## First, find the CPU-bound fraction The honest answer begins with measurement, and the ordering matters: 1. **How much data does a typical query actually read?** If a dashboard query scans a terabyte to return a thousand rows, execution efficiency is not the problem — pruning, partitioning, clustering and pre-aggregation are, and they operate on a different order of magnitude. No amount of vectorization beats not reading the bytes. 2. **Where does the remaining time go?** Object-storage fetch, decompression, filter and expression evaluation, join build and probe, shuffle, result serialization. Execution efficiency moves the middle terms only. 3. **What is the shape of the workload?** Long ETL batches are throughput problems and scale out beautifully. Interactive dashboards with tight latency SLOs and many concurrent users are where per-core efficiency shows up in ways more nodes cannot fix. Amdahl's law is the whole argument: a 5x improvement in a component consuming 20% of the time yields about 1.19x overall. ## What scale-out genuinely cannot fix Three things: - **Latency floor for a single query.** Adding nodes reduces per-node work, but coordination, planning, shuffle round trips and the final gather do not shrink proportionally, and eventually adding nodes makes a small query slower. If your target is sub-second dashboards, per-core efficiency and the ability to serve from cache are what get you there. - **Concurrency density.** With a fixed node budget, the number of simultaneous queries you can serve is inversely proportional to CPU-seconds each consumes. An inefficient engine needs more nodes for the same user count, which raises the floor of your always-on spend, not just the marginal cost. - **Tail behaviour under contention.** Efficient engines leave headroom; inefficient ones saturate cores sooner, and queueing effects turn a modest overload into a long tail. Conversely, scale-out fixes throughput cleanly, absorbs spiky batch load, and can be turned off when idle — which for many organizations makes an inefficient-but-elastic engine perfectly economic. ## What usually outweighs efficiency anyway A principal-level answer should be honest that raw execution efficiency rarely decides a platform: - **Elasticity and idle cost.** An engine that suspends to zero when nobody queries may beat a more efficient one that must stay warm. - **Operability.** Who runs it, who patches it, what the failure modes are, and whether your team can debug it at 3 a.m. - **Ecosystem and migration cost.** Existing SQL, BI connectivity, ingestion pipelines, and the cost of rewriting them. - **Governance and isolation.** Access control, data sharing, workload isolation between teams. - **Correctness and semantics.** Type handling, NULL and decimal semantics, transactional guarantees on write. - **Storage coupling.** Whether your data must live in the engine's proprietary format or can stay in open files that several engines read. Efficiency is a cost lever; these are risk levers, and risk usually dominates. ## How to benchmark without fooling yourself If efficiency does matter for your workload, measure it properly: - Use **your** queries and **your** data distribution — skew, cardinality and column widths change results more than engine internals do. - Test at **realistic concurrency**, not one query at a time. Single-query benchmarks systematically flatter engines that consume the whole machine per query. - Control for caching: run cold and warm, and know which layer cached. - Normalize to **cost**, not seconds — cost per query at a fixed latency target is the comparable number. - Separate the terms: measure bytes scanned alongside wall clock, so a win from better pruning is not misread as a win from better execution. - Beware benchmarks whose queries are trivially answerable from metadata or pre-aggregates on one system. ## The judgment to state out loud A defensible position sounds like this: for batch ETL, buy elasticity and pay for the extra CPU; for interactive, high-concurrency serving, per-core efficiency compounds with every concurrent user and deserves real weight; in both cases, fix layout and pruning first, because that lever is an order of magnitude larger and it is one you control regardless of which engine you pick. And be explicit that you would decide on measured cost-per-query at your latency target, not on vendor claims or published benchmark suites.
- Why can adding nodes make a small interactive query slower rather than faster?Because fixed costs do not shrink with width. Planning, task startup, coordination, shuffle round trips and the final gather all grow with the number of participants, while the per-node data shrinks toward triviality. Past a point you are paying coordination overhead to divide work that was already small, which is why interactive latency targets are met with efficiency and caching, not width.
- What would make you decide execution efficiency does not matter for your platform?Evidence that queries are dominated by data volume and I/O rather than CPU — huge scans returning small results, time concentrated in fetch and decompression, and large gains available from partitioning, clustering or pre-aggregation. In that regime the layout lever is an order of magnitude bigger, and I would choose on elasticity, operability and ecosystem instead.
- How do you compare two engines fairly when their pricing units differ?Normalize to cost per query at a fixed latency target and a realistic concurrency level, using your own data and queries. Run cold and warm, record bytes scanned alongside wall clock so pruning wins are not mistaken for execution wins, and include idle and always-on costs, since an engine that suspends to zero can win on the bill while losing on raw speed.
saying these in an interview costs you the question
- Treats faster per core as automatically cheaper without measuring the workload
- Assumes scale-out can meet any single-query latency target
- Compares engines on published benchmarks instead of own queries and data
- Benchmarks one query at a time and ignores concurrency
- Optimizes execution before fixing pruning and data layout