Impala is usually described as an in-memory engine, and the phrase is misread more often than it is understood. Impala does not load your tables into RAM and keep them there between queries. Every query reads its data from HDFS, S3, Ozone or Kudu when it runs. What is in memory is the intermediate data: rows travel from the scanner, through filters, joins and aggregations, across the network and back to the client without being written to disk in between, unless an operator runs out of memory and has to spill.

That design is why Impala answers interactive queries in seconds, and why its failure modes are memory failures. This article follows a row through the engine: row batches and tuples, the pull model, streaming and blocking operators, exchanges between hosts, and memory accounting. It ends with a worked diagnosis, failure modes and a checklist.

Advertisement

What in-memory execution means, and what it does not

Impala has three ways to avoid re-reading storage, and none of them is the execution engine: the operating system page cache, HDFS centralized caching for pinned tables, and the optional data cache that keeps recently read file ranges on each executor's local disk. All three cache file bytes. They make scans cheaper; they do not change how the query executes.

In-memory execution means that a plan is a pipeline of operators running concurrently on every executor, handing data to each other in memory buffers rather than files. A Hive query on MapReduce wrote every stage's output to HDFS for the next stage to read back. Impala has no such boundaries: a scan produces rows into memory, a join consumes them, an exchange streams them into another host's memory, and the coordinator streams final rows to the client. Memory is therefore the resource you plan around.

Row batches and tuples

The unit of data inside Impala is the row batch: a group of rows that an operator produces in one call and hands to the next. Batching is what makes the per-row cost small. A virtual function call, a lock acquisition or a memory-limit check costs the same whether it covers one row or a thousand, so the engine pays those costs once per batch. The BATCH_SIZE query option controls the row count; its default value of 0 means the built-in default of 1024 rows, and the documented maximum is 65,536.

Inside a batch, a row is not a self-contained record. It is an array of pointers to tuples, one pointer per table or intermediate result that contributes to the row. A tuple is a contiguous block of memory whose layout the planner fixes in advance: each column the query needs becomes a slot at a known byte offset, and null-ness is tracked in indicator bits rather than by sentinel values. Fixed-width values such as INT, DOUBLE and TIMESTAMP sit directly in their slots. Variable-length values such as STRING store a pointer and a length in the slot, and the bytes themselves live in a memory pool owned by the batch.

Two properties follow. The engine never materialises unused columns, because tuples only have slots for referenced ones. And a join does not copy data: its output row is two tuple pointers, one into each side, so joining a wide fact row to a dimension row costs a few pointer writes.

-- Row descriptor for: SELECT f.amount, d.region FROM sales f JOIN stores d ON f.store_id = d.id
--
-- row  = [ tuple* sales_tuple , tuple* stores_tuple ]
-- sales_tuple  : null bits | store_id (INT) | amount (DECIMAL)       -- fixed layout, slot offsets known at plan time
-- stores_tuple : null bits | id (INT)       | region (ptr, len) ---> bytes in the batch's memory pool

Because slot offsets are fixed at plan time, the per-row code that evaluates predicates and hashes join keys can be generated with those offsets baked in. That is the job of LLVM code generation, covered in Impala codegen; the tuple layout is what makes it pay off.

Advertisement

The pull model: Open, GetNext, Close

Operators, called exec nodes, form a tree inside each plan fragment, and they talk to each other through three calls. Open prepares the node and, for some nodes, does a great deal of work. GetNext fills a row batch supplied by the caller. Close releases everything. The root of the fragment calls GetNext on its child, which calls its own child, down to the scan at the leaf. Data is pulled upward one batch at a time, and each node returns as soon as its output batch is full or its input is exhausted.

// Simplified shape of a streaming operator (a filter). Not the literal Impala source.
Status SelectNode::GetNext(RuntimeState* state, RowBatch* out, bool* eos) {
  while (!out->AtCapacity()) {
    if (child_batch_pos_ == child_batch_->num_rows()) {     // need more input
      if (child_eos_) { *eos = true; return Status::OK(); }
      child_batch_->Reset();
      RETURN_IF_ERROR(child(0)->GetNext(state, child_batch_.get(), &child_eos_));
      child_batch_pos_ = 0;
      RETURN_IF_CANCELLED(state);                          // cancellation checked per batch
      RETURN_IF_ERROR(state->CheckQueryState());           // memory limit checked per batch
    }
    TupleRow* row = child_batch_->GetRow(child_batch_pos_++);
    if (EvalConjuncts(row)) out->CopyRowFrom(row);         // copies tuple pointers, not tuples
  }
  return Status::OK();
}

