Analytic functions, also called window functions, compute a value for every row from a set of related rows without collapsing them the way GROUP BY does. Running totals, rankings, the previous event of the same user, a moving average: all of them are one analytic call in Impala instead of a self-join. They are also a common source of slow queries and quietly wrong answers, because the defaults are subtle and Impala supports a narrower window syntax than some other engines.

This article explains the model from first principles, lists exactly what Impala supports and what it rejects, works through the queries people actually write, and shows how the executor runs them so you can predict memory use and skew. The general window-function semantics shared with Hive are covered in Hive window functions; this page is about Impala's implementation, and it assumes the daemon roles described in Impala architecture.

Advertisement

From GROUP BY to OVER

Start with a table of orders. SELECT customer_id, SUM(amount) FROM orders GROUP BY customer_id returns one row per customer; the individual orders are gone. Now suppose you want every order and the customer's total next to it. Without analytic functions you would aggregate in a subquery and join it back. With them you write SUM(amount) OVER (PARTITION BY customer_id): the aggregate is computed per partition, and every input row survives with the result attached.

An analytic call has up to three parts inside OVER:

  • PARTITION BY divides rows into independent groups. Omit it and the whole result set is one partition.
  • ORDER BY orders rows within each partition. It is required for ranking and offset functions and gives running aggregates their meaning.
  • Window frame (ROWS or RANGE BETWEEN ... AND ...) selects which rows relative to the current one feed the aggregate.

-- orders(customer_id, order_id, order_ts, amount)
SELECT customer_id,
       order_id,
       amount,
       SUM(amount)  OVER (PARTITION BY customer_id ORDER BY order_ts
                          ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS running_total,
       ROW_NUMBER() OVER (PARTITION BY customer_id ORDER BY order_ts, order_id) AS nth_order,
       LAG(order_ts) OVER (PARTITION BY customer_id ORDER BY order_ts)          AS prev_order_ts
FROM orders;

This one query produces, for each order, the customer's running spend, the order's sequence number and the timestamp of the previous order, which is the raw material for inter-purchase time. Each of these would otherwise be a self-join.

The default frame, and why ties bite

When an analytic call has ORDER BY but no frame, the default is RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW. The word RANGE matters: it includes every row whose ORDER BY value equals the current row's, its peers. With ROWS, the frame ends at the physical current row. The difference only appears on ties, which is why it slips through testing on clean data:

-- Two orders share order_ts = 10:00 for the same customer, amounts 5 and 7.
-- Default frame when ORDER BY is present: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
SUM(amount) OVER (PARTITION BY customer_id ORDER BY order_ts)              -- 12, 12 (peers included)
SUM(amount) OVER (PARTITION BY customer_id ORDER BY order_ts
                  ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)        -- 5, 12 (order among peers arbitrary)
SUM(amount) OVER (PARTITION BY customer_id ORDER BY order_ts, order_id
                  ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)        -- deterministic

The rule to adopt: for running totals, write ROWS explicitly and make ORDER BY unique by adding a tie-breaker column such as a primary key. The same tie-breaker rule applies to ROW_NUMBER: if two rows tie on the ORDER BY key, Impala may number them differently in two runs, and a dedup query built on it becomes nondeterministic.

A second default catches everyone once. LAST_VALUE(x) OVER (PARTITION BY k ORDER BY t) returns the current row's value (or the last peer's), not the partition's last value, because the default frame ends at the current row. Write RANGE BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING, or use FIRST_VALUE with the order reversed.

Advertisement

What Impala supports

FunctionWhat it returnsWindow clause
ROW_NUMBER, RANK, DENSE_RANKPosition in the ordered partition; RANK leaves gaps after ties, DENSE_RANK does notNot allowed; ORDER BY required
NTILE(n), CUME_DIST, PERCENT_RANKBucket number, cumulative distribution, relative rankNot allowed; ORDER BY required
LAG, LEAD (expr, offset, default)Value from a row before or after the current oneNot allowed; ORDER BY required
FIRST_VALUE, LAST_VALUEFirst or last value in the frameAllowed; ORDER BY required
SUM, COUNT, AVGAggregate over the frameAllowed
MIN, MAXAggregate over the frameAllowed only if the frame starts at UNBOUNDED PRECEDING

Impala's documentation adds restrictions that shape how you write queries:

  • RANGE supports only three combinations: UNBOUNDED PRECEDING to CURRENT ROW (the default), CURRENT ROW to UNBOUNDED FOLLOWING, and UNBOUNDED PRECEDING to UNBOUNDED FOLLOWING. RANGE with numeric offsets, the value-based sliding window such as a seven-day interval, is not supported. Use ROWS over a gap-free series instead.
  • Placement. Analytic calls are allowed only in the SELECT list and the outermost ORDER BY. You cannot use them in WHERE, GROUP BY or HAVING, and there is no QUALIFY clause, so filtering on a rank needs an inline view or WITH clause.
  • No DISTINCT inside. DISTINCT cannot be combined directly with an analytic call; aggregate first or wrap the analytic result and apply DISTINCT outside.

