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.
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.
| Property | HDFS | S3 through S3A |
|---|---|---|
| Where bytes live | Blocks replicated on DataNodes, often the same hosts as impalad | Remote service; no host is local to any object |
| Read path | Local or short-circuit read straight from disk | HTTPS GET with a byte range per request |
| Rename | Metadata operation on the NameNode, atomic for a file | No rename: copy then delete, cost grows with bytes |
| Listing a directory | One NameNode call, fast | Paged LIST requests, slower and billed |
| Consistency | Strong | Strong read-after-write since December 2020 |
| Caching options | OS page cache, HDFS centralized cache | Impala data cache on executor-local disks |
| Cost model | Pay for disks and nodes all the time | Pay 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.
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.
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 ... PURGEfor 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
| Symptom | Likely cause | Fix |
|---|---|---|
| HDFS scans slow, high remote bytes | Executors not on DataNode hosts, or short-circuit reads disabled | Co-locate, enable dfs.client.read.shortcircuit and a domain socket path |
| S3 scans slow after scaling executors | Hash ring changed, data caches cold | Scale in steps, warm caches with the common queries |
| Rows from a failed insert appear | S3_SKIP_INSERT_STAGING writes final files directly | Rebuild the partition; write partitions from one idempotent job |
| Queries stall on table load | Full INVALIDATE METADATA on a large S3 table | Partition-level REFRESH; compact small files |
| S3 throttling errors (503 Slow Down) | Too many requests to one prefix | Fewer, larger files; spread partitions across prefixes |
| SET CACHED fails on cold partitions | HDFS caching does not apply to S3 | Use 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
- Pull five recent query profiles and record local, short-circuit, remote and cached bytes for each scan.
- 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.
- Replace every table-wide INVALIDATE METADATA in pipelines with partition-level REFRESH or ADD PARTITION.
- Decide whether partial writes from skipped staging are acceptable; if not, plan a move to Iceberg.
- Enable the data cache on executors with a quota sized from the daily working set.
- 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.