Partitioning is chosen when a table is small and the queries are guesses. A year later the table is a hundred times larger, the queries filter on something else, and the daily partitions are either too big to scan or too small to be efficient. Partition evolution is the ability to change the layout without breaking readers, and ideally without rewriting the data already written.

Classic Hive tables and Iceberg tables answer this very differently. This article explains both from the storage model up: what classic Hive can evolve cheaply, what forces a full rewrite, how Iceberg keeps several partition layouts alive in one table, and how queries prune across them. Iceberg tables in Hive are introduced in Hive Iceberg tables; this page focuses on changing partitioning over time.

Advertisement

Two models of a partition

In a classic Hive table, a partition is a directory such as .../web_events/dt=2026-09-01/ plus a row in the metastore. The partition value lives in the path, not in the files. Queries prune by matching their WHERE clause against the partition column, so the user has to filter on dt explicitly, and changing what defines a partition means moving files into different directories. The metastore side is explained in Hive metastore.

In an Iceberg table, a partition is a value computed from a source column by a transform such as day(ts), hour(ts), bucket(16, id) or truncate(5, code). The table's metadata stores one or more partition specs, each with an id, and each data file is recorded in a manifest along with the spec it was written under and its partition value. Directories are an implementation detail. Queries filter on the source column (ts), and Iceberg derives the partition filter. That indirection is what makes evolution possible.

What classic Hive can evolve

Each Hive partition carries its own storage descriptor: column list, SerDe, input and output format and location. A table's metadata describes how new partitions are created, but existing partitions keep what they were created with. That gives you a few cheap evolutions and one important trap.

-- add a column; CASCADE updates the table AND every existing partition
ALTER TABLE web_events ADD COLUMNS (referrer STRING) CASCADE;

-- the default is RESTRICT: only the table metadata changes
ALTER TABLE web_events ADD COLUMNS (campaign STRING);

-- different file formats per partition are allowed
ALTER TABLE web_events PARTITION (dt='2026-09-01') SET FILEFORMAT PARQUET;

-- change the declared type of a partition column (metastore only)
ALTER TABLE web_events PARTITION COLUMN (dt DATE);

The trap is RESTRICT. After adding campaign without CASCADE, existing partitions still list the old columns. Rewriting one of those partitions with INSERT OVERWRITE can write the new column into the files while the partition's own metadata still omits it, so queries return NULL for it. Use CASCADE when you mean every partition, or fix individual partitions with ALTER TABLE ... PARTITION (...) ADD COLUMNS.

Mixed formats work because each partition is read with its own descriptor, which is how teams move a table from text to ORC or Parquet one partition at a time. Changing a partition column's type only changes the metastore; existing directory values must still parse as the new type.

What classic Hive cannot do is change the partition keys. Going from dt to (dt, hour), or from date to region, means a new table and a rewrite.

Advertisement

The classic rewrite, and what it costs

The standard approach is to create a new table with the new keys and copy with a dynamic-partition insert, then swap names.

SET hive.exec.dynamic.partition=true;
SET hive.exec.dynamic.partition.mode=nonstrict;      -- default is strict
SET hive.exec.max.dynamic.partitions=20000;           -- default 1000
SET hive.exec.max.dynamic.partitions.pernode=2000;    -- default 100

CREATE TABLE web_events_v2 (user_id BIGINT, url STRING, ts TIMESTAMP)
PARTITIONED BY (dt STRING, hr INT) STORED AS ORC;

INSERT OVERWRITE TABLE web_events_v2 PARTITION (dt, hr)
SELECT user_id, url, ts, dt, hour(ts) FROM web_events;

ALTER TABLE web_events RENAME TO web_events_old;
ALTER TABLE web_events_v2 RENAME TO web_events;