Worked examples you will actually write

Top N per group and latest record per key. Both are ROW_NUMBER in an inline view with a filter outside. The dedup form is the standard way to collapse change-data-capture rows into current state; note the second ORDER BY column that makes it deterministic when two changes carry the same timestamp.

-- Top 3 products by revenue in each category. Analytic calls are not allowed in WHERE,
-- so compute the rank in an inline view and filter outside it.
SELECT category, product_id, revenue
FROM (
  SELECT category, product_id, revenue,
         ROW_NUMBER() OVER (PARTITION BY category ORDER BY revenue DESC, product_id) AS rn
  FROM product_revenue
) ranked
WHERE rn <= 3;

-- Keep only the latest version of each record (deduplication of CDC data)
SELECT *
FROM (
  SELECT c.*,
         ROW_NUMBER() OVER (PARTITION BY c.customer_id ORDER BY c.updated_at DESC, c.change_seq DESC) AS rn
  FROM customer_changes c
) latest
WHERE rn = 1;

Moving averages without RANGE offsets. Because Impala rejects RANGE BETWEEN INTERVAL 6 DAYS PRECEDING, the frame must count rows. ROWS BETWEEN 6 PRECEDING AND CURRENT ROW equals seven days only if there is exactly one row per day, so build a dense series by cross-joining a calendar table and left-joining the facts. Skipping this step silently turns a 7-day average into a 7-active-days average.

-- 7-day moving average of daily revenue per store.
-- RANGE with a numeric offset is not supported, so use ROWS over a gap-free daily series.
WITH days AS (
  SELECT s.store_id, d.day
  FROM stores s CROSS JOIN calendar d
  WHERE d.day BETWEEN '2026-09-01' AND '2026-09-29'
),
daily AS (
  SELECT dy.store_id, dy.day, COALESCE(SUM(o.amount), 0) AS revenue
  FROM days dy
  LEFT JOIN orders o ON o.store_id = dy.store_id AND to_date(o.order_ts) = dy.day
  GROUP BY dy.store_id, dy.day
)
SELECT store_id, day, revenue,
       AVG(revenue) OVER (PARTITION BY store_id ORDER BY day
                          ROWS BETWEEN 6 PRECEDING AND CURRENT ROW) AS avg_7d
FROM daily;

Sessionization. The two-step pattern flags each event that starts a new session with LAG, then turns the flags into session numbers with a running SUM. The partition filter sits inside the inner query so that partition pruning happens before any rows are shuffled.

