Impala does not read Parquet through a Java library. Its executors, written in C++, include a Parquet scanner built for the format: it reads the footer, decides which row groups and pages can be skipped, decompresses and decodes only the column chunks a query needs, and hands row batches to the rest of the plan, with LLVM code generation specialising parts of the work for each query. That is what native means here, and it is why Parquet is the format Impala is fastest on.

The scanner is only as fast as the files allow. A query that should touch a few megabytes reads gigabytes when files carry no useful statistics, are unsorted, or hold thousands of tiny row groups. This page follows a scan from planning to row batches, names the query option behind each skipping level, explains schema resolution, and shows how to write files the reader can skip. The file format itself, including footers, encodings and page index structures, is explained in Parquet Format; ORC behaviour is in Impala ORC Support, in depth. Option names and defaults were checked against the Apache Impala and Cloudera documentation on 2026-10-02.

Advertisement

The scan pipeline

A Parquet scan happens in two places. The planner on the coordinator lists the files of the partitions that survive partition pruning, breaks them into scan ranges and assigns ranges to executors, preferring local replicas on HDFS. On object stores, where there are no blocks, the split size comes from PARQUET_OBJECT_STORE_SPLIT_SIZE, 256 MB by default. Each row group is read by exactly one scan range, so a file with one large row group is read by one scanner thread regardless of its size.

On the executor, the scanner reads the file footer first. The footer gives the schema, the row groups, the location of every column chunk and the statistics written for them. From this point the work is a series of filters, each cheaper than the next is expensive, followed by decoding of what survives.

Impala skips Parquet data at five levels before it decodes a valuePlannerfiles, splits, rangesFooterschema, row groupsRow-group filtermin/max, dict, bloomPage filterpage indexColumn readersdecompress, decode dict/RLEsurviving pagesLate materializationpredicate columns firstPredicates + runtime filtersconjuncts, join filtersRow batches to the rest of the plan: joins, aggregations, exchangesremaining columns are decoded only for rows that survivedPlanning picks files; the footer and statistics pick row groups and pages; predicates pick rows.Each level only works if the writer left something to skip: sorted data, statistics, dictionaries, page indexes.
The path of a Parquet scan in Impala, from planning to row batches. Each box is a chance to avoid reading or decoding data.

Two consequences follow. Every file the plan touches costs a footer read before any row is skipped, so small files are a reader problem too. And every skipping level depends on metadata the writer produced.

Five levels of skipping and their switches

LevelWhat it provesControlled byNeeds from the writer
PartitionWhole directories cannot matchPlanner; partition columns in WHEREA partition layout aligned with filters
Row group by statisticsMin and max of a column exclude the predicatePARQUET_READ_STATISTICS, default trueColumn statistics in the footer
Row group by dictionaryNo dictionary value satisfies the predicatePARQUET_DICTIONARY_FILTERING, default trueChunk fully dictionary encoded
Row group by bloom filterAn equality value is definitely absentBloom filter reading options; written via table propertyBloom filters for chosen columns
PagePage-level min and max exclude the predicatePARQUET_READ_PAGE_INDEX, default truePage index, PARQUET_WRITE_PAGE_INDEX on

Statistics filtering compares the predicate with each row group's min and max. It is extremely effective when data is sorted or clustered on the filtered column and almost useless when values are spread uniformly, because every row group's range then covers nearly the whole domain. Dictionary filtering helps exactly there for low-cardinality columns: if a row group's dictionary for country holds twelve values and none equals the one requested, the group is skipped even though its min and max straddle it. The scan node's profile reports skipped groups in NumDictFilteredRowGroups.

Bloom filters answer equality questions on high-cardinality columns, such as an ID, where neither ranges nor dictionaries help. They cost space and write time, so they are opt-in per column. The page index applies the same min and max logic per page, typically tens of kilobytes to about a megabyte, inside a row group. That matters because Impala files usually hold a single large row group; without the page index, a selective filter on a large file would still decode the whole group. The page index also lets the reader skip pages of other columns that hold only filtered-out rows.

Advertisement

Late materialization and predicate order

After skipping, the scanner decodes surviving pages. With late materialization, the scanner decodes predicate columns first, evaluates the predicates, and decodes the remaining columns only for rows that pass. The Cloudera documentation describes it this way: only fields referenced by predicates, and rows not filtered out by predicates, are fully materialized. The option parquet_late_materialization_threshold sets the minimum run of consecutive filtered-out rows for which materialization is skipped; the default is 20, and a negative value disables the feature. It works only for Parquet.

