skip to content

questions

27

The same amount of data is scanned by two queries, yet one starts returning rows almost immediately while the other returns nothing for several seconds and then delivers everything at once. What in the execution plan explains the difference?

level: juniorimportance: must knowfreq 48%

answer

  1. first row early = fully pipelined path
  2. one breaker gates the whole query's first row
  3. limit helps streaming plans, not sorts
  4. index-supplied order deletes the sort
  5. first-row cost vs total cost in the optimizer

basics

~20 s

A fully pipelined plan passes each row up to the client as it is produced, so the first row arrives early. If the plan contains a blocking operator — a sort, a hash build, a grouping step — nothing can be emitted until that operator has read all its input, so results appear only at the end.

solid answer

~60 s

It is pipelining versus blocking. In a **pipelined** plan every operator can emit output after seeing one input row: scan → filter → project → client. The first result row travels the whole plan while the scan is still on the second page, so the client sees output almost immediately, and if the client stops early the engine stops working. In a plan containing a **blocking (pipeline-breaking)** operator — a sort, a hash-join build, a hash aggregation, DISTINCT — that operator must consume its entire input before it knows its first output row. Everything above it waits, so the query looks frozen and then floods. This is why engines and cost models distinguish **first-row cost** from **total cost**: an optimizer told that only a few rows are wanted may pick an index scan that streams in order rather than a cheaper-overall scan-plus-sort. When latency to the first row matters, look for breakers in the plan and try to remove them, typically by providing an index that already supplies the required ordering or grouping.

go deeper

for a junior

Say clearly that a plan with no blocking operator streams rows immediately, while a sort or grouping step forces the engine to read everything first.

for a middle

Add the row-limit consequence and the index-supplied-ordering fix, and name the common breakers.

for a senior

Bring in first-row versus total cost, cursor and interactive-fetch behaviour, top-N sorts, and how to measure time-to-first-row in production.

for a principal

Treat it as a latency-versus-throughput contract with the optimizer and design schemas, indexes and API paging so the hot interactive paths stay pipelined.

## Two ways a plan can behave in time A plan is a tree of operators. Rows flow from the scans at the leaves up to the client at the root. Whether the client sees rows early depends entirely on whether anything in that path has to wait for its whole input. **Pipelined behaviour.** Operators like filters, projections, index lookups and limits are *streaming*: give them one row and they can immediately produce zero or one output row. In a plan built only from these, row 1 is delivered to the client after roughly one row's worth of work. Total time to the last row may be large, but the *first* row is nearly free. **Blocking behaviour.** Some operators cannot produce anything correct until they have seen everything: - A **sort** cannot know which row comes first until it has seen the last one — the minimum might be the final row read. - A **hash aggregation** cannot finalize any group's total while more rows for that group might still arrive. - A **hash join's build side** must have a complete hash table before probing is correct. - **DISTINCT** implemented by hashing or sorting must see all rows to know what is unique. One of these anywhere below the root converts the whole query, from the client's point of view, from "streams" to "waits, then dumps". ## Why this shows up as a support ticket Users report it as inconsistent responsiveness: adding an ordering clause or a grouping clause to a previously snappy query makes it "hang". Nothing got slower in total; the *shape* of the time changed. Two rules of thumb follow: 1. **A row limit only helps a pipelined plan.** Limiting to ten rows above a streaming plan stops the scan after ten qualifying rows are found. Limiting to ten rows above a sort still reads and sorts everything — the limit only reduces what is emitted afterwards. (The common optimization, a top-N sort keeping a small heap, cuts memory dramatically but still consumes the whole input.) 2. **An index can delete the breaker.** If an index already returns rows in the requested order, the optimizer can drop the sort and walk the index, streaming rows and stopping after the tenth. Same for grouping when the input is already clustered on the grouping key: a streaming aggregate emits each group as soon as the key changes. ## First-row cost versus total cost Optimizers cost a plan twice: the cost to produce the first row and the cost to produce all rows. When the engine knows only a few rows are wanted — an explicit limit, a cursor that will be fetched from interactively, an optimizer hint or session setting for fast-first-row behaviour — it can prefer a plan with a worse total cost but a much better first-row cost. Typically that means an ordered index scan with a nested-loop join rather than a full scan with a hash join and a sort. If all rows will be consumed anyway (a batch export), the total-cost plan is the right one, and the blocking plan is frequently the faster of the two overall. ## What else pipelining buys - **Bounded memory.** A pipelined plan holds one row (or one batch) per operator, not the whole intermediate result, so it neither needs a large memory grant nor risks spilling to temporary files. - **Cheap abandonment.** Close the cursor after ten rows and the work below simply stops. - **Better perceived performance.** Interactive UIs can render the first page while the rest arrives. ## What pipelining does not buy It does not reduce total work, and it is not always achievable — sorting and grouping genuinely require seeing all input, and no plan shape changes that. It also does not help when the client can only use a complete result anyway (an aggregate returning one row, a report that must be sorted). ## How to investigate Read the plan from the bottom up and mark every sort, hash, aggregate and materialize node. The topmost such node is where your first row is gated. Then ask whether it can be removed (an index providing order or clustering), shrunk (top-N), or accepted (a true aggregate). Measuring only total query time hides this entirely — measure time-to-first-row separately when latency is what users feel.

  • Does a row limit reduce the work a plan does?
    Only for the streaming part of it. Above a fully pipelined plan the limit stops execution once enough rows are produced, so scanning stops early. Above a blocking operator such as a sort or a hash aggregation, the input has already had to be fully consumed before any row could be emitted, so the limit only trims the output. Removing the breaker — usually with an index that already provides the order — is what turns the limit into real savings.
  • How would you measure whether a query is pipelined?
    Measure time-to-first-row separately from total time, for example by timing the first fetch from a cursor rather than the full result. A pipelined plan returns row one in roughly the time to find one qualifying row; a blocked plan returns nothing until near the end. Reading the plan for sort, hash, aggregate and materialize nodes tells you the same thing statically.

