A partitioned table stores each value of its partition key in its own directory: one directory per day, per region or per tenant. Partition pruning is the planner's decision to read only the directories a query can possibly need. On a table with three years of daily partitions, a query for last week touches 7 of 1,095 directories, and the other 99 percent of the data is never opened, never listed and never scheduled. No index, file format or cache gives a larger saving, which is why the first question about any slow Impala scan is whether it pruned.

Pruning in Impala happens mostly at plan time. The coordinator compares the query's predicates with the partition values it holds in catalog metadata, keeps the partitions that might match, and builds scan ranges only for their files. A second, run-time form uses join results to skip more partitions as the query executes. This page explains how both work, which predicates prune and which silently do not, how to confirm pruning in EXPLAIN, how stale metadata turns pruning into missing rows, and how to choose keys that prune well. Facts come from the Apache Impala documentation.

How static pruning works

When Impala plans a query against a partitioned table, it splits the WHERE clause into conjuncts, the pieces joined by AND. A conjunct that refers only to partition key columns and constants can be decided per partition: the planner checks it against each partition's key values and drops partitions where it is false. Conjuncts that refer to ordinary columns are left for the scan to evaluate row by row. Only the surviving partitions contribute files to the plan, so pruning reduces planning work for scan ranges, the number of files opened, the bytes read and the number of hosts involved.

This depends entirely on what the catalog knows. The planner does not list storage directories at query time; it uses partition and file metadata cached from the Hive Metastore and the file system. A partition directory that exists in storage but not in metadata is invisible, and a partition in metadata whose files changed is read from a stale file list. That dependency is the source of the most damaging pruning bugs, covered below.

Static pruning happens in the planner; dynamic pruning happens in the scansSQL queryWHERE dt BETWEEN ...Planner (coordinator)evaluate partition predicatesCatalog metadatapartition list, files, statspartitions7 of 1,095 keptScan rangesfiles of kept partitions onlyExecutors: SCAN HDFS nodesread Parquet/ORC from HDFS or S3Hash join build sideproduces RF000 on dtruntime filterDynamic pruning: a scan skips a partition when itskey value fails the filter it received from the joinPartitions unknown to the catalog are never scanned
Static pruning removes partitions before any scan range is built. Dynamic pruning, driven by runtime filters from a join, removes more while the query runs.

Which predicates prune

Direct comparisons prune: equality, ranges, BETWEEN and IN lists. Conjunctions prune further, so a predicate on year and a range on month together keep only the matching months of that year. Impala also applies predicate propagation, documented since release 1.2.2: if a query says the partition column equals another column and that column equals a constant, the planner infers the constant for the partition column and prunes on it.

-- Table: events(user_id BIGINT, kind STRING, payload STRING)
--        PARTITIONED BY (year INT, month INT, day INT), stored as Parquet

-- Prunes to one day
SELECT count(*) FROM events WHERE year = 2026 AND month = 9 AND day = 30;

-- Prunes to one quarter
SELECT count(*) FROM events WHERE year = 2026 AND month BETWEEN 7 AND 9;

-- Prunes by propagation: census_year is pinned, so year is pinned too
SELECT count(*) FROM events e JOIN census c ON e.year = c.census_year
WHERE c.census_year = 2025;

-- Does NOT prune: the OR mixes a partition column with a data column,
-- so no partition can be ruled out before reading rows
SELECT count(*) FROM events WHERE year = 2026 OR kind = 'refund';

-- Does NOT prune statically: the value is only known after a subquery runs.
-- It can prune dynamically through a runtime filter (see below)
SELECT count(*) FROM events
WHERE year IN (SELECT max(year) FROM fiscal_calendar);

The OR case is the most common surprise. Rewrite it as a UNION ALL of two queries, each with its own prunable predicate, or restrict the OR to partition columns only. Expressions over partition columns, such as concatenating year and month into a string, are evaluated per partition by the planner only when the conjunct references nothing but partition columns and constants; whether a given expression prunes is easy to check and not worth guessing, so check it in EXPLAIN. Two further rules come from the Impala documentation. Do not partition on a TIMESTAMP column; split the time into year, month and day columns instead. And an analytic function prunes only on columns named in its PARTITION BY clause, so a query that filters on year above a window function should include year in that clause.