Every byte is read and written again, the partition count multiplies (a year of hourly partitions is 8,760 directories and metastore rows), and every query and job that filtered on the old columns must be checked. The limits above exist to stop an accidental explosion of partitions; raise them deliberately. Mechanics of dynamic inserts are covered in Hive dynamic partitioning. In practice teams often rewrite only recent history and keep the old table for the rest, which pushes the problem onto every query that has to union the two.

How Iceberg keeps several layouts alive

An Iceberg table's metadata holds a list of partition specs and a pointer to the default spec for new writes. Changing the spec adds a new spec and moves the pointer. Nothing else happens: per the Iceberg documentation, partition evolution is a metadata operation and does not eagerly rewrite files. Old data stays in the old layout, new data uses the new one.

Two partition specs in one Iceberg table, planned separatelymetadata.jsonspecs: [0: day(ts), 1: hour(ts)]current snapshotmanifest listmanifests, spec_id 0partition value: ts_daymanifests, spec_id 1partition value: ts_hourdata files before changeone group per daydata files after changeone group per hourquery: ts BETWEEN 22:00 and 02:00filter projected per spec: day range | hour rangeOld files are never rewritten by the spec change; each spec prunes with its own derived filter.
The manifest list points at manifests written under different specs; a query's filter on ts is projected into each spec's partition space separately.

Reads use what the Iceberg docs call split planning: each partition layout plans its files separately, using the filter derived for that layout. A query filtering on ts is turned into a day-range predicate for files written under the day spec and an hour-range predicate for files under the hour spec. Results are correct across the boundary because both sets are filtered against the same source column, and the user never references a partition column.

The change is cheap in a second sense too. Per Cloudera's documentation, a spec change produces a new metadata.json and a commit but not a new snapshot, so it does not show up as a data change in time travel.

The statements: Hive, Impala and Spark

Hive 4 and Impala share the SET PARTITION SPEC statement. It states the complete new spec, not a change to the old one.

-- Hive 4: create an Iceberg table partitioned by day of ts
CREATE TABLE web_events (user_id BIGINT, url STRING, ts TIMESTAMP)
PARTITIONED BY SPEC (day(ts))
STORED BY ICEBERG;

-- later: switch new writes to hourly and add a bucket on user_id
ALTER TABLE web_events SET PARTITION SPEC (hour(ts), bucket(16, user_id));

Cloudera's documentation gives the general form ALTER TABLE t SET PARTITION SPEC (TRUNCATE(5, level), HOUR(event_time), BUCKET(15, message), price), and an example that uses VOID(col) for fields being retired. Spark with the Iceberg SQL extensions offers finer-grained statements instead:

ALTER TABLE db.web_events ADD PARTITION FIELD bucket(16, user_id);
ALTER TABLE db.web_events DROP PARTITION FIELD bucket(16, user_id);
ALTER TABLE db.web_events REPLACE PARTITION FIELD day(ts) WITH hour(ts);

The Iceberg Spark DDL page notes that dropping a partition field leaves the source column in the schema, changes the schema of metadata tables such as files (which can break metadata queries), and that dynamic partition overwrite behaviour changes whenever partitioning changes. Impala's side of this is described in Impala and Iceberg.

Migrating a classic Hive table first

To evolve an existing classic table's partitioning without the rewrite, convert it to Iceberg in place. Cloudera documents an ALTER TABLE that switches the storage handler:

ALTER TABLE web_events SET TBLPROPERTIES (
  'storage_handler'='org.apache.iceberg.mr.hive.HiveIcebergStorageHandler',
  'format-version'='2');

Per that documentation, Hive reads the footers of the existing data files to build Iceberg metadata and commits all files in a single commit; the supported source formats are Avro, Parquet and ORC. The original partition columns become identity partitions in the first spec. From there SET PARTITION SPEC works as above. Take a backup of the metastore entry before converting, and test on a copy: text-format tables are not covered, and any engine that reads the table must support Iceberg.

Worked example: day to hour

