Impala was built to read, and its write support shows that history. Whether a statement such as UPDATE or DELETE works depends entirely on what is underneath the table: plain Parquet or text files on HDFS or S3 accept only appends and overwrites; Kudu tables accept row-level changes; Iceberg tables accept row-level changes in recent releases; Hive transactional tables accept inserts only. Teams that discover this mid-project end up with either a stalled migration or a cluster full of tiny files.

This page explains how Impala writes data to file-based tables, step by step, from the coordinator to the files on disk. It covers INSERT INTO and INSERT OVERWRITE, partitioned inserts and the hints that control them, why INSERT ... VALUES is a trap, LOAD DATA and CTAS, and the partition-rewrite pattern that stands in for UPDATE on file tables. Kudu and Iceberg row-level DML each have their own page, linked from the matrix below.

What DML means depends on the storage

Start with this table, because it answers the question most people are actually asking.

Table typeINSERT INTOINSERT OVERWRITEUPDATE / DELETEUPSERT / MERGERead more
HDFS / S3 files (Parquet, uncompressed text)YesYesNoNoThis page
HDFS / S3 files (ORC, Avro, RCFile, SequenceFile, compressed text)No (LOAD DATA, or write from Hive)NoNoNoThis page
KuduYesNoYesUPSERT yesImpala + Kudu
Iceberg (format v2)YesYesDELETE since 4.3; UPDATE in later 4.x (check your release)MERGE since 4.5Impala + Iceberg
Hive insert-only transactionalYesCheck your releaseNoNoThis page
Hive full-ACIDNo (write from Hive)NoNoNoWrite from Hive

Impala writes only Parquet and uncompressed text files itself; for ORC, Avro and the other read-only formats, write the files with Hive, then run REFRESH table_name in Impala.

Every Impala DML statement auto-commits at the end of the statement. There are no multi-statement transactions, so "insert into the fact table and the audit table atomically" is not something Impala can promise on any storage. For insert-only transactional tables the Impala documentation gives an ordering guarantee per table: if a reader does not see a committed insert, it will not see any insert committed after it. Insert-only tables must be managed, file-format tables such as Parquet, Avro or text, and their transactional property cannot be altered later.

The write path of an INSERT

An INSERT ... SELECT into a file-based table runs like a query whose final operator writes files instead of returning rows. The coordinator plans it, the scan and compute fragments produce rows on every executor, an optional exchange redistributes rows by partition key, and a table sink on each executor writes data files. Those files go first into a hidden staging directory under the table directory, named _impala_insert_staging in Impala 2.0.1 and later. When every fragment has finished, the coordinator moves the files into their final partition directories and tells the catalog about any new partitions and files, which then broadcasts the change to every coordinator.

ClientINSERT ... SELECTCoordinatorplan, fragmentsCatalog servicetable, partitionsExecutors (one fragment instance per node)Scan + computesource rowsExchange / sortSHUFFLE, CLUSTEREDTable sinkParquet writertable_dir/_impala_insert_staging/...new files written here firsttable_dir/year=2026/month=10/files moved into place at commitOther enginesHive, Spark: REFRESHmetadatacommitcatalog update
The write path of an INSERT into a partitioned Parquet table. Readers in Impala see the new files only after the catalog update; Hive and Spark need their own refresh.

Three consequences follow. First, the work is parallel: each executor writes its own files, so one statement produces many files, and file names are unique so several INSERT INTO statements can run against the same table at once. Second, the files are owned by the operating-system user the Impala daemons run as, typically impala, whatever the connected user's identity; that user needs write permission on the destination and read permission on any source files, and downstream jobs that expect other ownership will fail. Third, the commit is a set of file moves, not an atomic swap; on an HDFS-style filesystem those moves are cheap, but on object stores a "move" is a copy and delete, which is slower and can leave a partial set visible to engines that list directories directly.

INSERT INTO, INSERT OVERWRITE and partitions

INSERT INTO appends: existing files stay, new files are added. INSERT OVERWRITE replaces: the old files of the target are deleted as part of the statement, and the documentation is explicit that overwritten files are removed immediately rather than going through the HDFS trash. Treat every INSERT OVERWRITE against a production table as an irreversible delete followed by a load, and back up or snapshot before running one by hand.

Partitioned tables add a choice between static and dynamic partition values.

-- Static: every row goes to one named partition.
INSERT OVERWRITE sales PARTITION (year=2026, month=10)
SELECT order_id, customer_id, amount, ts FROM staging_sales WHERE ts >= '2026-10-01';

