Impala was built for HDFS. Its scheduler assumed that every block of every file had a few replicas sitting on the same machines as the query engine, and that reading a block meant reading a local disk. Then storage moved: object stores became cheaper per terabyte, compute was separated from storage, and the same SQL started running against files reached over HTTP. Impala still runs the same plan on S3, but almost every assumption underneath that plan changes.

This article explains what changes and what to do about it. It compares the two storage contracts, follows a scan from planner to bytes on each, configures S3A for Impala, explains why writes and metadata cost more on an object store, and shows how to run one table whose recent partitions stay on HDFS while older ones age out to S3. The remote-read cache is covered in depth in Impala data cache and the catalog in Impala metadata in depth; this page cites them where they matter.

Advertisement

Two storage contracts behind one table

An Impala table is a directory of files plus metadata in the Hive Metastore. The planner does not care whether the directory starts with hdfs:// or s3a://. The storage layer underneath it cares a great deal, because HDFS and S3 make very different promises.

PropertyHDFSS3 through S3A
Where bytes liveBlocks replicated on DataNodes, often the same hosts as impaladRemote service; no host is local to any object
Read pathLocal or short-circuit read straight from diskHTTPS GET with a byte range per request
RenameMetadata operation on the NameNode, atomic for a fileNo rename: copy then delete, cost grows with bytes
Listing a directoryOne NameNode call, fastPaged LIST requests, slower and billed
ConsistencyStrongStrong read-after-write since December 2020
Caching optionsOS page cache, HDFS centralized cacheImpala data cache on executor-local disks
Cost modelPay for disks and nodes all the timePay per GB stored plus per request

The Impala S3 documentation still recommends S3Guard and warns about eventual consistency. That advice predates S3's move to strong read-after-write consistency, and S3Guard has since been removed from Hadoop's S3A connector. Ignore both on current versions; the real costs on S3 today are latency, the lack of rename and the price of listing.

How a scan reaches the bytes

When a query arrives, the coordinator asks its catalog cache for the table's partitions, files and, on HDFS, the block locations of each file. The planner splits each file into scan ranges and the scheduler assigns those ranges to executors. On HDFS the scheduler prefers an executor on a host that holds a replica of the block. That executor then reads the block through a short-circuit read: the HDFS client gets a file descriptor from the DataNode over a Unix domain socket and reads the block file directly, skipping the DataNode's network path.

On S3 there are no block locations, so there is nothing local to prefer. A naive scheduler would spread ranges at random, and every query would read every byte from S3 again. Impala instead hashes each file to a small set of candidate executors with consistent hashing, so the same file tends to land on the same executor query after query. That stability is what makes the executor-local data cache useful: the second query that touches a file finds its bytes on local SSD.

One query, two storage paths: the planner is the same, the bytes travel differentlyCoordinatorplan + scheduleCatalog cachefiles, blocks, statsmetadataExecutor AHDFS DataNode hostExecutor Bcompute-only hostlocal rangeshashed remote rangesLocal disk replicashort-circuit readno networkData cachelocal SSD, LRU/LIRShit?S3 bucket via s3a://GET with byte rangemissPartition locationshot on HDFS, cold on S3per-partition LOCATIONHDFS ranges are scheduled to hosts holding a replica; S3 ranges have no host, so files are hashed to executorsConsistent hashing keeps a file on the same executors, which is what makes the data cache hit on the next query
HDFS ranges go to executors that hold a replica and read it without the network. S3 ranges are hashed to executors, checked against the local data cache, and fetched with ranged GETs on a miss.

Two profile counters tell you which path a scan took. BytesReadLocal and BytesReadShortCircuit show HDFS locality; on S3 they stay at zero by design, and the bytes read from the data cache are what you watch instead. If an HDFS table shows high remote bytes, executors and DataNodes are not co-located or the scheduler could not place ranges locally, and every scan is paying for the network it was designed to avoid.

Advertisement

Configuring S3A for Impala

Impala supports only the s3a:// scheme, not s3:// or s3n://. Configuration lives in core-site.xml for every impalad and catalogd. The documented settings that matter most are the connection pool and the split size. Credentials should come from an instance profile or another credential provider, never from keys in table URLs, which would end up in logs and query profiles.

