Upgrading Impala looks simple: it is a set of stateless daemons, so you replace binaries and restart. The difficulty is everything around the binaries. Impala depends on a specific Hive Metastore generation, on an authorization service, on the CPU instruction set of every host, and on SQL semantics that queries and dashboards have quietly relied on for years. A major version can change all of those at once, and the first sign of trouble is often a report whose numbers changed rather than an error.

This article explains what an Impala upgrade really involves, which breaking changes matter when moving from the 3.x line to 4.x, how to inventory a workload before you start, how to choose between upgrading in place and running a parallel cluster, the procedure itself, and how to validate and roll back. Facts about specific changes are taken from the Apache Impala incompatible-changes page and release notes; at the time of writing the newest release listed on the Apache downloads page is 4.5.2. Vendor distributions package Impala with their own versions and procedures, so check your vendor's notes for the exact build you run. The Hive side of the same project, including the metastore schema upgrade, is covered in the Hive upgrade path article.

What an upgrade actually touches

An Impala cluster has three daemon types. statestored tracks cluster membership and broadcasts updates. catalogd loads table metadata from the Hive Metastore and storage and distributes it. impalad runs on each worker, plans queries when acting as a coordinator, and executes query fragments. None of them stores table data: the data lives in HDFS, object storage or Kudu, and the table definitions live in the Hive Metastore. That is good news for upgrades, because there is no Impala data format to migrate, and it is the reason rollback is usually possible.

It also means the upgrade is really about compatibility at five boundaries: with the metastore, with the authorization service, with the hardware, with the SQL that users run, and with the clients that connect. The diagram shows them. A plan that only covers the binaries covers the easiest part.

What an Impala upgrade touchesImpala binariesimpalad, catalogd, statestoredHive Metastore4.x needs Hive 3 HMSAuthorization4.x: Ranger onlyHost CPUs4.x: AVX minimumData formatsLZO text removedSQL semantics||, HAVING ordinals, DECIMALClients and toolsimpala-shell, JDBC/ODBC, UDFsRestart order after replacing binaries: statestored, then catalogd, then every impalad
The binaries are the smallest part of the change. Each surrounding box is a compatibility boundary that a major version can move.

Breaking changes from 3.x to 4.x

Most clusters still on an old line are moving from 3.x to 4.x, often skipping several releases. Because incompatible changes accumulate, read every release's list between your current version and the target, not only the target's. The ones most likely to affect a production workload are these:

ReleaseChangeWhat breaks
3.0DECIMAL_V2 on by defaultDecimal precision, rounding and overflow behaviour differ; SET DECIMAL_V2=FALSE restores the old behaviour per session
3.0Fine-grained privileges; REFRESH needs its own privilegeUsers with only SELECT or INSERT can no longer refresh tables until granted
3.0Reserved word list updatedColumn or table names that became reserved need backticks
3.2:shutdown() targets the KRPC portDrain scripts that pass the old backend port fail
3.3Default table format becomes ParquetCREATE TABLE without STORED AS no longer creates text tables; set DEFAULT_FILE_FORMAT if a pipeline expects text
3.4PARQUET_OBJECT_STORE_SPLIT_SIZE (default 256 MB) used for object storesScan parallelism on S3 or ADLS changes for Parquet tables
4.0Hive 2.x support removedA Hive 2 metastore must be upgraded to Hive 3 first
4.0Sentry support removed; Ranger onlyAuthorization policies must exist in Ranger before cutover
4.0Minimum x86_64 CPU raised from SSSE3 to AVXDaemons will not run on older hosts unless started with --enable_legacy_avx_support
4.0Impala-lzo support removedLZO-compressed text tables become unreadable to Impala
4.0|| concatenates strings instead of meaning ORQueries that used || as a logical OR return different results
4.0Ordinals in HAVING disallowed by default; dateless timestamps droppedSome queries fail to parse; time-only timestamp values are rejected

Two of these change results without raising errors, which makes them the most dangerous: DECIMAL_V2 and the meaning of ||. The rest fail loudly, which is better, because you will find them in testing.

Inventory your workload first

Before choosing a date, collect evidence about your own workload. Four inventories are enough for most clusters.

Queries. Export the last 30 days of query text from the coordinators' query logs or your monitoring system. Then search it for constructs affected by the breaking changes:

import re
from collections import Counter

CHECKS = {
    "pipe_operator": re.compile(r"\|\|"),
    "having_ordinal": re.compile(r"\bHAVING\b[^;]*?\b\d+\s*(?:[<>=!]|$)", re.I),
    "create_without_stored_as": re.compile(r"\bCREATE\s+TABLE\b(?![^;]*\bSTORED\s+AS\b)", re.I | re.S),
    "decimal_math": re.compile(r"\bDECIMAL\s*\(", re.I),
    "refresh": re.compile(r"\bREFRESH\b", re.I),
}

def scan(queries):
    hits = Counter()
    examples = {}
    for q in queries:
        for name, rx in CHECKS.items():
            if rx.search(q):
                hits[name] += 1
                examples.setdefault(name, q[:200])
    return hits, examples

These patterns are deliberately rough; they find candidates for a human to review, not proof of breakage. A query that uses || on string columns is fine after 4.0; one that used it as boolean OR is not, and only reading it tells you which.

Tables. List table formats and compression, looking specifically for LZO text and any formats your target version handles differently. Hosts. On every worker run grep -c avx /proc/cpuinfo; a zero means the host lacks AVX. Security. Export every role and grant from your current authorization system, so you can rebuild them in Ranger and compare. The Ranger and Kerberos side is covered in Impala security with Ranger and Kerberos.

In place or side by side

There are two broad strategies, and the inventory usually decides between them.

In place. Stop the cluster, replace the binaries on every host, restart. It needs no extra hardware and is the procedure the Apache documentation describes. The cost is a full outage for the duration, and every incompatibility you missed appears in production at once. The documented procedure stops every daemon before replacing binaries, so plan on all of them changing version together.

Side by side. Build a second Impala cluster at the target version, pointed at the same storage and a metastore it is compatible with, and move users over gradually. Because Impala holds no data, both clusters can read the same tables. You can replay real queries on the new cluster and compare results before anyone depends on it, and rollback is a DNS or load-balancer change. The costs are hardware, and care with writes: two clusters writing the same tables need clear ownership, and the old cluster will not see new files written by the new one until it refreshes metadata.

For a 3.x to 4.x jump that also needs a metastore upgrade and a move to Ranger, side by side is usually worth it. For a minor upgrade within 4.x with a clean inventory, in place with a short maintenance window is usually fine.

The procedure: order, drain, restart

Whichever strategy you pick, the order of dependencies is fixed. Upgrade the metastore first, because the new Impala needs it. Whether your current Impala keeps working against the upgraded metastore depends on how that build was made (3.x could be built against Hive 2 or Hive 3), so prove it on a staging copy of the metastore before touching production. Put authorization in place second, so the new cluster enforces the same access from its first query. Upgrade Impala last.

For the Impala step, drain work before stopping daemons. The :shutdown() statement asks a daemon to stop accepting new fragments, wait for a grace period, and exit once running work finishes or a deadline passes; --shutdown_deadline_s defaults to one hour. Since 3.2 it addresses the daemon by its KRPC port.

# Drain executors one by one (run from impala-shell connected to a coordinator).
# Host and port are examples; use each executor's KRPC address.
:shutdown('worker-07.example.com:27000');

# After all daemons are stopped and binaries replaced, restart in this order:
#   1. statestored   on the statestore host
#   2. catalogd      on the catalog host
#   3. impalad       on every coordinator and executor
# Starting impalad before the statestore produces "Not connected" errors.

After the restart, confirm that every daemon reports the new version on its web UI, that the statestore lists every expected member, and that the catalog has loaded metadata. Run INVALIDATE METADATA on specific tables if their metadata looks stale; avoid running it without a table name on a large cluster, because it forces a reload of everything. The behaviour of the catalog after restart is described in the Impala catalog article.

Validating with query replay

Validation should compare behaviour, not just confirm that queries run. The approach that works is query replay: take a few hundred representative read-only queries from the inventory, run each against the old and new clusters, and compare row counts, checksums of the results and elapsed time.

import hashlib, time
from impala.dbapi import connect        # pip install impyla