A clickstream table was created in 2025 partitioned by day(ts) (spec 0). By September 2026 it receives about 2 TB a day, and most dashboard queries look at the last two hours. Each query scans a whole day's files, around 2 TB, to use about 170 GB.

On 2026-09-01 at 00:00 the team runs ALTER TABLE clicks SET PARTITION SPEC (hour(ts)). Writers pick up spec 1 on their next commit. Nothing is rewritten.

Now a query runs with ts BETWEEN '2026-08-31 22:00' AND '2026-09-01 01:59'. Iceberg plans spec-0 files with the projected predicate day = 2026-08-31 (and 2026-09-01, which has no spec-0 files), so it reads the whole of the 31st, about 2 TB, and then filters rows. It plans spec-1 files with hour in [00, 01] on the 1st, about 170 GB. The answer is correct; the old half is no faster than before.

A week later queries rarely touch August, so the improvement is real where it matters. If old data still needs hourly pruning, rewrite just those days with Iceberg's data-file rewrite procedure so they are written under the current spec, and expire old snapshots afterwards so storage is reclaimed. Inspect the result with the files metadata table, which records each file's spec_id.

Failure modes

  • Dynamic overwrite after the change. INSERT OVERWRITE in dynamic mode replaces the partitions the new rows land in; under a new spec those are different partitions, so a backfill job can silently stop replacing the old data. Use explicit filters or MERGE for corrections.
  • Filtering on a derived column. A query on a stored event_date column instead of ts cannot use the hour partitions. Filter on the transform's source column.
  • Expecting old data to speed up. The spec change only affects new writes. Plan targeted rewrites if history needs the new layout.
  • Bucket count changes. Moving from bucket(16, id) to bucket(32, id) is a new field; old and new files are not co-bucketed, so bucket-aware joins lose that property across the boundary.
  • Too-fine partitions. Hourly partitions on a table with little data per hour create tiny files and large manifests. Check bytes per partition before evolving, and schedule compaction.
  • Classic RESTRICT drift. In non-Iceberg tables, columns added without CASCADE read as NULL in rewritten old partitions.

Trade-offs

ApproachRewrite costReader impactBest when
Classic Hive CASCADE / per-partition formatNoneTransparentAdding columns, migrating file formats
Classic rewrite to new keysFull tableQueries must change partition filtersTable cannot move to Iceberg
Migrate to Iceberg, then SET PARTITION SPECMetadata onlyQueries filter on source columnsLong-lived tables whose access pattern changes
Iceberg evolution plus targeted rewriteOnly the rewritten rangesOld ranges gain new pruningHot history needs the new layout

What to do next

  1. List your largest classic Hive tables and the columns their top queries filter on; mark any where those differ from the partition keys.
  2. Check recent ADD COLUMNS statements for missing CASCADE, and compare partition column lists with the table's in the metastore.
  3. For tables to evolve, convert a copy to Iceberg with the storage-handler ALTER and confirm every reading engine handles it.
  4. Pick the new spec from measured bytes per partition; aim for partitions large enough to avoid small files.
  5. Run SET PARTITION SPEC (or the Spark ADD/REPLACE PARTITION FIELD equivalents) in a quiet window and verify new files carry the new spec_id.
  6. Audit jobs that use dynamic INSERT OVERWRITE on the table before the change.
  7. Decide whether any history needs rewriting under the new spec, and schedule compaction and snapshot expiry.
Key takeaway: Classic Hive stores partition values in directory paths and a storage descriptor per partition, so you can add columns (with CASCADE), mix file formats and retype partition columns cheaply, but changing partition keys means rewriting the table. Iceberg stores partitioning as transforms of source columns in versioned specs, so changing it is a metadata commit: old files keep their layout, new files use the new one, and split planning prunes each layout with its own derived filter. Migrate long-lived tables to Iceberg in place, evolve with SET PARTITION SPEC, watch dynamic overwrites and metadata queries, and rewrite only the history that needs the new layout.