<!-- core-site.xml on every impalad and catalogd -->
<property>
  <name>fs.s3a.connection.maximum</name>
  <value>1500</value>   <!-- value recommended in the Impala S3 docs -->
</property>
<property>
  <name>fs.s3a.block.size</name>
  <value>268435456</value> <!-- 256 MB: matches Parquet files written by Impala -->
</property>
<!-- No fs.s3a.access.key / secret.key here: use the host's IAM role -->

S3 has no blocks, so the S3A connector reports a fake block size and the planner splits files at those boundaries. For Parquet, Impala 3.4 and later use the PARQUET_OBJECT_STORE_SPLIT_SIZE query option instead, with a 256 MB default. The aim is that one scan range covers one Parquet row group: if splits cut row groups in half, two executors each fetch the footer and part of a group, and you pay for extra requests and wasted reads. Impala writes 256 MB Parquet files by default; files written by Hive or Spark often have 128 MB row groups, which is why the docs suggest 128 MB for those tables.

Writes without rename

On HDFS, an INSERT writes its files into a hidden staging directory and the coordinator moves them into the table directory when the query finishes. A move is a cheap, atomic NameNode operation, so readers never see half a query's output. On S3 a move is a copy of every byte followed by a delete, which can take longer than the query that produced the data.

Impala therefore skips staging on S3 by default: the S3_SKIP_INSERT_STAGING query option is true, and executors write final files directly into the table location. The trade is speed for atomicity. If a query fails halfway, some of its files may already be in the table, and readers that list the directory after a REFRESH can see them. INSERT OVERWRITE and LOAD DATA still have to copy or delete existing objects, so they are slow on S3 in proportion to the data they touch.

  • Write each partition once, from one job, and treat a failed write as a partition to rebuild rather than a query to retry in place.
  • Prefer appending new partitions to overwriting old ones; overwrite on S3 deletes and rewrites objects.
  • Use DROP TABLE ... PURGE for managed tables on S3, which the docs note is much faster than the default drop.
  • If you need atomic commits on object storage, move the table to Iceberg, where a commit is a metadata pointer swap and not a directory move.

Metadata is the hidden cost

Impala caches file listings in the catalog, so a query does not list directories on every run. The cost appears when that cache has to be rebuilt. On HDFS, listing ten thousand partitions is a burst of NameNode calls measured in seconds. On S3 the same listing is thousands of paged LIST requests, each a round trip, and a full INVALIDATE METADATA on a large partitioned table can take minutes and stall queries that wait for the table to load.

The rule is to refresh as narrowly as possible. When an external job adds data to one partition, refresh that partition instead of the table.

-- After a Spark job writes dt=2026-09-30 to S3:
REFRESH sales.events PARTITION (dt='2026-09-30');

-- After it adds a brand-new partition directory:
ALTER TABLE sales.events ADD IF NOT EXISTS PARTITION (dt='2026-10-01');
REFRESH sales.events PARTITION (dt='2026-10-01');

-- Avoid in pipelines; reloads every file listing of the table:
-- INVALIDATE METADATA sales.events;

Small files multiply the problem: each file is an entry in the catalog, a LIST result, a footer GET and a scan range. A table with a million 2 MB files on S3 is slow to load, slow to plan and slow to scan. Compact to files in the hundreds of megabytes before data reaches the long-term location.

The data cache on S3 executors

The data cache stores byte ranges read from remote storage on executor-local disks. It is enabled per daemon with a flag that lists directories and a quota per directory, for example --data_cache=/data/0,/data/1:500GB, which allows up to 500 GB in each directory. Impala 3.4 and later can choose between LRU and the scan-resistant LIRS policy with --data_cache_eviction_policy. The cache key is the file name, its modification time and the offset, so rewriting a file under the same name invalidates its entries naturally.

Size the cache from the working set, not the table: the partitions your dashboards hit every day, divided across executors. A cache smaller than the daily working set churns and behaves like no cache at all,. The data cache article covers sizing and its metrics; the point for this page is that on S3 the cache is what stands in for locality, and it only works when scheduling stays stable.

One table, two storage tiers

Impala lets each partition have its own LOCATION, and locations in one table can mix HDFS and S3. That makes a simple tiering design possible: keep the last month on HDFS, where locality and short-circuit reads make dashboards fast, and move older partitions to S3, where storage is cheap and nodes do not need to hold it. The query does not change; partition pruning picks the partitions and each scan uses the path its location implies.