-- Dynamic: partition values come from the last columns of the SELECT, in order.
INSERT INTO sales PARTITION (year, month)
SELECT order_id, customer_id, amount, ts, year(ts), month(ts) FROM staging_sales;

-- Mixed: year fixed, month dynamic.
INSERT INTO sales PARTITION (year=2026, month)
SELECT order_id, customer_id, amount, ts, month(ts) FROM staging_sales WHERE year(ts) = 2026;

With a dynamic INSERT OVERWRITE, the partitions that receive rows are replaced and partitions that receive no rows are left alone; confirm this on a scratch table for your release before relying on it. That makes it the natural tool for reprocessing a window of days, and also the source of a classic mistake: a source query with a bug that returns no rows for one day leaves the old, wrong data for that day in place without any error. Verify row counts per partition after reprocessing rather than assuming the overwrite touched everything you meant it to.

Partitioned inserts: hints and memory

A dynamic-partition insert into Parquet is memory-hungry, because a Parquet writer buffers a whole row group per open file. If each executor receives rows for 400 partitions, it holds 400 open writers at once. Impala provides hints, written in comment form after the PARTITION clause, to control how rows reach the writers. The older square-bracket form still parses but is deprecated.

-- Redistribute rows by partition key so each partition is written by one executor:
-- fewer, larger files and less memory per node, at the cost of a network exchange.
INSERT INTO sales PARTITION (year, month) /* +SHUFFLE */
SELECT order_id, customer_id, amount, ts, year(ts), month(ts) FROM staging_sales;

-- Write where the rows already are: no exchange, but every node may open every partition.
INSERT INTO sales PARTITION (year, month) /* +NOSHUFFLE */
SELECT ...;

-- Sort rows by partition key before the sink, so each writer opens one partition at a time.
INSERT INTO sales PARTITION (year, month) /* +CLUSTERED */
SELECT ...;

The rule of thumb: when there are many partitions, use SHUFFLE and keep CLUSTERED in place so memory stays bounded; when the source is already laid out by partition key and the partition count is small, NOSHUFFLE saves the exchange. Check the plan with EXPLAIN before trusting a hint, and watch the per-node peak memory in the query profile on the first run.

INSERT VALUES and the small-file problem

The INSERT ... VALUES form is fine for a ten-row lookup table and wrong for anything else. Each statement runs on one node and writes at least one new data file, so a loader that inserts one row per statement creates one file per row. HDFS keeps every file's metadata in NameNode memory, Impala's catalog tracks every file, and a scan pays a per-file cost for opening, footer reads and scheduling. A table with a million 2 KB files is slower to query and slower to refresh than the same data in forty files.

The fix is to land small writes somewhere designed for them and move data into Parquet in bulk: stream rows into Kudu, or into a text or Avro staging table, then run a periodic INSERT ... SELECT into the Parquet table. The general problem, and compaction strategies for tables already damaged, are in Hive Small Files. Target file size for Parquet writes is controlled by the PARQUET_FILE_SIZE query option; set it at session level for big batch loads rather than changing the cluster default.

LOAD DATA and CREATE TABLE AS SELECT

LOAD DATA INPATH '/landing/2026-10-03/' INTO TABLE raw_events does not read the files; it moves them from the source directory into the table directory and updates metadata. It is fast, but the files must already be in the table's format, and the source directory is emptied, which surprises anyone who expected a copy. Use it for landing data that another tool has already written correctly.

CREATE TABLE ... AS SELECT (CTAS) combines DDL and an insert in one statement, and is the cleanest way to rewrite data in a new layout:

CREATE TABLE sales_v2
PARTITIONED BY (year, month)
STORED AS PARQUET
AS SELECT order_id, customer_id, amount, ts, year(ts) AS year, month(ts) AS month
FROM sales;

Worked example: correcting rows without UPDATE

Worked example. A billing table, invoices, is Parquet on HDFS, partitioned by year and month, about 2 billion rows. Finance discovers that 12,000 invoices in September 2026 carry the wrong tax code and must be corrected. UPDATE invoices SET ... fails on a file table, so the correction is a partition rewrite.

-- 1. Build the corrected September partition in a scratch table and check it.
CREATE TABLE invoices_fix_2026_09 STORED AS PARQUET AS
SELECT i.invoice_id, i.customer_id, i.amount,
       COALESCE(f.correct_tax_code, i.tax_code) AS tax_code, i.issued_at
FROM invoices i LEFT JOIN tax_fixes f ON i.invoice_id = f.invoice_id
WHERE i.year = 2026 AND i.month = 9;

