Apache Impala 4 is not one release but a line of them: 4.0.0, 4.1.0, 4.2.0, 4.3.0, 4.4.0 and 4.5.0, with maintenance releases such as 4.1.2 and 4.5.2 on the project's download page. Read as a whole, the line changes what Impala is. It started 4.0 as a fast MPP engine over Hive tables with Sentry security and finished 4.5 as an engine that writes Iceberg tables with DELETE, UPDATE and MERGE, survives the loss of its catalog and statestore daemons, records every query in a SQL table, and authenticates with JWT and OAuth.
This article walks through the 4.x features by area rather than by release, says which release introduced each one with its Apache JIRA number so you can check it in the official changelogs, and turns the list into an upgrade plan. Distributions such as Cloudera's ship their own Impala builds with backports, so map features to your vendor's version as well as to Apache's. For the project's direction see the future of Impala.
The 4.x line at a glance
Impala keeps three kinds of daemon: impalad plans and executes queries, acting as coordinator, executor or both; catalogd loads table metadata from the Hive Metastore and storage and broadcasts it; statestored tracks membership and relays catalog and admission updates. Almost every 4.x feature belongs to one of five areas: execution, SQL, table formats, the control plane and the security and operations edge.
| Release | Headline changes (Apache JIRA) |
|---|---|
| 4.0 | MT_DOP in all operators (IMPALA-3902); experimental Iceberg (IMPALA-10149); full-ACID ORC reads (IMPALA-9042); ROLLUP, CUBE, GROUPING SETS; INTERSECT and EXCEPT; Ranger replaces Sentry; row filtering; S3 spilling; aarch64 |
| 4.1 | UTF-8 aware string functions (IMPALA-2019); ALTER TABLE and SET PARTITION SPEC for Iceberg; Kudu multi-row transactions; event-driven automatic invalidation (IMPALA-7954); JWT authentication |
| 4.2 | Ozone (IMPALA-9400); BINARY columns; Hive generic UDFs; Iceberg V2 position-delete reads and snapshot expiry; FILE__POSITION virtual column; Kerberos over HTTP in impala-shell |
| 4.3 | Iceberg DELETE (IMPALA-11877), rollback, LOAD DATA and table migration; codegen caching (IMPALA-11470); catalogd HA (IMPALA-12155); Java 17 support; query timeline in the web UI |
| 4.4 | Iceberg UPDATE (IMPALA-12313) and DROP PARTITION; statestore HA (IMPALA-12156); SQL interface to completed and live queries (IMPALA-12426, IMPALA-12540); SHOW VIEWS |
| 4.5 | Iceberg MERGE (IMPALA-12732); KILL QUERY (IMPALA-12648); OAuth (IMPALA-13288); event-processor control by SQL; ANSI trim(); RPM and DEB packages |
Execution: MT_DOP, codegen caching and spilling
Before 4.0, multi-threaded execution through the MT_DOP query option applied only to scans and aggregations; joins and other operators ran one fragment instance per node. Impala 4.0 extended it to all operators (IMPALA-3902). With MT_DOP=4, each executor runs four instances of every fragment, so a join-heavy query can use more than one core per node.
-- Compare one heavy query at two degrees of parallelism.
SET MT_DOP=1;
SELECT c.region, sum(o.amount) FROM orders o JOIN customers c ON o.cust_id = c.id GROUP BY c.region;
SUMMARY; -- per-operator time and rows, in impala-shell
SET MT_DOP=4;
SELECT c.region, sum(o.amount) FROM orders o JOIN customers c ON o.cust_id = c.id GROUP BY c.region;
SUMMARY;More instances mean more memory: each instance has its own hash tables and buffers, so peak per-node memory can rise roughly with MT_DOP, which interacts with admission control. Raise it per pool or per query after measuring, not globally. In the two SUMMARY outputs, compare the join and aggregation rows: if their average time per instance falls while the scan rows stay flat, the query was CPU-bound in those operators and MT_DOP helps; if the scans dominate either way, the bottleneck is I/O and extra instances only add memory. Record peak memory per node from the profile for both runs before choosing a pool default. See Impala admission control for how memory estimates gate queries.
Code generation, Impala's LLVM-based compilation of expressions and operators, became cheaper in 4.3, which can cache generated functions across queries (IMPALA-11470) and added codegen for structs. Short repetitive dashboard queries benefit most, because codegen time was a large share of their runtime. Version 4.0 also allowed spilling to S3 (IMPALA-9867), useful for compute nodes with small local disks, at the price of slower spills.
SQL additions and one behaviour change
Impala 4.0 caught up on analytic SQL: ROLLUP, CUBE and GROUPING SETS (IMPALA-7204), INTERSECT and EXCEPT, and || as string concatenation when the left operand is a STRING. Later releases added UTF-8 aware string functions and masking functions (4.1), SHA2 and MD5 (4.1), BINARY columns (4.2), statistical functions (4.3), SHOW VIEWS (4.4) and an ANSI trim() (4.5). Complex types grew too: structs and arrays in the SELECT list for Parquet and ORC (4.1 and 4.2) and complex types nested inside each other (4.3).
-- One pass instead of three UNION ALL queries.
SELECT region, product, sum(amount) AS revenue
FROM sales
GROUP BY ROLLUP (region, product);
-- Rows in this month's load that were not in last month's.
SELECT cust_id FROM load_2026_10
EXCEPT
SELECT cust_id FROM load_2026_09;
-- 4.5: stop a runaway query from any coordinator, by id.
KILL QUERY '4f4c1e5a2b7d9e10:7a3b2c1d00000000';One behaviour change hides here: 4.0 disabled ordinals in HAVING by default (IMPALA-7844), so a number in a HAVING clause no longer refers to a select-list column. Search saved queries for it before upgrading.
Table formats: Iceberg, ORC, Kudu and storage
Iceberg is the biggest story of the line. Support arrived as experimental in 4.0 and grew release by release: table and partition-spec evolution in 4.1; V2 position-delete reads and snapshot expiry in 4.2; DELETE, rollback, LOAD DATA and in-place migration of external Hive tables in 4.3; UPDATE and DROP PARTITION in 4.4; and MERGE in 4.5. Impala writes row-level changes as position delete files, merge-on-read, so reads slow down until compaction. Syntax and maintenance are covered in Impala and Iceberg.
-- 4.5: upsert a day of changes into an Iceberg V2 table.
MERGE INTO sales t
USING sales_changes s ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET amount = s.amount, updated_at = s.updated_at
WHEN NOT MATCHED THEN
INSERT (order_id, region, product, amount, updated_at)
VALUES (s.order_id, s.region, s.product, s.amount, s.updated_at);
-- Bound the time-travel window so snapshot metadata does not grow forever.
ALTER TABLE sales EXECUTE expire_snapshots(now() - interval 7 days);ORC and Hive ACID. 4.0 reads full-ACID Hive 3 ORC tables (IMPALA-9042) with a vectorized, codegen-enabled ORC scanner; 4.1 pushes min-max, equality, IS NULL and IN-list predicates into ORC; 4.2 resolves ORC columns by name. Impala reads these tables but leaves writes and compaction to Hive.
Kudu gained multi-row transactions (4.1), non-unique primary keys (4.3), BINARY columns and soft-deleted tables (4.4). Storage gained Apache Ozone in 4.2 and Huawei OBS in 4.3, alongside HDFS and S3.
Control plane: event-driven metadata and HA
In 3.x, a Hive or Spark write outside Impala left stale metadata until someone ran REFRESH or INVALIDATE METADATA. 4.1 closed IMPALA-7954, automatic invalidation from Hive Metastore notification events: catalogd polls the event log and applies table and partition changes itself. 4.3 added event-processor errors to the catalogd web UI and 4.5 added SQL commands to control the event processor. Event lag is now a metric to watch, covered in Impala metadata in depth.
The other control-plane change is high availability. catalogd and statestored were single points of failure: losing catalogd stopped DDL and metadata updates, and losing statestored stopped membership and admission updates. 4.3 added an active-passive catalogd pair (IMPALA-12155) and 4.4 an active-passive statestored pair (IMPALA-12156). Operationally, treat each pair like any other failover service: place the two members on different hosts and racks, alert on which one is active, and rehearse a failover in staging by killing the active daemon during a DDL-heavy test. Watch what clients see while it happens; a short pause in DDL or in admission updates is expected, while failed queries or a split-brain are bugs to report before production depends on the pair.
Security and operations
Security. 4.0 removed Sentry and made Apache Ranger the only authorization provider, adding row-filtering policies alongside column masking, plus SAML authentication, Apache Knox integration and FIPS compliance. 4.1 added JWT authentication, 4.2 Kerberos over HTTP in impala-shell and deferred view authorization, and 4.5 OAuth. See Impala with Ranger and Kerberos.
Operations. 4.3 added a query timeline to the web UI. 4.4 added the biggest operational feature of the line: a SQL interface to completed and running queries (IMPALA-12426, IMPALA-12540). With --enable_workload_management set on the coordinators and catalogd, completed queries land in sys.impala_query_log, an Iceberg table, and running ones are visible in sys.impala_query_live. Questions that used to need profile scraping become SQL.
-- Slowest query shapes per pool over the last day.
SELECT resource_pool, db_user, count(*) AS runs,
avg(total_time_ms) AS avg_ms, max(pernode_peak_mem_max) AS peak_mem
FROM sys.impala_query_log
WHERE start_time_utc > now() - interval 1 day
GROUP BY resource_pool, db_user
ORDER BY avg_ms DESC
LIMIT 20;
-- The log is kept indefinitely; expire and compact it on a schedule.
ALTER TABLE sys.impala_query_log EXECUTE expire_snapshots(now() - interval 7 days);Platform changes to check: 4.0 requires AVX-capable CPUs (with a --enable_legacy_avx_support escape hatch) and supports ARM; 4.3 supports Java 17 and builds without Python 2; 4.5 installs from RPM and DEB packages.
Worked example: upgrading from 3.4 to 4.5
Worked example. A team runs Impala 3.4 on Hive 2 with Sentry, some LZO text tables, and nightly Spark jobs followed by a scripted INVALIDATE METADATA. They want 4.5 for Iceberg MERGE and the query log. The plan, in order:
- Preflight the hard breaks from 4.0. Hive 2 support is gone, so upgrade the metastore to Hive 3 first. LZO is removed, so list LZO tables and rewrite them as Parquet. Check every host with
grep -c avx /proc/cpuinfo. - Migrate authorization. Export Sentry roles and grants, translate them into Ranger policies, and test with real accounts before cutover; there is no fallback to Sentry.
- Fix query behaviour changes. Search saved SQL for HAVING ordinals and run the top hundred dashboard queries on a 4.5 staging cluster, comparing results and profiles.
- Turn on events, then retire scripts. Enable event processing, watch event lag for a week, and only then remove the scripted invalidations.
- Add HA and workload management. Deploy catalogd and statestored pairs, enable workload management, and schedule expiry and compaction of the query log.
- Adopt Iceberg gradually. Migrate one high-churn table, move its upserts to MERGE, and schedule compaction before extending.
Failure modes
| Failure | Cause | Fix |
|---|---|---|
| Daemons refuse to start after upgrade | CPU without AVX | Replace hosts, or use the legacy AVX flag as a stopgap |
| Users lose access at cutover | Sentry grants not fully translated to Ranger | Diff effective permissions per role before switching |
| HAVING queries return different rows | Ordinals now disabled | Rewrite with column names or aliases |
| Stale tables after Spark writes | Event processor paused, lagging or in error | Alert on event lag; check the catalogd UI |
| Memory admission rejections rise | MT_DOP raised globally | Set it per pool and re-measure peak memory |
| Iceberg reads slow down over weeks | Accumulated position delete files | Schedule compaction and snapshot expiry |
| Query log grows without bound | No expiry on sys.impala_query_log | Schedule expiry and compaction like any Iceberg table |
Per-query diagnosis is covered in Impala troubleshooting.
Trade-offs
Upgrading across the 4.x line buys Iceberg writes, HA control-plane daemons, self-updating metadata and SQL-queryable history. The costs are concentrated in 4.0: a forced Ranger migration, a Hive 3 metastore, CPU requirements and removed LZO. MT_DOP trades memory for latency. Iceberg row-level writes trade read speed for write convenience until compaction. Event-driven metadata trades manual refreshes for one more component to monitor. If you are on 3.x, the hard part is reaching 4.0; after that, each release is mostly additive, so go straight to the newest release your distribution supports.
What to do next
- Find your exact version with
SELECT version()and map it to the release table above, or to your vendor's release notes. - If you are on 3.x, inventory the 4.0 blockers: Hive 2 metastore, Sentry, LZO tables, non-AVX hosts and HAVING ordinals.
- Enable event-based metadata sync and alert on event lag before deleting any invalidation script.
- Deploy catalogd and statestored as HA pairs if you are on 4.3 or 4.4 and later respectively.
- Turn on workload management and build your first dashboard from
sys.impala_query_log, with expiry scheduled from day one. - Test MT_DOP on your five heaviest join queries and set it per pool where it helps.
- Pick one table with frequent corrections and pilot Iceberg with MERGE, with compaction scheduled.