# Age out one partition from HDFS to S3 (run from an edge node)
SRC=hdfs:///warehouse/sales.db/events/dt=2026-08-31
DST=s3a://corp-lake/sales/events/dt=2026-08-31

hadoop distcp -update "$SRC" "$DST"
hdfs dfs -count "$SRC"; hadoop fs -count "$DST"   # file count and bytes must match

impala-shell -q "ALTER TABLE sales.events PARTITION (dt='2026-08-31')
                 SET LOCATION '$DST';
                 REFRESH sales.events PARTITION (dt='2026-08-31');"

# Keep the HDFS copy for a grace period, then remove it
# hdfs dfs -rm -r -skipTrash "$SRC"

Two details matter. The copy must finish and be verified before the location changes, or queries will read an incomplete partition. The table should be external, or the move must account for Impala deleting data on a managed-table drop.

Worked example: a 30-day hot tier

A clickstream table holds two years of daily partitions, about 400 GB of Parquet a day, on a 20-node cluster where every node runs a DataNode and an impalad. HDFS with three replicas needs about 880 TB of raw disk for the two years, and capacity is the reason the team wants more nodes. Ninety-five percent of queries touch the last 30 days.

The team moves partitions older than 30 days to S3 with the procedure above, run nightly. HDFS now holds about 12 TB of data, 36 TB with replication, and the existing nodes have room to spare. Dashboard queries keep their local reads. Monthly reports that scan a year now read about 140 TB from S3; the first run of each report is slower, so the team gives each executor a 1 TB data cache and schedules the reports after the nightly move so the cache holds the newest cold partitions. These figures are an illustration of the arithmetic, not a benchmark; measure your own read latency with profiles before and after.

Failure modes

SymptomLikely causeFix
HDFS scans slow, high remote bytesExecutors not on DataNode hosts, or short-circuit reads disabledCo-locate, enable dfs.client.read.shortcircuit and a domain socket path
S3 scans slow after scaling executorsHash ring changed, data caches coldScale in steps, warm caches with the common queries
Rows from a failed insert appearS3_SKIP_INSERT_STAGING writes final files directlyRebuild the partition; write partitions from one idempotent job
Queries stall on table loadFull INVALIDATE METADATA on a large S3 tablePartition-level REFRESH; compact small files
S3 throttling errors (503 Slow Down)Too many requests to one prefixFewer, larger files; spread partitions across prefixes
SET CACHED fails on cold partitionsHDFS caching does not apply to S3Use the data cache for S3 partitions

For diagnosing which of these you are seeing, start with the query profile: scan node times, bytes read by source and the time spent in the catalog. The Impala troubleshooting guide walks through reading a profile end to end.

Trade-offs

HDFS gives you the fastest scans and atomic writes, but you pay for replicated disks all the time and you cannot scale compute without scaling storage. S3 separates the two, so you can run executors only when queries run, and it is far cheaper per terabyte, but every scan is remote, listings cost money and writes lose atomicity unless you adopt a table format such as Iceberg. A split design gets most of both at the price of a tiering job and two sets of failure modes. Whatever you choose, read fewer bytes first: partition pruning, column projection and runtime filters cut remote reads more than any cache.

What to do next

  1. Pull five recent query profiles and record local, short-circuit, remote and cached bytes for each scan.
  2. On S3 tables, check that Parquet row groups match the split size, and set PARQUET_OBJECT_STORE_SPLIT_SIZE or fs.s3a.block.size to match.
  3. Replace every table-wide INVALIDATE METADATA in pipelines with partition-level REFRESH or ADD PARTITION.
  4. Decide whether partial writes from skipped staging are acceptable; if not, plan a move to Iceberg.
  5. Enable the data cache on executors with a quota sized from the daily working set.
  6. If HDFS capacity is the constraint, pilot the per-partition move to S3 on one table, with copy verification and a grace period before deletion.
Key takeaway: On HDFS, Impala's speed comes from locality: scans are scheduled to replica hosts and read disks directly. On S3 there is no locality, so speed comes from reading fewer, larger files, splitting at row-group boundaries, stable hash-based scheduling with a local data cache, and narrow metadata refreshes. Writes lose rename and so lose atomicity. Per-partition locations let one table keep its hot data on HDFS and its history on S3 with no change to queries.