Impala is a massively parallel SQL engine, but on a Hadoop cluster it is also a guest. It does not own the data, which lives in HDFS. It does not own the schema, which lives in the Hive Metastore. It shares every worker host with a DataNode and usually with a YARN NodeManager. Most Impala performance and correctness problems on Hadoop come from these boundaries rather than from the engine itself: scans that read remotely when they should read locally, queries that miss rows another engine just wrote, and hosts that run out of memory because Impala and YARN both assumed the RAM was theirs.
The daemons themselves, impalad, statestored and catalogd, and the roles they play are covered in the companion page on Impala's daemon architecture. This article covers Impala's architecture as seen from the cluster: where each piece runs, where metadata and data come from, how locality is achieved and measured, and how to share a host with the rest of the Hadoop stack.
The deployment shape
A typical on-premises layout puts one impalad on every host that runs a DataNode. That is the design decision everything else depends on: Impala reaches its speed by reading data from local disks rather than over the network, so the query process must run where the blocks are. The two singleton services run on master hosts: catalogd, which loads metadata from the Hive Metastore and the NameNode, and statestored, which broadcasts membership and catalog updates to every impalad. The Hive Metastore and the NameNode are not part of Impala, but Impala depends on both at all times.
Inside each impalad there are two runtimes. The frontend is Java, loaded into the process through JNI; it parses SQL, analyses it against the cached catalog and produces a distributed plan. The backend is C++; it schedules fragments, scans files, runs operators and exchanges rows with other impalads. Knowing which side a problem is on tells you whether to look at JVM heap or native memory.
Where metadata comes from, and why it goes stale
Impala needs two kinds of metadata. Logical metadata, meaning databases, tables, columns, partitions and storage formats, comes from the Hive Metastore, the same service Hive and Spark use. Physical metadata, meaning the files in each partition directory, their blocks and which DataNodes hold each replica, comes from the NameNode. catalogd loads both, caches them, and publishes changes through the statestore. In the newer on-demand mode, enabled with the use_local_catalog and catalog_topic_mode flags, coordinators fetch metadata from catalogd when a query needs it rather than receiving everything up front.
Caching is what makes planning fast, and it is also why metadata goes stale. When Impala itself runs INSERT, CREATE TABLE or ALTER TABLE, catalogd updates its cache immediately. When another engine changes the data, Impala does not know until it is told or until it learns from metastore events. The commands differ in cost:
-- Hive or Spark (through the metastore) added data to an existing partition
REFRESH sales PARTITION (sale_date='2026-09-30');
-- Files were copied into the table directory with hdfs dfs -put: no metastore event
REFRESH sales;
-- A table was created outside Impala, or its structure changed outside Impala
INVALIDATE METADATA analytics.new_events;
-- catalogd flag that enables automatic sync from metastore notification events
-- (documented upstream as off by default, value 0; check your release and distribution)
-- --hms_event_polling_interval_s=2REFRESH reloads the file and block list for a table Impala already knows, and it can be limited to a single partition. INVALIDATE METADATA discards the cached table so it is reloaded in full on next use, which is far more expensive for a table with thousands of partitions; use it for new tables and structural changes made elsewhere, never as a habit after every load.
Automatic sync reduces the manual work. With --hms_event_polling_interval_s set on catalogd, Impala reads the metastore's notification log and applies table, partition and insert events itself. The upstream documentation describes it as off by default and lists its limits: files added directly to the filesystem, and Spark writes that bypass the metastore, generate no event, so they still need REFRESH.
Life of a scan on HDFS
Follow one query to see how the pieces combine. A client connects to an impalad acting as coordinator, typically over the HiveServer2 protocol on port 21050. The frontend parses and analyses the statement, prunes partitions from the predicate on sale_date, and turns the surviving files into scan ranges, slices of files that can be read independently. For each scan range the catalog already knows which DataNodes hold a replica.
The scheduler then assigns each scan range to an executor. The rule that matters on Hadoop is locality: a range is assigned to an impalad on a host that stores a replica, so the read stays on local disk. The REPLICA_PREFERENCE query option controls the preference, with values CACHE_LOCAL (the default), DISK_LOCAL and REMOTE. Fragments are sent to executors, the C++ scanners read their ranges, and partial results move between impalads only at exchange points such as a shuffle for a join or the final merge of an aggregation. The coordinator streams the result back to the client.
-- Prefer replicas in the HDFS cache, then on local disk (the default is CACHE_LOCAL)
SET REPLICA_PREFERENCE=CACHE_LOCAL;
-- Spread scans of a hot, highly replicated table across all hosts holding a replica
SET SCHEDULE_RANDOM_REPLICA=TRUE;
SELECT COUNT(*) FROM dim_customer WHERE segment = 'enterprise';When many queries hit the same small, highly replicated table, choosing among equally local replicas without randomisation can concentrate the work on a few hosts. SCHEDULE_RANDOM_REPLICA, available since Impala 2.5, spreads such scans across all hosts that hold a replica.
Short-circuit reads: the local read that skips the DataNode
A local read through the normal HDFS protocol still streams bytes from the DataNode process over a socket. A short-circuit read lets the client open the block file directly after the DataNode passes it a file descriptor over a Unix domain socket, removing a copy and a context switch per buffer. Impala relies on this for its scan throughput. Two settings enable it, and they must be visible to both the DataNodes and every impalad, because impalad is the HDFS client:
<!-- hdfs-site.xml, visible to the DataNodes AND to every impalad -->
<property>
<name>dfs.client.read.shortcircuit</name>
<value>true</value>
</property>
<property>
<name>dfs.domain.socket.path</name>
<value>/var/run/hdfs-sockets/dn</value>
</property>The socket directory must exist on every host with the permissions the DataNode expects. If the configuration is wrong, HDFS silently falls back to ordinary reads. Queries still work; they are just slower, which is why a misconfiguration can go unnoticed for months. See the short-circuit read guide for the HDFS side.
Worked example: diagnosing a scan that got slower
A daily report that scanned the last 30 days of sales in 40 seconds now takes three minutes. Nothing in the SQL changed. The plan is the same, so the question is where the bytes came from. Run the query and print its profile:
-- in impala-shell, run the query and then print its profile
SELECT region, SUM(amount)
FROM sales
WHERE sale_date >= '2026-09-01'
GROUP BY region;
PROFILE;
-- Illustrative excerpt from the HDFS_SCAN_NODE section (values invented for the example):
-- BytesRead: 12.40 GB
-- BytesReadLocal: 12.40 GB
-- BytesReadShortCircuit: 12.40 GB
-- BytesReadDataNodeCache: 0
-- BytesReadRemoteUnexpected: 0In a healthy run every byte is local and short-circuit, as in the excerpt. Read the counters as a decision tree. If BytesReadLocal is well below BytesRead, ranges are being read remotely: either some DataNodes have no impalad, or the cached block locations are wrong. If BytesReadLocal is high but BytesReadShortCircuit is low, the reads are local but go through the DataNode, which points at the short-circuit settings or the socket path. A non-zero BytesReadRemoteUnexpected means the scheduler planned a local read that turned out to be remote.
In this example the investigation found that three new worker hosts had been added with DataNodes but no impalad, and the HDFS balancer had moved a large share of recent blocks onto them. Every block that now lived only on those hosts was read across the network. Two fixes applied: deploy impalad to the new hosts, and run REFRESH sales so catalogd reloaded the moved block locations. Note that some of these counters are only finalised when a query finishes successfully, so compare completed runs.
HDFS caching for the hottest data
HDFS centralised caching pins chosen blocks in DataNode memory. Impala can request it per table or per partition, and a cached replica becomes the most preferred under the default CACHE_LOCAL setting. With a cache replication greater than one, Impala picks randomly among the hosts holding a cached copy, which also avoids a CPU hotspot:
# HDFS side: create a cache pool that Impala may use (limit in bytes)
hdfs cacheadmin -addPool impala_hot -owner impala -limit 64424509440
-- Impala side: pin a dimension table and the newest partition of a fact table
ALTER TABLE dim_customer SET CACHED IN 'impala_hot' WITH REPLICATION = 3;
ALTER TABLE sales PARTITION (sale_date='2026-09-30')
SET CACHED IN 'impala_hot' WITH REPLICATION = 2;
-- check what is cached and how much
-- hdfs cacheadmin -listDirectives -stats
-- release it when the partition cools down
ALTER TABLE sales PARTITION (sale_date='2026-09-30') SET UNCACHED;Cache small dimension tables joined by many queries and the most recent partitions of fact tables. Do not cache large historical data: the memory comes from the same hosts that run Impala and YARN, and every gigabyte pinned is a gigabyte neither can use. See HDFS centralised caching for the pool and directive model.
Memory: sharing a host with YARN
Impala does not run inside YARN. impalad is a long-running daemon that manages its own memory, bounded by the --mem_limit startup flag, and admission control decides whether a query may start based on the memory it is expected to need. YARN's NodeManager separately advertises yarn.nodemanager.resource.memory-mb to the ResourceManager. Neither knows about the other, so the split must be planned explicitly.
A worked budget for a 256 GB worker: reserve about 16 GB for the operating system and page cache headroom, about 8 GB for the DataNode and NodeManager processes, and the DataNode's cache memory, dfs.datanode.max.locked.memory, if you use HDFS caching, say 8 GB; a cache pool's limit, by contrast, is a cluster-wide total across all hosts. That leaves roughly 224 GB to divide. Giving Impala a mem_limit of 120 GB and YARN 100 GB keeps the total under the physical memory. If both are set to 200 GB because each was sized alone, the host is overcommitted, and when a Spark job and a large Impala join peak together the kernel's out-of-memory killer picks a victim. Pair the static split with admission control pools so Impala queues queries instead of exceeding its share.
Failure modes
- Missing rows after another engine writes. Stale file lists in catalogd. Refresh the affected partitions after loads, or enable event polling and still refresh after direct filesystem writes.
- File-not-found errors during a scan. A compaction or
INSERT OVERWRITEelsewhere deleted files Impala still lists. Refresh, and schedule rewrites away from peak query windows. - Remote reads after rebalancing or decommissioning. Blocks moved; the cached locations did not. Refresh large tables after balancer runs, and keep impalad on every DataNode.
- Silent loss of short-circuit reads. A configuration change or a missing socket directory. Alert when
BytesReadShortCircuitfalls well belowBytesReadLocalon routine queries. - Catalog pressure from small files. Every file and block is metadata held in catalogd's JVM heap and loaded through NameNode calls. Millions of small files slow loads and refreshes; compact them at write time.
- Overcommitted hosts. Impala and YARN sized independently. Budget memory per host, as above.
- Statestore or catalogd outage. Running impalads keep serving queries with the state they last received, but metadata changes and membership updates stop flowing until the service returns. Monitor both and restart quickly.
Trade-offs: co-located HDFS versus object storage
| Aspect | Impala on co-located HDFS | Impala on object storage |
|---|---|---|
| Scan locality | Local and short-circuit reads | Always remote; mitigated by a local data cache |
| Scaling | Compute and storage scale together | Compute scales independently of data |
| Metadata cost | Block locations from the NameNode | Object listings; no block locations |
| Operations | Run HDFS, balance disks, plan memory beside YARN | No DataNodes to run; pay per request and egress |
Co-location is the fastest design for steady, scan-heavy workloads on owned hardware. Object storage wins when compute demand varies or storage grows faster than compute; Impala's local data cache recovers much of the locality benefit there. The metastore stays shared in both cases; see the Hive Metastore for its own scaling and availability.
What to do next
- List every host running a DataNode and confirm each also runs an impalad; fix any gaps first.
- Verify short-circuit reads on every host: the two
hdfs-site.xmlsettings, the socket directory, and the profile counters on a routine query. - Add a refresh step to every pipeline that writes Impala tables through Hive, Spark or direct file copies, scoped to the partitions written.
- Decide whether to enable
--hms_event_polling_interval_s; if you do, keep refreshes for writes that bypass the metastore. - Write down a per-host memory budget covering the OS, DataNode, NodeManager, HDFS cache, Impala
mem_limitand YARN, and make the numbers add up. - Cache only small, hot tables and recent partitions, and review cache directives monthly.
- Refresh large tables after balancer or decommission runs, and track remote bytes in your query monitoring.