Ordering at a coffee counter: a pipelined plan hands you each drink as it is made; a blocking plan makes every drink for the whole queue before serving anyone.

saying these in an interview costs you the question

  • Blaming the network or client rendering when the plan clearly contains a sort or hash aggregate.
  • Claiming a row limit always makes a query cheap.
  • Thinking pipelining reduces total execution time — it changes when rows appear, not how much work is done.
  • Believing every plan can be made pipelined; sorting and grouping fundamentally need all input.
  • Judging responsiveness solely by total query duration and never measuring time-to-first-row.

context

open as a page

Describe how a naive nested loop join produces its result, and what its cost is in terms of the sizes of the two inputs.

level: juniorimportance: must knowfreq 64%

basics

~20 s

For each row of the outer input, scan the entire inner input and emit the pairs that satisfy the join predicate. Comparisons are O(n*m), and the naive form re-reads the inner input once per outer row, so I/O is the real problem.

open as a page

When a database sorts or groups a large result set, the plan may report that the operator "spilled to disk". What does spilling mean, what triggers it, and how does it show up in query performance?

level: juniorimportance: must knowfreq 60%

basics

~20 s

Each sort or grouping operator gets a limited memory budget. If the rows it must hold exceed that budget, the engine writes partial results into temporary files and reads them back later. Spilling replaces memory work with disk I/O, so latency jumps sharply.

open as a page

Explain how a database execution engine performs a hash join between two tables, describing what happens in each phase and when the optimizer is likely to choose this join method.

level: middleimportance: must knowfreq 70%

basics

~20 s

Phase one (build): read the smaller input fully and hash its join-key values into an in-memory hash table. Phase two (probe): stream the larger input, hash each row's key, look it up, and emit matches. It is chosen for equality joins over large inputs with no useful index.

open as a page

What is vectorized query execution, and why does processing a batch of rows per operator call typically run several times faster than processing one row per call?

level: middleimportance: must knowfreq 45%

basics

~20 s

Each operator call returns a batch — commonly around 1024 rows held column-wise — instead of one row. Per-call dispatch is amortized over the batch, and the inner loops become tight, branch-free, cache-resident and SIMD-friendly, so far more cycles go into real work.

open as a page

Walk through how the classic Volcano (iterator) execution model runs a query plan — what open(), next() and close() do on each operator — and explain where its CPU overhead comes from.

level: middleimportance: must knowfreq 52%

basics

~20 s

Every plan operator implements open/next/close. The root's next() pulls one row from its child, which pulls from its child, down to the scans. It is simple and streams rows early, but costs an indirect call per row per operator, so big scans spend most cycles on machinery rather than data.

open as a page

In a query execution plan, what distinguishes a streaming (pipelined) operator from a blocking, pipeline-breaking one? Give examples of each and explain the practical consequences of a blocking operator.

level: middleimportance: must knowfreq 55%

basics

~20 s

A streaming operator can emit an output row after seeing one input row — filters, projections, index lookups, nested-loop joins. A blocking operator must consume its entire input first — sorts, hash-join builds, hash aggregation, distinct, some window functions. Blocking costs memory, risks spilling to disk, and delays the first row until all input is read.

open as a page

How does a sort-merge join algorithm produce its result, and what must be true of its two inputs before the merge step can begin?

level: middleimportance: must knowfreq 62%

basics

~20 s

Both inputs must arrive sorted on the join key. Two cursors then walk them once in parallel: advance whichever side holds the smaller key, and emit rows when the keys are equal. One pass over each input; the expensive part is getting the inputs sorted.