-- Sessionize clickstream: a new session starts after 30 minutes of inactivity.
SELECT user_id, event_ts,
       SUM(is_new) OVER (PARTITION BY user_id ORDER BY event_ts, event_id
                         ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS session_no
FROM (
  SELECT user_id, event_id, event_ts,
         CASE WHEN LAG(event_ts) OVER (PARTITION BY user_id ORDER BY event_ts, event_id) IS NULL
                OR unix_timestamp(event_ts)
                   - unix_timestamp(LAG(event_ts) OVER (PARTITION BY user_id ORDER BY event_ts, event_id)) > 1800
              THEN 1 ELSE 0 END AS is_new
  FROM clicks
  WHERE event_date = '2026-09-29'          -- prune partitions BEFORE windowing
) flagged;

Working around the DISTINCT restriction.

-- Not allowed: COUNT(DISTINCT page) OVER (PARTITION BY user_id)
-- Workaround: aggregate first, then join or window over the aggregate
WITH per_user AS (
  SELECT user_id, COUNT(DISTINCT page) AS distinct_pages
  FROM clicks GROUP BY user_id
)
SELECT c.user_id, c.page, p.distinct_pages
FROM clicks c JOIN per_user p ON p.user_id = c.user_id;

-- DISTINCT over analytic output: compute in an inline view, apply DISTINCT outside
SELECT DISTINCT customer_id, first_order_ts
FROM (SELECT customer_id,
             FIRST_VALUE(order_ts) OVER (PARTITION BY customer_id ORDER BY order_ts) AS first_order_ts
      FROM orders) v;

How Impala executes an analytic query

Understanding the physical plan explains most performance problems. Conceptually, and visible in EXPLAIN output as an exchange, a sort and an analytic node, the steps are:

  1. Executors scan and filter. WHERE predicates, partition pruning and runtime filters all apply here, before any window work.
  2. Rows are redistributed with a hash exchange on the PARTITION BY columns, so all rows of one partition land on one executor.
  3. Each executor sorts its rows by the partition keys and the ORDER BY columns.
  4. The analytic evaluator streams through the sorted rows, resetting state at each partition boundary and emitting one result per row.

Several consequences follow. The sort is the expensive, memory-hungry step, and if it exceeds its memory reservation it spills to disk, which is correct but slow; see Impala spill to disk and Impala memory limits. A query with no PARTITION BY sends every row to a single node, so a global ROW_NUMBER over a billion rows runs on one executor no matter how large the cluster is. And one heavy partition key, such as a bot user with millions of clicks, is processed by one executor while the rest sit idle.

When several analytic calls share the same PARTITION BY and ORDER BY, the planner can usually evaluate them over a single sort; calls with different specifications need additional sorts. Check EXPLAIN rather than assuming, and where you can, align the specifications of calls in one SELECT.

How an analytic function executes across Impala executors (simplified)Scan + filterWHERE runs firstHash exchangeon PARTITION BY keysSortpartition keys, ORDER BYAnalytic evaluationone pass per partitionSpill to diskif sort exceeds memoryOuter queryfilter on the resultNo PARTITION BYall rows go to one node: no parallelismRows for one partition key always meet on one executor, so one giant key means one busy node
Filtering happens before the exchange; rows are hash-distributed on the PARTITION BY keys, sorted, and evaluated in one pass. Without PARTITION BY every row flows to one node.

Performance guidance

  • Reduce rows before windowing. Filter partitions and pre-aggregate in an inner query. A window over daily aggregates is thousands of times cheaper than one over raw events.
  • Always partition when the logic allows it. A global ranking can often be replaced by a per-day or per-region ranking followed by a small final merge.
  • Watch skew. Query the row count per partition key before running a large window. Handle outliers separately or add a salt when the function tolerates it, which sums and counts do but rankings do not.
  • Select only needed columns. Sorted rows are materialised, so wide rows multiply memory use and spill volume.
  • Keep statistics current. COMPUTE STATS helps the planner size joins and exchanges around the analytic step, even though the window itself is not reordered.
  • Read the profile. In the query profile, compare the time and peak memory of the sort and analytic nodes across executors; one executor far above the others is skew.

Failure modes

SymptomCauseFix
Dedup keeps a different row each runROW_NUMBER ORDER BY has tiesAdd a unique tie-breaker column
Running total jumps on duplicate timestampsDefault RANGE frame includes peersWrite ROWS BETWEEN explicitly
LAST_VALUE equals the current valueDefault frame ends at the current rowUse an UNBOUNDED FOLLOWING frame
Moving average wrong around gapsROWS counts rows, not daysDensify with a calendar table
Query fails with a memory limit error or crawlsSort too large, spilling, or skewed partitionFilter earlier, narrow columns, fix skew, raise the limit last
Analysis error on an analytic call in WHEREPlacement restrictionMove the call into an inline view

NULLs in ORDER BY are a quieter trap: whether NULL sorts first or last affects rankings and LAG results. State NULLS FIRST or NULLS LAST explicitly whenever the ordering column can be NULL, rather than relying on the engine default, especially when the same SQL also runs on Hive or Spark.

Trade-offs and portability

Analytic functions nearly always beat the self-join they replace, but they are not free: every distinct window specification costs a shuffle and a sort. When a result is needed repeatedly, such as a customer's lifetime order number, it is often cheaper to compute it once during ingestion and store it as a column. For portability, the RANGE restrictions are the main difference to plan around: SQL that uses value-based RANGE offsets on another engine must be rewritten with ROWS and a dense series before it runs on Impala, and the result should be checked row for row.

What to do next

  1. Grep your Impala SQL for ROW_NUMBER and RANK and confirm every ORDER BY has a unique tie-breaker.
  2. Replace implicit frames in running totals with explicit ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW.
  3. Build or verify a calendar table and use it to densify any series fed to a ROWS-based moving window.
  4. Run EXPLAIN on your three heaviest analytic queries and count the sorts; align OVER clauses where possible.
  5. Check row counts per PARTITION BY key for skew before scheduling large window jobs.
  6. Push partition filters and pre-aggregation into inner queries, then compare peak memory in the profile before and after.
Key takeaway: An Impala analytic function computes a per-row value over a partition, ordered and framed, without collapsing rows. Impala supports the standard ranking, offset and aggregate functions but restricts RANGE to three unbounded or current-row combinations, limits MIN and MAX frames, and allows analytic calls only in the SELECT list and outermost ORDER BY. Execution is a hash exchange on the partition keys, a sort and a streaming pass, so filter early, always partition, beware skew, and make ORDER BY deterministic with a tie-breaker.