def run(host, sql):
    conn = connect(host=host, port=21050)
    cur = conn.cursor()
    start = time.monotonic()
    cur.execute(sql)
    rows = cur.fetchall()
    elapsed = time.monotonic() - start
    conn.close()
    digest = hashlib.sha256(repr(sorted(map(repr, rows))).encode()).hexdigest()
    return len(rows), digest, elapsed

def compare(sql, old="impala-old.example.com", new="impala-new.example.com"):
    a, b = run(old, sql), run(new, sql)
    status = "same" if a[:2] == b[:2] else "DIFF"
    return status, a[0], b[0], round(b[2] / max(a[2], 1e-3), 2)

Sorting the row representations makes the comparison insensitive to row order, which differs between runs anyway. Investigate every DIFF. Expect some: a decimal rounding change, a different result for a query that relied on ||, or a query that now fails to parse. For each, decide whether to fix the query, set a compatibility query option for the session, or accept the new result. For slowdowns of more than about 1.5 times, compare query profiles from both versions; the usual causes are missing statistics, changed defaults for scan parallelism, or different join strategies. The Impala troubleshooting guide covers reading profiles.

Worked example: 3.4 on Hive 2 to 4.x

Consider a team running Impala 3.4 on a Hive 2 metastore with Sentry, 40 workers and roughly 15,000 queries a day. The inventory finds 37 distinct query shapes using || (31 string concatenations, 6 boolean ORs inside a BI tool's generated SQL), 4 LZO text tables used by a legacy ingest job, 3 older worker hosts without AVX, and around 200 Sentry roles.

The plan follows the dependency order. First, test the 3.4 build against a staging copy of a Hive 3 metastore; it passes, so the production metastore is upgraded. Second, convert the four LZO tables to Parquet with a one-off rewrite while 3.4 can still read them, and change the ingest job to write Parquet. Third, rebuild the roles as Ranger policies and compare effective access for a sample of users. Fourth, build a 4.x cluster on new hosts, replacing the three non-AVX machines rather than relying on the legacy flag, and replay 500 queries. The replay finds the six OR expressions, which the BI team rewrites, and two queries using HAVING ordinals. Finally, move dashboards over one team at a time, keeping the old cluster read-only for two weeks before decommissioning it.

Failure modes

The failures seen most often in Impala upgrades:

  • Silent result changes from || or DECIMAL_V2. Only replay with result comparison finds them before users do.
  • Authorization gaps. Users who could query yesterday get permission errors, or worse, users who should not see a table can. Compare access before cutover.
  • Daemons that will not start on hosts without AVX. Check every host, including spare and standby machines.
  • Wrong restart order producing "Not connected" errors until the statestore is up.
  • Unreadable tables in a format the new version dropped, discovered by a monthly job weeks after cutover.
  • Stale metadata after cutover in a side-by-side setup, because each cluster's catalog caches what it saw. Refresh tables written by the other cluster.
  • Untested clients. Old JDBC or ODBC drivers, shell scripts that parse impala-shell output, and native UDFs compiled against the old version.

Rollback

Because Impala keeps no data, rolling back the binaries is usually straightforward: reinstall the old version and restart in the same order. What makes rollback hard is everything else that moved. Once the metastore is on Hive 3, an older Impala that only supports Hive 2 cannot use it, so keep a metastore backup and a tested restore procedure. Once Sentry is gone, the old version has no authorization. Tables written by the new version may use features the old reader handles differently. Decide your rollback point for each step separately, and keep the side-by-side cluster or metastore backup until you are past it. Admission control pools and memory limits also need to be carried across; see Impala admission control.

What to do next

  1. Record your current Impala, metastore and authorization versions, and read every incompatible-change list between them and your target.
  2. Export 30 days of queries and run the pattern scan; review every hit by hand.
  3. Check every host for AVX, list LZO or other dropped formats, and export all roles and grants.
  4. Choose in place or side by side, and write the dependency order: metastore, authorization, then Impala.
  5. Build a replay harness, run it against a test cluster at the target version, and resolve every DIFF.
  6. Write the drain, restart order and rollback steps as a runbook, and rehearse it once on the test cluster.
Key takeaway: An Impala upgrade is mostly about the boundaries around the daemons: the Hive metastore, authorization, CPU features, SQL semantics and clients. Read every incompatible-change list on the path, inventory queries, tables, hosts and grants, upgrade the metastore and authorization before Impala, drain and restart in order, and prove results match with query replay before you cut over.