SELECT COUNT(*) FROM invoices WHERE year = 2026 AND month = 9;        -- must match
SELECT COUNT(*) FROM invoices_fix_2026_09;
SELECT COUNT(*) FROM invoices_fix_2026_09 f JOIN tax_fixes t USING (invoice_id)
WHERE f.tax_code = t.correct_tax_code;                                  -- expect 12000

-- 2. Replace exactly that partition. Overwritten files skip the trash: back up first.
INSERT OVERWRITE invoices PARTITION (year=2026, month=9)
SELECT invoice_id, customer_id, amount, tax_code, issued_at FROM invoices_fix_2026_09;

-- 3. Refresh statistics for the changed partition, then drop the scratch table.
COMPUTE INCREMENTAL STATS invoices PARTITION (year=2026, month=9);
DROP TABLE invoices_fix_2026_09;

The rewrite touches one partition of perhaps 60 million rows instead of the whole table, which is why partitioning by time pays off even for tables that are "never updated". If corrections like this happen weekly rather than yearly, that frequency is the signal to move the table to Iceberg or Kudu, where a row-level UPDATE or MERGE replaces the whole procedure.

After the write: stats, metadata and other engines

Impala itself sees its own inserts immediately after the statement returns, because the coordinator updated the catalog. Two other things do not happen automatically. Statistics are not recomputed, so the planner keeps using old row counts and can pick a bad join order; the documentation recommends running COMPUTE STATS after inserts, and incremental stats keep that cheap for partitioned tables, as described in Impala Stats and Metadata. And other engines know nothing: if Hive or Spark wrote the files instead, Impala needs REFRESH table (or REFRESH table PARTITION (...)) to see them, and Hive cannot safely read the table while Impala is still staging an insert into it. The catalog mechanics are covered in Impala Metadata, in depth.

In a cluster with several coordinators behind a load balancer, set SYNC_DDL=1 in pipelines where the next statement may land on a different coordinator; it makes the statement wait until all coordinators have received the metadata change.

Failure modes

  • A failed insert leaves debris. If an executor dies mid-write, the statement fails and the staging directory is normally cleaned up, but a crash of the coordinator itself can leave files under _impala_insert_staging. Monitor for stale staging directories and delete them once no insert is running.
  • Out of memory on partitioned inserts. Too many open Parquet writers per node. Use /* +SHUFFLE */ with clustering, or split the load by partition range.
  • Permission errors downstream. Files owned by impala that a Spark job cannot overwrite. Fix with directory ACLs or a Ranger policy, not by running Impala as a superuser.
  • Empty-result overwrite. A dynamic overwrite whose source returns nothing for a partition leaves stale data in place; a static overwrite with an empty source empties the partition. Both are silent.
  • Stale plans after big loads. Missing stats after a large insert produce broadcast joins of huge tables. Make stats part of the load job.
  • Expecting UPDATE to work. The statement fails at analysis on file tables; the outage is the migration plan that assumed it would work.

Trade-offs

File tables with append and overwrite are the simplest and fastest option to scan, and every engine can read them, but corrections cost partition rewrites and small writes cost small files. Kudu gives true row-level changes and fresh data at the price of a separate storage service. Iceberg gives row-level DML, snapshots and engine-neutral metadata, at the price of delete files that must be compacted and features that vary by Impala release. Insert-only transactional tables add atomic visibility per insert without row-level changes. Choose by how often rows change after they land, not by which syntax looks familiar.

What to do next

  1. List every Impala table you write to and classify it in the matrix above; flag any job that expects UPDATE or DELETE on a file table.
  2. Replace every INSERT ... VALUES loader with a staging table or Kudu plus a periodic INSERT ... SELECT.
  3. Add SHUFFLE and CLUSTERED hints to large dynamic-partition inserts and check peak memory in the profile.
  4. Wrap every INSERT OVERWRITE in a runbook with a backup step and per-partition row-count checks.
  5. Make COMPUTE INCREMENTAL STATS the last step of each load job.
  6. Add REFRESH steps wherever Hive or Spark writes tables that Impala reads, and SYNC_DDL where coordinators are load balanced.
  7. If row corrections are routine, plan a move to Iceberg or Kudu instead of scripting more partition rewrites.
Key takeaway: Impala DML depends on storage: file tables on HDFS or S3 accept INSERT INTO and INSERT OVERWRITE only, Kudu and recent Iceberg releases accept row-level changes, and every statement auto-commits on its own. Inserts write through a staging directory and are moved into place, owned by the impala user. Shuffle and cluster large partitioned inserts, avoid INSERT VALUES, treat overwrites as irreversible, compute stats after loads, and move tables that need routine corrections to Iceberg or Kudu.