Cancellation and the memory limit are checked once per batch, not once per row. And because a filter's output points into its input's tuples, the input's memory stays charged to the query until the rows referencing it have moved on; batches pass ownership of their memory upward with their rows.

Scans are the exception to pure pull. I/O threads and scanner threads work ahead of the consumer, decoding Parquet, ORC or text into a queue of ready batches, and GetNext takes the next one. That read-ahead overlaps I/O with computation, and it is memory outside the spillable operators, which matters later.

HDFS / S3 / Kudufile bytes, page cacheSCAN (fragment F01)I/O threads + scannersEXCHANGE senderhash-partition rowsEXCHANGE receiverper-sender queuesSCAN dim (build)broadcast or shuffleHASH JOIN (F02)build: blocking, probe: streamingAGGREGATEblocking, spillableCOORDINATORresult spooling, client fetchBUFFER POOLreservations per operatorread rangesrow batchesnetwork (KRPC)probe rowsbuild rowsjoined batchesfinal rowspagesBatches stream upward; only the build side, aggregation and sorts hold data, and those draw from reserved buffers
Data flow for a broadcast join and aggregation: row batches stream between operators and across exchanges; blocking operators draw buffers from the buffer pool.

Streaming versus blocking operators

Whether a query fits in memory depends almost entirely on which operators must hold their input and which can pass it straight through. A streaming operator consumes a batch and emits a batch, holding only a bounded amount of state. A blocking operator, a pipeline breaker, must see some or all of its input before it can emit anything, and must keep what it has seen.

OperatorBehaviourWhat it holds in memorySpillable
Scan, filter, project, limitStreamingIn-flight batches and read-ahead buffersNo (bounded)
Exchange (send / receive)StreamingQueued batches per sender and receiverNo (bounded by queue limits)
Hash joinBuild side blocking, probe side streamingEntire build input as a hash tableYes
Aggregation (GROUP BY)BlockingOne entry per distinct groupYes
Pre-aggregationStreaming, best effortPartial groups; passes rows through when it stops helpingNot needed
Sort (ORDER BY without LIMIT)BlockingAll input rowsYes
Top-N (ORDER BY with LIMIT)BlockingOnly N rowsNot needed
Analytic functionsBlocking per partition windowRows of the current partitionYes

Bring this table to every EXPLAIN. Most Impala memory incidents are a blocking operator holding something far larger than the author expected.

A join's Open consumes the whole build input and builds a hash table before GetNext streams the probe side through it. Runtime filters exploit that ordering: once the build is complete, Impala ships Bloom or min-max filters of its keys to the probe-side scans, which skip rows, row groups or files before reading them.

Exchanges: in-memory execution across hosts

A distributed plan is cut into fragments wherever data must move between hosts, and each cut becomes an exchange. The sender partitions batches by hash of the key, or broadcasts them, serialises them and sends them over KRPC; the receiver keeps a queue per sender and hands deserialised batches to its parent through GetNext.

Broadcast exchanges replicate the whole build side, so a broadcast join's hash table costs its full size on every host. Partitioned exchanges give each host its share, but only if keys are evenly spread: a hot key sends all its rows to one host. Queues are bounded, so a slow consumer pushes back on senders rather than growing memory, and one slow host slows the whole fragment.

At the top, the coordinator returns rows to the client. With result spooling enabled through the SPOOL_QUERY_RESULTS query option (check its default on your release), finished rows are buffered on the coordinator so executors can release resources even if the client fetches slowly.

Where the memory comes from

Every allocation in a query is charged to a tree of memory trackers: the process, the resource pool, the query, each fragment instance, and each exec node. The query-level limit comes from MEM_LIMIT or the resource pool's settings, as described in Impala memory limits. When a charge would push a tracker past its limit, the allocation fails and the query fails with a memory-limit-exceeded error that names the tracker. The profile's per-node peak memory figures come from the same trackers.

The operators that can hold unbounded data, namely hash joins, aggregations, sorts and analytic functions, do not allocate freely. They take fixed-size buffers from the buffer pool against a reservation. Before a query starts, the planner computes the minimum reservation every spillable operator needs to make progress in the worst case, and Impala claims that amount up front. You can see it in the EXPLAIN header as the per-host resource reservation, next to the per-host memory estimate. The default spillable buffer size is 2 MB, controlled by DEFAULT_SPILLABLE_BUFFER_SIZE.

