Impala is fast partly because it does not ask the metastore or the file system anything while it plans a query. Schemas, partitions, file lists, block locations and statistics are cached in advance by the catalog daemon and pushed to every coordinator. That cache is also the source of the most common Impala support question: why can Hive see my new data when Impala cannot?
This article explains the metadata lifecycle end to end: what exactly is cached, how it travels, how it goes stale, what REFRESH and INVALIDATE METADATA cost, what event-based automatic sync covers and what it misses, and how to keep catalog memory under control. The daemon architecture itself is covered in the Impala catalog service and the statestore; this page is about operating the metadata that flows through them.
Why metadata is the hard part
To plan a scan, a coordinator needs much more than a table schema. It needs the list of partitions to prune, every file in the surviving partitions with its size so it can build scan ranges, block locations on HDFS so it can schedule reads near the data, and statistics so it can order joins and estimate memory for admission control.
Consider a table partitioned by day and hour for two years: about 17,500 partitions. If each partition holds 40 files, the table has 700,000 files. Listing those at plan time would mean hundreds of thousands of metastore and storage calls before a single row is read, and on object stores listings are slow and throttled. Caching turns that into a memory lookup. The price is that the cache must be kept correct, and must fit in memory.
What exactly is cached
The catalog caches several distinct kinds of object, and they go stale for different reasons. Knowing which kind changed tells you which command fixes it.
| Object | Source | Used for | Refreshed by |
|---|---|---|---|
| Database and table schema | Hive Metastore | name resolution, column types, SerDe and format | ALTER TABLE event, INVALIDATE METADATA |
| Partition list and locations | Hive Metastore | partition pruning | partition events, REFRESH, ALTER TABLE ... RECOVER PARTITIONS |
| File descriptors (name, size, modification time) | file system listing | scan ranges, split assignment | REFRESH, INSERT event |
| Block locations (HDFS) | NameNode | local reads and scheduling | REFRESH |
| Table and column statistics | Hive Metastore | join order, memory estimates, admission control | COMPUTE STATS, ALTER TABLE event |
| Functions | Hive Metastore and catalog | UDF resolution | CREATE and DROP FUNCTION, REFRESH FUNCTIONS |
Tables load lazily. After a catalog restart, a table exists only as a name until the first query or DDL touches it, and that first query pays the full load. For very large tables this is the source of the familiar complaint that the first query of the morning takes a minute and every later one takes a second.
How changes travel to coordinators
Only catalogd mutates cached metadata. Each change gets a new catalog version, and coordinators learn about changes in one of two modes. In the legacy mode, catalogd publishes full objects on the statestore catalog topic and every coordinator keeps a complete copy of every loaded table. In on-demand mode, catalogd publishes only invalidations and coordinators fetch what they need at partition granularity and cache it with eviction.
On-demand mode is enabled with --catalog_topic_mode=minimal on catalogd and --use_local_catalog=true on every coordinator. Its coordinator cache is sized by local_catalog_cache_mb, where the default of -1 means 60 percent of the JVM heap, and entries expire after local_catalog_cache_expiration_s, default 3600 seconds. Cloudera's documentation notes that HDFS caching is not supported in this mode. The architectural trade-off between the two modes is argued in the catalog article; for operations, the important point is that on-demand mode shrinks coordinator memory and avoids broadcasting a huge table after a small partition change.
How metadata goes stale
When Impala itself runs DDL or DML, catalogd updates its cache as part of the statement, and the change is visible to that coordinator immediately. Staleness comes from everything else:
- Hive, Spark or another engine adds partitions or inserts data through the metastore.
- An ingest job writes files straight into a table or partition directory without telling the metastore.
- Someone changes a schema, a location or table properties from Hive.
- Statistics are recomputed by another engine.
In the first and third cases the metastore knows about the change, so an event-based mechanism can find it. In the second case only the storage layer knows, and nothing short of a REFRESH will reveal the files.
REFRESH versus INVALIDATE METADATA
The two commands are often used interchangeably, which is how clusters end up with multi-minute catalog stalls. They do very different amounts of work.
| Command | What it does | Cost | Use when |
|---|---|---|---|
| REFRESH t | reloads t's partitions and file metadata, reusing unchanged entries where possible | proportional to files in t | another engine added or rewrote data in existing partitions |
| REFRESH t PARTITION (k=v) | reloads one partition's files | small | a job wrote exactly one known partition |
| ALTER TABLE t RECOVER PARTITIONS | adds partition directories that exist on storage but not in the metastore | listing of the table root | directories were created without metastore calls |
| INVALIDATE METADATA t | discards t entirely; full reload on next access | full reload, paid by the next query | schema changed outside Impala, or the cached object is corrupt |
| INVALIDATE METADATA | discards every table in the catalog | every table reloads lazily; very high | almost never in production |
The rule of thumb: REFRESH when data changed, INVALIDATE when the table's definition changed outside Impala, and never run a global INVALIDATE METADATA as a routine step. After an invalidation, the next query against the table blocks while catalogd reloads it, and on a large table that is the slow query users will report.
-- A Spark job appended files to one existing partition
REFRESH sales.orders PARTITION (dt='2026-09-29');
-- A job created new partition directories without metastore calls
ALTER TABLE sales.orders RECOVER PARTITIONS;
-- Hive added a column to the table
INVALIDATE METADATA sales.orders;
-- Check what Impala now sees
SHOW PARTITIONS sales.orders;
SHOW FILES IN sales.orders PARTITION (dt='2026-09-29');
Event-based automatic sync
Instead of relying on every pipeline to issue REFRESH, catalogd can poll the Hive Metastore's notification event log and apply changes itself. The switch is the catalogd flag --hms_event_polling_interval_s; a positive value enables polling at that interval, and the upstream documentation recommends a value under 5 seconds. The upstream page lists the default as 0, meaning off, and some distributions enable it by default, so check the value your build actually runs with. The metastore must also be configured to record notification events.
According to the Impala documentation, the event processor handles these operations: ALTER TABLE invalidates the table; adding, altering or dropping partitions refreshes those partitions; CREATE and DROP of tables and databases add and remove entries; INSERT refreshes the affected table and partitions if the table is loaded; and ALTER DATABASE updates its properties. It does not see files added to or removed from storage directly, and it does not see Spark writes that save to a path with write.save() rather than through the metastore.
Sync can be turned off per database or table with the property impala.disableHmsSync, which is useful for staging tables that churn constantly and that no one queries from Impala.
-- Opt a noisy staging table out of event processing
ALTER TABLE staging.raw_clicks
SET TBLPROPERTIES ('impala.disableHmsSync'='true');Event sync is asynchronous. There is always a lag of at least the polling interval, plus the time to apply the event, and a burst of events from a large backfill can build a backlog. Pipelines that must guarantee visibility before a downstream step should still issue an explicit REFRESH as their last action.
Consistency across coordinators: SYNC_DDL
A DDL or DML statement is visible on the coordinator that ran it as soon as it returns, but other coordinators receive the change only when the next catalog update reaches them. A load balancer that sends the next statement to a different coordinator can therefore produce a table not found error just after a successful CREATE.
Setting the query option SYNC_DDL=1 makes the statement wait until all coordinators have received the new catalog version. It adds latency to each DDL, so enable it in sessions that create objects and then query them through a load balancer, such as ETL scripts, rather than globally.
SET SYNC_DDL=1;
CREATE TABLE mart.daily_revenue STORED AS PARQUET AS
SELECT dt, SUM(amount) AS revenue FROM sales.orders GROUP BY dt;
-- the next statement can land on any coordinator and still see the table
Worked incident: the missing partition
A dashboard showed no sales for yesterday. Hive returned the rows, Impala did not. SHOW PARTITIONS in Impala listed the partition, so the metastore change had arrived. SHOW FILES for that partition returned nothing. The Spark job had been changed to write Parquet files directly to the partition path with df.write.mode("append").parquet(path) after the partition was created empty by an earlier step. The partition creation reached Impala through events; the files, written straight to storage, never generated an event.
The immediate fix was REFRESH sales.orders PARTITION (dt='2026-09-29'). The durable fix was to make the pipeline publish through the metastore, using insertInto on the table, and to add a final REFRESH step for the partitions it wrote, as below.
from impala.dbapi import connect
def refresh_partitions(table, partitions, host="impala-lb.internal", port=21050):
"""Make freshly written partitions visible to every coordinator."""
conn = connect(host=host, port=port)
cur = conn.cursor()
cur.execute("SET SYNC_DDL=1")
for spec in partitions: # e.g. "dt='2026-09-29'"
cur.execute(f"REFRESH {table} PARTITION ({spec})")
cur.close()
conn.close()
Catalog memory and large tables
Every cached file descriptor, partition and statistics entry lives in the catalogd JVM heap, and in legacy mode also in every coordinator's heap. Tables with millions of small files are the usual cause of catalog memory trouble, long full reloads and oversized catalog topic updates. The file count is the number to watch, more than the data volume.
Reduce it at the source: compact small files, avoid over-partitioning on high-cardinality keys, and keep incremental statistics in check, because per-partition incremental stats add to metadata size, as explained in Impala table statistics. Size catalogd's heap from measured usage, not guesswork, and alert on garbage collection time. On Iceberg tables, file lists come from table metadata rather than directory listings, which changes the refresh story but not the need to keep file counts sane.
Failure modes
- Invisible data. Files written directly to storage with no REFRESH. Check with SHOW FILES before assuming a query bug.
- First query after restart is slow. Lazy loading of large tables. Warm critical tables with a cheap query after restarts, and schedule restarts outside business hours.
- Catalog stall after a global INVALIDATE METADATA. Every table reloads on first touch. Remove the command from scripts and replace it with targeted REFRESH.
- Event backlog. A big backfill in Hive generates thousands of partition events and Impala lags. Watch the event processor metrics your version exposes, and opt noisy staging tables out.
- Table not found just after CREATE. Statement routed to another coordinator; use SYNC_DDL in that session.
- Wrong plans after data changes. Statistics stale or missing, not metadata. Recompute statistics after large loads.
- Out of memory in catalogd. Too many files or partitions. Compact, prune partitions and consider on-demand mode.
What to do next
- Find out whether event polling is on: read the
hms_event_polling_interval_svalue your catalogd runs with. - List every pipeline that writes Impala-queried tables without going through the metastore, and add a targeted REFRESH at the end of each.
- Remove global INVALIDATE METADATA from scripts and runbooks.
- Enable SYNC_DDL in ETL sessions that create and immediately query objects through a load balancer.
- Rank tables by file and partition count, and compact the worst offenders.
- If coordinators run short of heap, test on-demand metadata mode on one coordinator group.
- Add SHOW PARTITIONS and SHOW FILES to your missing-data runbook, before anything else.
- Review the Hive Metastore setup, since its health and notification log bound Impala's freshness.