The gain is largest for wide tables with a selective filter: reading twenty columns where one predicate keeps one row in a thousand, most values of the other nineteen are never decoded.

Runtime filters are the other row-level mechanism. When a join's build side is small, Impala builds bloom or min and max filters from its keys and pushes them to the probe-side scan, where they can drop rows, and for Parquet whole row groups and pages, before the join sees them. How they are planned and why they sometimes arrive too late is covered in Impala Runtime Filters.

Skipping, as code

The row-group decision is simple enough to write down, and writing it down makes clear why sort order matters so much. This sketch mirrors the order of checks for a range predicate on one column:

def row_groups_to_read(row_groups, col, lo, hi, use_stats=True, use_dict=True):
    # Mirrors the order of checks: statistics first, then dictionary, else read.
    keep = []
    for rg in row_groups:
        s = rg.stats.get(col)
        if use_stats and s and (s.max < lo or s.min > hi):
            continue                               # whole row group proven empty for the range
        d = rg.dictionary.get(col)                 # only if every page of the chunk is dict-encoded
        if use_dict and d is not None and not any(lo <= v <= hi for v in d):
            continue                               # no dictionary value can match
        keep.append(rg)
    return keep

Two conditions in that sketch carry real restrictions. Dictionary filtering applies only when the whole column chunk is dictionary encoded; writers fall back to plain encoding when a dictionary grows too large, and from that point the chunk cannot be filtered by dictionary. And statistics must exist and be trustworthy: some older writers produced incorrect statistics for certain types and orderings, and readers ignore statistics they cannot rely on, so older files may be scanned in full.

Schema resolution: position or name

When the table schema in the metastore and the file schema differ, Impala must decide which file column feeds which table column. By default it resolves by position: the first table column reads the first file column. PARQUET_FALLBACK_SCHEMA_RESOLUTION can be set to name to match by column name instead.

Position resolution is fast and works for files Impala wrote into a table that only ever had columns added at the end. It produces wrong answers, not errors, when another engine writes files with a different column order or when a column is dropped from the middle of the table: a STRING column can silently read another STRING column. Name resolution survives reordering and dropped columns but treats renamed columns as missing, returning NULL. For tables written by several engines, prefer name resolution set consistently for the table's users, test it on a copy, and avoid renames; table formats such as Iceberg solve this properly with column IDs, as discussed in Impala + Iceberg Tables.

Timestamps are the other cross-engine trap. Hive and Spark historically wrote INT96 timestamps adjusted to UTC while Impala did not, so one file could show different times in different engines. Impala has startup flags for legacy Hive timestamps; agree one convention for all writers.

Writing files the reader can skip

Impala as a writer aims for large files with a single row group: its documentation notes that Impala-written Parquet files typically contain one row group and are sized to match the block size, so a file is one HDFS block. It recommends inserting roughly 256 MB per INSERT, or a multiple of it, per partition. Many small inserts produce many small files, each with its own footer to read and its own row group to filter.

Sorting is the single most effective writer setting. A SORT BY clause on the table makes inserts sort rows by those columns before writing, so each file and each page covers a narrow range of the sort keys and statistics become selective. Choose the columns most often filtered with ranges or equality, high-selectivity first, and keep the list short because each extra column sorts within the previous ones. Compression is set per session with COMPRESSION_CODEC; snappy is the default, zstd typically gives smaller files for some extra CPU, and both decompress fast enough that I/O usually dominates.

Bloom filters are written only for columns named in the parquet.bloom.filter.columns table property, as a comma-separated list of column names with optional bitset sizes in bytes, and only according to PARQUET_BLOOM_FILTER_WRITE: NEVER, IF_NO_DICT to write them when the row group is not fully dictionary encoded, or ALWAYS. IF_NO_DICT avoids duplicating what the dictionary already proves. Supported types are integers, FLOAT, DOUBLE and STRING; DECIMAL, TIMESTAMP, DATE, CHAR and VARCHAR are not.

Worked example: one tenant, one hour

-- Table written by Impala, sorted so that statistics are selective.
CREATE TABLE events (
  event_time   TIMESTAMP,
  tenant_id    BIGINT,
  event_type   STRING,
  country      STRING,
  payload      STRING
)
PARTITIONED BY (event_date STRING)
SORT BY (tenant_id, event_time)
STORED AS PARQUET
TBLPROPERTIES ('parquet.bloom.filter.columns'='tenant_id');

SET COMPRESSION_CODEC=zstd;
SET PARQUET_BLOOM_FILTER_WRITE=IF_NO_DICT;
INSERT INTO events PARTITION (event_date)
SELECT event_time, tenant_id, event_type, country, payload, event_date
FROM staging_events;