Confirming pruning in EXPLAIN

Never assume pruning; read the plan. The scan node in EXPLAIN reports the partitions it will read out of the total. Recent releases print it as partitions=kept/total along with the file count and size; older documentation shows the same figure as #partitions=kept/total.

EXPLAIN SELECT kind, count(*) FROM events
WHERE year = 2026 AND month = 9 GROUP BY kind;

...
00:SCAN HDFS [analytics.events]
   partition predicates: year = 2026, month = 9
   HDFS partitions=30/1095 files=240 size=61.20GB
   row-size=12B cardinality=...

That fragment is an abridged, illustrative shape, not captured output, but the fields are the ones to check. If the kept count equals the total, nothing was pruned; look for an OR, a predicate on the wrong column, a type mismatch or a value only known at run time. If the kept count is right but files and size are much larger than expected, pruning worked and the problem is partition size or small files. The bigger guide to plan reading is Impala query plans. After running a query, the profile shows what each scan actually read, which is the way to confirm dynamic pruning.

Dynamic pruning with runtime filters

Static pruning needs constants at plan time. Many real queries filter a fact table through a dimension instead: join sales to a date dimension and keep rows where the date is a holiday. The planner cannot know which dates those are. Dynamic partition pruning, available since Impala 2.5, handles this with runtime filters. The hash join builds its table from the small side, produces a filter over the join key values it saw, and sends it to the scans of the large side. A scan whose partition key value cannot pass the filter skips that partition's files entirely.

02:HASH JOIN [LEFT SEMI JOIN, BROADCAST]
|  runtime filters: RF000 <- year
|
00:SCAN HDFS [default.yy]
|  runtime filters: RF000 -> year

The arrow pointing into the join marks the filter's producer; the arrow pointing out of the scan marks its consumer. The query option RUNTIME_FILTER_MODE chooses LOCAL or GLOBAL scope, and GLOBAL has been the default since Impala 2.6, so filters can travel between hosts. RUNTIME_FILTER_WAIT_TIME_MS bounds how long a scan waits for filters before starting without them. Two limits matter in practice: a host whose join spilled to disk does not produce filters for that join, and a scan that starts before its filter arrives reads partitions it could have skipped. The full mechanics, including Bloom and min-max filters, are in Impala runtime filters.

Metadata: pruning can only keep what the catalog knows

Because the planner trusts metadata, partitions that metadata does not know about are not slow; they are absent. A Spark job that writes year=2026/month=10/day=3 directly to storage, or a Hive job on a cluster whose events Impala has not processed, produces data that an Impala query for that day returns as zero rows, with no error. This is the single most expensive pruning failure, because the result looks plausible.

-- New partition directories written outside Impala: discover them
ALTER TABLE events RECOVER PARTITIONS;

-- New or replaced files inside a known partition: reload just that partition
REFRESH events PARTITION (year = 2026, month = 10, day = 3);

-- Table changed in ways REFRESH does not cover (schema, many partitions): heavier
INVALIDATE METADATA events;

Clusters that run automatic metastore event processing pick up partitions added through the Hive Metastore, but files written straight to storage without a metastore call still need RECOVER PARTITIONS or REFRESH. Make the writer responsible: the job that lands a partition registers it and refreshes it as its last step, and a freshness check compares expected and visible partitions. Statistics are metadata too; see Impala table statistics, since a pruned scan with stale row counts can still get a bad join order.

Designing keys that prune

Pruning only helps if queries filter on the key. Choose partition columns from the predicates that real queries use, nearly always time and sometimes a coarse tenant or region. Store keys as integers or as zero-padded strings. A STRING month partition written as 03 by one job and 3 by another creates two partitions for the same month, and a query for month = '03' silently misses half the data. With integer keys that cannot happen. Rows whose key is NULL land in the Hive default partition, which an equality predicate excludes and an IS NULL predicate selects, so decide deliberately whether that partition should exist.