Above the minimum, an operator grows its reservation as needed. When it cannot, it writes a hash partition or sorted run to scratch disk and frees the buffers. That is spilling to disk, and only reservation-based operators do it. Scans, exchanges, string data and expression evaluation cannot spill, which is why a query can still hit a memory error with spilling enabled.

Worked example: reading a plan for memory

Take a typical reporting query on a cluster of ten executors. The fact table sales has two billion rows; stores has five thousand.

SET MEM_LIMIT=4g;
EXPLAIN
SELECT d.region, SUM(f.amount) AS revenue
FROM sales f JOIN stores d ON f.store_id = d.id
WHERE f.sale_date >= '2026-09-01'
GROUP BY d.region;

-- Read the plan bottom-up and classify each node:
--   00:SCAN sales          streaming; partition pruning on sale_date; runtime filter on store_id
--   01:SCAN stores         streaming; feeds the build side
--   04:EXCHANGE BROADCAST  stores rows copied to every executor (5,000 rows: cheap)
--   02:HASH JOIN           build holds 5,000 tuples per host; probe streams the fact rows
--   03:AGGREGATE STREAMING pre-aggregation per host: a handful of regions
--   05:EXCHANGE HASH(region)
--   06:AGGREGATE FINALIZE  blocking but tiny: one entry per region

After running the query, type SUMMARY in impala-shell. The columns that matter are #Rows, Est. #Rows, Peak Mem and Est. Peak Mem for each operator. In a healthy run, the join's peak memory is a few megabytes, the scan's peak is its read-ahead buffers, and the aggregations are negligible. Time is dominated by the scan, which is the right place for it.

Now a colleague rewrites it as FROM stores d JOIN sales f with a STRAIGHT_JOIN hint. The build side is now the filtered fact rows; even partitioned, each host builds a hash table on about a tenth of them, gigabytes instead of megabytes. The join exhausts its reservation, spills, re-reads from scratch disk and runs many times slower, or fails without scratch space. The profile shows the join's peak memory near the limit and spilled partitions on that node.

The diagnosis generalises: list the blocking operators, estimate what each holds, compare with the per-host limit, and check estimates against reality. If Est. #Rows is off by orders of magnitude, run COMPUTE STATS before any memory tuning.

Failure modes

The failure modes follow directly from the design.

  • Build side larger than expected. Missing statistics make a large table look small, so it is broadcast or built. The join has the largest Peak Mem and Est. #Rows far below #Rows. Compute statistics; hint only as a last resort.
  • Memory error despite spilling. The failing tracker is a scan, exchange or expression. Causes: wide rows with long strings, many scanner threads on wide Parquet files, large GROUP_CONCAT results. Project fewer columns or raise the limit.
  • Skew. A hot key puts a partitioned join or aggregation on one host, which spills or fails while others idle; Max Time far exceeds Avg Time. Pre-aggregate, salt the key, or handle the hot key separately.
  • Spill storms. Concurrent spilling saturates scratch disks. Admit fewer queries by memory and use fast local scratch disks.
  • Slow clients. Without result spooling, a slow fetcher keeps executors' memory alive. Enable spooling and set fetch timeouts.

Trade-offs

Pipelined execution buys latency at the cost of resilience. A batch engine that materialises every stage can retry a failed stage; Impala restarts the query, because no intermediate state is durable. That suits queries of seconds or minutes, not multi-hour transformations, which belong on Hive or Spark. Spilling turns would-be failures into slow successes but depends on scratch disk throughput, and reservations claimed up front reduce concurrency in exchange for predictable admission.

What to do next

  1. Run EXPLAIN on your five most expensive queries and mark every blocking operator and what it holds.
  2. Run each query, then SUMMARY; flag every operator whose #Rows differs from Est. #Rows by more than ten times, and compute statistics for the tables involved.
  3. Check the per-host resource reservation in each plan header against your pool's memory limit.
  4. Confirm scratch directories are configured on fast local disks and watch for spilled-partition counters in profiles.
  5. Verify result spooling is enabled and clients fetch promptly.
  6. Read Impala query execution for the coordinator's side of the lifecycle, then set per-pool memory limits from what your profiles show.
Key takeaway: Impala's in-memory execution means intermediate rows stream through a pipeline of operators in memory rather than being written to disk between stages; tables are still read from storage on every query. Rows move in batches of tuple pointers, operators pull batches through Open, GetNext and Close, and only blocking operators such as hash-join builds, aggregations and sorts hold large state, drawing it from buffer-pool reservations that can spill. Diagnose memory by classifying operators in EXPLAIN, comparing estimates with actual rows in SUMMARY, and fixing statistics before tuning limits.