open as a page

What is an index nested loop join, and how does its cost model differ from a nested loop join over an inner table with no usable index?

level: middleimportance: must knowfreq 58%

basics

~20 s

Instead of scanning the inner input, it performs an index lookup on the inner table for each outer row's join key. Cost becomes outer rows times the cost of one lookup (an index descent plus row fetches) rather than outer rows times the whole inner table, so it scales with the outer side, not the product.

open as a page

Explain how a relational engine sorts a data set that does not fit in the memory available to the sort operator. Describe the algorithm and what determines its I/O cost.

level: middleimportance: must knowfreq 55%

basics

~20 s

It uses external merge sort: fill memory, sort that chunk, write it out as a sorted run, repeat; then merge the runs with a heap, taking the smallest current row across runs. Cost is roughly two I/Os per row per pass, and passes grow when runs exceed the merge fan-in.

open as a page

Relational engines can implement grouping and duplicate elimination either by hashing on the grouping key or by sorting the input and collapsing adjacent equal rows. Compare the two strategies and explain when a planner picks each.

level: middleimportance: must knowfreq 55%

basics

~20 s

Hashing builds one in-memory entry per distinct key and needs memory proportional to the number of groups, but no ordering. Sorting needs no per-group memory but costs an O(n log n) sort, yet is free when input is already ordered and it delivers ordered output. Planners choose by estimated group count versus memory budget.

open as a page

Why does an execution engine build the hash table on the smaller of the two join inputs, and what goes wrong at runtime when the optimizer's row estimate for that input turns out to be far too low?

level: seniorimportance: must knowfreq 55%

basics

~30 s

The build side's size determines the hash table's memory footprint, so the smaller input keeps it in memory. Choice is made from estimated rows, not actual. A large underestimate produces a hash table far bigger than its memory grant, forcing spilling to disk, or leaves the engine probing a huge table against a tiny one — either way, a plan that looked cheap runs for a very long time.

open as a page

What happens when an operator that must buffer its input — such as a sort or a hash-join build side — is granted less memory than its input needs, and how do you recognize and fix that in production?

level: seniorimportance: must knowfreq 47%

basics

~20 s

It spills: the sort writes sorted runs to temporary storage and merges them later; the hash build partitions the input to disk and processes partitions in turn. The query stays correct but does extra passes of I/O. Diagnose via spill counters and temp-file usage; fix the cardinality estimate, the memory grant, or the plan.

open as a page

A report query that returned in milliseconds during testing now runs for hours in production on the same plan shape — an index-driven nested loop join. What property of that operator makes it so sensitive, and how would you reason about the regression?

level: seniorimportance: must knowfreq 50%

basics

~20 s

Its cost is the outer row count times a per-lookup constant, and the optimizer chose it believing the outer side would produce few rows. If the real outer count is far larger, cost grows linearly with no fallback: millions of random index probes instead of one sequential scan plus a hash join.

open as a page

A query joins two large tables on the condition that one table's timestamp falls between two columns of the other, and the plan will not use a hash join no matter how the query is written. Explain why a hash join can only implement equality join conditions.

level: middleimportance: should knowfreq 45%

basics

~20 s

A hash function maps a value to a bucket and destroys ordering and locality: values that are close numerically land in unrelated buckets. So a hash table can answer "which rows have exactly this key?" but never "which rows are less than, or within a range of, this key." Range joins need other algorithms.

open as a page

When a query plan shows a merge join with no sort step beneath it, where does the required ordering come from, and why is that usually cheaper than sorting the rows at query time?

level: middleimportance: should knowfreq 46%

basics

~20 s

The order comes from an ordered access path — typically scanning a B-tree index whose leading columns are the join key, or the output of an earlier order-preserving operator. Reading in index order costs nothing extra, while an explicit sort costs O(N log N) plus temp-storage I/O when it does not fit in memory.

open as a page

What is a block nested loop join, and which weakness of the row-at-a-time version does it address?

level: middleimportance: should knowfreq 40%

basics

~20 s

Instead of one outer row at a time, it buffers a block of outer rows that fits in memory and scans the inner input once per block, comparing each inner row against every buffered outer row. Inner scans drop from the outer row count to the number of blocks; the comparison count stays the same.

open as a page

A table has 50 million rows and a query asks for the 10 newest by timestamp with ORDER BY plus a small LIMIT. If no index provides that order, what does a good execution engine do instead of fully sorting 50 million rows, and what are the limits of that technique?

level: middleimportance: should knowfreq 48%

basics

~20 s

It runs a top-N heap sort: keep a bounded heap of only N rows, compare each incoming row against the heap's worst element, and discard immediately if it loses. Memory is O(N) instead of O(input), so no spill — but the whole input is still scanned, and the trick dies once N gets large.