Granularity is the other lever. Every partition costs catalog memory, metastore rows and planning time, and every partition holds at least one file, so very fine keys produce many small files that read slowly even when pruned perfectly. A useful test is the bytes per partition after compaction: partitions that hold only a few small files are too fine, and partitions that every query must read most of are too coarse. Hive-side layout and the arithmetic of choosing keys are covered in Hive table partitions.

Worked example: a 90-second dashboard tile

An analytics team runs a dashboard over a 1,095-partition events table partitioned by year, month and day, about 2 GB per day. The slowest tile shows refunds per day for the last 30 days and takes 90 seconds. Its query filters with WHERE to_date(event_ts) >= date_sub(now(), 30), where event_ts is an ordinary TIMESTAMP column inside the files. EXPLAIN shows partitions=1095/1095: nothing pruned, about 2 TB scanned, because the predicate references a data column, not a partition key.

The rewrite adds partition predicates that the dashboard computes before sending the query. A 30-day window usually crosses a month boundary, so a single month range plus a day range is wrong: it either over-reads or, with the day bounds reversed, matches nothing. Instead it emits one conjunct per month, using partition columns only: (year = 2026 AND month = 9 AND day >= 4) OR (year = 2026 AND month = 10 AND day <= 3). The original timestamp predicate stays, ANDed on, for exact edges. EXPLAIN now shows partitions=30/1095 and roughly 60 GB. The tile drops to about 3 seconds. A week later the tile shows zero refunds for today. The ingest job had switched to writing files directly to storage, so today's partition was missing from metadata; adding RECOVER PARTITIONS and a partition REFRESH to the end of the ingest job fixed it, and a freshness alert on the newest visible partition now catches the next regression.

Failure modes

  • Predicate on a data column. Filtering on a timestamp inside the files instead of the partition columns reads every partition. Add explicit partition predicates.
  • OR across partition and data columns. No partition can be excluded. Split into UNION ALL or restructure the predicate.
  • Inconsistent string keys. '3' and '03' become different partitions; queries miss data silently. Use integer keys or enforce zero padding at write time.
  • Invisible partitions. Data written outside Impala is not in metadata, so queries return zero rows. Register and refresh as part of the write.
  • Over-partitioning. Hundreds of thousands of tiny partitions slow planning and metadata loading and leave small files. Coarsen the key and compact.
  • Late runtime filters. A scan starts before the filter arrives, or the producing join spills, so dynamic pruning does not happen. Check the profile, not just the plan.

Trade-offs

ChoiceGainCost
Finer partitions (day, hour)Short time-range queries read lessMore files, metadata and planning time
Coarser partitions (month)Fewer, larger files; cheaper metadataShort queries read more
Second key (tenant, region)Prunes per-tenant queriesPartition count multiplies
Rely on dynamic pruningWorks with dimension filtersDepends on filter timing and join spills
Precompute partition bounds in the clientReliable static pruningQuery generation logic to maintain

What to do next

  1. Pick the ten most expensive queries on your largest partitioned tables and run EXPLAIN on each; record kept/total partitions.
  2. For any query where kept equals total, find the predicate that should have pruned and rewrite it onto partition columns.
  3. Replace OR conditions that mix partition and data columns with UNION ALL or partition-only conditions.
  4. Check partition key types; move STRING month and day keys to integers or enforce zero padding at write time.
  5. Make every ingest job register and refresh the partitions it writes, and alert when the newest visible partition is older than expected.
  6. Measure bytes per partition; coarsen keys whose partitions hold only a few small files.
  7. For star-schema queries, confirm runtime filters appear in the plan and check the profile to see that partitions were skipped.
Key takeaway: Impala prunes partitions mostly at plan time by evaluating predicates on partition key columns against the partition values in catalog metadata, and at run time by using runtime filters from joins. Prune by filtering directly on partition columns with equality, ranges or IN lists, avoid OR conditions that mix in data columns, and confirm every important query with the partitions kept over total figure in EXPLAIN. Keep metadata current, because a partition the catalog does not know about returns no rows rather than slow rows, and choose integer keys at a granularity that leaves partitions holding large files.