COMPUTE STATS events;

-- The query: one tenant, one hour, a few columns out of many.
SELECT event_type, count(*)
FROM events
WHERE event_date = '2026-10-01'
  AND tenant_id = 4711
  AND event_time BETWEEN '2026-10-01 09:00:00' AND '2026-10-01 10:00:00'
GROUP BY event_type;

Suppose a day partition holds 300 GB written as about 1,200 files of 256 MB, and the query wants one tenant's events for one hour. Partition pruning keeps only 2026-10-01. Because the table is sorted by tenant then time, each file covers a narrow band of tenants, so statistics on tenant_id exclude almost every file; a handful remain whose ranges include 4711. Inside those, the page index on event_time skips pages outside the hour, and late materialization decodes event_type only for rows that pass both predicates. The scan reads tens of megabytes of column data instead of the 300 GB a careless layout would need.

Remove the SORT BY and the same query degrades sharply: every file's tenant range spans the whole domain, statistics prune nothing, and only the bloom filter on tenant_id can exclude files that truly lack the tenant. Verify the difference in the query profile: compare bytes read and rows read at the scan node with the table size, look at the filtered row-group counters, and use the techniques in Impala Query Plans, in depth to compare estimates with actuals. When results are unexpected, inspect the files themselves:

# Check what a file gives Impala to skip with, before blaming Impala for scanning it.
import pyarrow.parquet as pq

pf = pq.ParquetFile("part-00000.parq")
md = pf.metadata
print("created_by:", md.created_by, " row groups:", md.num_row_groups, " rows:", md.num_rows)
for i in range(md.num_row_groups):
    rg = md.row_group(i)
    print(f"rg {i}: rows={rg.num_rows} bytes={rg.total_byte_size}")
    for j in range(rg.num_columns):
        col = rg.column(j)
        st = col.statistics
        rng = (st.min, st.max) if st is not None and st.has_min_max else "no stats"
        print(f"   {col.path_in_schema:<12} {col.compression:<6} dict={col.has_dictionary_page} "
              f"column_index={col.has_column_index} range={rng}")

Failure modes and trade-offs

  • Many small files from streaming writers or frequent inserts: footer reads and per-file overhead dominate. Compact into large files with a periodic INSERT OVERWRITE or a table-format compaction.
  • Unsorted data: statistics exist but prune nothing. Add SORT BY and rewrite the hot partitions.
  • External writers with many small row groups or without page indexes: the reader works but skips less. Check created_by and the row-group layout before tuning Impala.
  • Position resolution on a table written by several engines: silent wrong answers. Decide name or position per table and test schema changes.
  • Dictionary fallback on high-cardinality columns: dictionary filtering stops working for those chunks. Use a bloom filter for equality lookups instead.
  • Timestamp conventions differing between engines: hour-shifted results. Standardise the writer convention.
  • Turning off reader options while debugging and leaving them off: each option exists to skip data, so keep defaults in production and change them only per session.

The trade-offs are mostly on the write side. Sorting and bloom filters make inserts slower and use memory; large files reduce parallelism for small tables; zstd saves storage and I/O at some CPU. Most clusters gain by paying those costs once at write time, since data is usually read many more times than it is written.

What to do next

  1. Pick your three most expensive Parquet queries and record bytes read and rows read at each scan node from their profiles.
  2. Inspect a few files from each table with the pyarrow script: row groups per file, statistics, dictionary encoding and page index presence.
  3. Add SORT BY on the columns those queries filter on, and rewrite the hottest partitions with an INSERT OVERWRITE.
  4. Keep insert volumes near 256 MB per partition, and compact tables that have accumulated small files.
  5. Add bloom filters with IF_NO_DICT only for high-cardinality equality columns, then check that profiles show fewer row groups read.
  6. Set schema resolution deliberately for tables shared with other engines, and agree a timestamp convention.
  7. Re-measure the same queries and keep the before and after profiles as evidence for the change.
Key takeaway: Impala reads Parquet with its own C++ scanner that skips data at five levels before decoding: partitions, row groups by statistics, dictionaries and bloom filters, then pages by the page index, followed by late materialization and runtime filters at row level. Every level depends on what the writer left in the files, so the main tuning work is on the write side: large files with one row group, SORT BY on filtered columns, page indexes on, and bloom filters for high-cardinality equality columns. Keep the reader options at their defaults, choose schema resolution deliberately for multi-engine tables, and prove every change with scan-node bytes and rows in the query profile.