open as a page

What is a grace hash join (and its hybrid variant), and how does an execution engine complete a hash join correctly when the build input does not fit in the memory budget?

level: seniorimportance: should knowfreq 40%

basics

~20 s

When the build side exceeds memory, the engine hash-partitions both inputs into matching partitions written to temporary storage, then joins one partition pair at a time in memory. That is grace hash join. The hybrid variant keeps the first partition in memory and joins it during partitioning, so part of the work avoids disk entirely.

open as a page

Compare pull-based (demand-driven) and push-based (data-driven) query execution engines: what changes in how an operator is written, and what does a push-based design buy you?

level: seniorimportance: should knowfreq 28%

basics

~20 s

Pull: the consumer calls next() on its child and rows come back as return values. Push: the producer drives the loop and hands each tuple or batch to its parent's consume() method. Push fuses operators into one tight loop, suits code generation and multi-consumer plans; pull makes early termination and multi-input operators easier.

open as a page

Materializing an intermediate result costs memory and delays the first row, so why would a query optimizer deliberately insert a node that buffers an intermediate result instead of recomputing or re-scanning it?

level: seniorimportance: should knowfreq 30%

basics

~20 s

Because the intermediate is consumed more than once, or must be frozen. Buffering it turns repeated re-execution of a subtree into repeated cheap reads, lets one result feed several consumers, stabilizes a result that would otherwise change while the statement modifies the same rows, and shields an expensive or non-deterministic subplan from repeated evaluation.

open as a page

In a sort-merge join, what has to happen when the join key has duplicate values on both sides, and why does that force the algorithm to re-read part of one input?

level: seniorimportance: should knowfreq 32%

basics

~20 s

Equal keys on both sides form groups whose full cross product must be emitted. The algorithm marks the start of the inner group and rewinds to that mark for every row of the outer group, so inner rows are read once per outer duplicate. Output for that key is the product of the two group sizes.

open as a page

A nightly reporting query that sorts and groups a large table ran in 40 seconds for months and now takes 25 minutes, with no code change and only steady data growth. Walk through how you would confirm that sort/aggregate memory is the cause and what you would change.

level: seniorimportance: should knowfreq 42%

basics

~20 s

Get the executed plan with actual row counts and per-operator memory. Look for a sort or hash aggregate that switched from in-memory to external/batched with large temp bytes, and for actual rows far above estimates. Then cut rows and columns feeding it, supply order via an index, refresh statistics, and only then raise the memory budget for that job.

open as a page

For an equality join between two very large tables, when would you expect a sort-merge join to be the better physical operator than a hash join, and what properties of the data, the storage layout and the machine drive that judgement?

level: principalimportance: should knowfreq 36%

basics

~20 s

Merge join wins when the order is already there or cheap: both sides indexed or clustered on the join key, memory too small to hold a build side, an ordering needed downstream, or a large-versus-large join whose hash build would spill anyway. Hash join wins for one-shot joins on unordered inputs with a build side that fits.

open as a page

For a system that runs both short transactional queries and large reporting queries against the same relational engine, how would you decide how much working memory to allow hash join operators, and what are the consequences of setting that budget too high or too low?

level: principalimportance: nice to knowfreq 28%

basics

~20 s

Budget from total memory minus the buffer cache, divided by realistic peak concurrency and the number of memory-hungry operators per query — not per query in isolation. Too low means constant spilling and IO; too high means a few concurrent reports exhaust memory, evict the cache, or crash the process.

open as a page

An execution engine can interpret a plan operator by operator, or compile the plan into machine code at runtime with adjacent operators fused into one loop. How would you decide which strategy a new engine should adopt, and what does each cost?

level: principalimportance: nice to knowfreq 22%

basics

~20 s

Decide by workload and engineering budget. Compilation fuses operators into register-resident loops and wins on long, compute-heavy analytical queries, but costs milliseconds of compile time per query plus hard debugging. Vectorized interpretation gets most of the win with precompiled primitives, no warm-up, and far less complexity — the safer default.

open as a page

You own a database serving both interactive OLTP traffic and heavy batch reporting on the same instance. How do you decide the per-operation memory budget for sorts and hash aggregates (PostgreSQL's work_mem, SQL Server's memory grants, Oracle's PGA target), and what policy keeps one workload from starving the other?

level: principalimportance: nice to knowfreq 32%

basics

~20 s

Budget from the top down: total RAM minus OS and buffer pool leaves a working-memory pool; divide by realistic peak concurrency times operators-per-plan to get a conservative default. Then differentiate by workload — small default for OLTP, elevated per-role or per-session for batch — and treat spilling as acceptable degradation, not failure.

open as a page