TT Lab
Get started
Learn Learning paths Courses

Lakehouse Table Format — Understanding Apache Iceberg Through Its Metadata

A partition is a rule, not a column — hidden partitioning and partition evolution

Continue in TT Lab

In one line

Iceberg computes partition values not as a column that people create but through a transform rule applied to a source column (for example days(order_ts)), records them in the manifests, and lets readers skip files with conditions on the source column alone. Rules pile up with spec IDs, so even if you change them midway, old files are still read with the old rule.

Why Hive-style partitioning was a problem

The example in the Partitioning docs is typical. To split a log table by date in Hive, you create a separate column called event_date, and the writer has to compute that value from event_time and insert it. Problems follow one after another.

How it works — transforms and the partition spec

An Iceberg partition spec has, for each field, a source column ID, a partition field ID, a transform, and a name. When the engine writes rows, it computes the transform and records it in the manifest as the partition value of that file. All the rows of one file have the same partition value. The list of transforms is as follows.

Transform Result Used for
identity The value as it is Category columns with few distinct values
year · month · day · hour Years, months, days, or hours counted from 1970 Timestamp columns
bucket[N] The remainder of the hash divided by N Keys with many distinct values (customer number)
truncate[W] The value truncated to width W Numeric ranges, the beginning of a string
void Always null When removing a field in version 1

bucket drops the sign bit of a 32-bit Murmur3 (x86, seed 0) hash and divides it by N. Because the spec fixes the hash, the same value goes to the same slot even if Spark writes and Python reads.

Readers use only source column conditions such as order_ts >= X. Scan planning turns that condition into a partition condition (order_ts_day >= day(X)) and compares it with the partition values in the manifests. This transform is computed in the "inclusive" direction, so a file that may contain rows matching the condition is never left out. Users do not need to know how the partitions were divided — that is why it is called "hidden" partitioning.

Partition evolution — old files stay as they are

Suppose the data has grown and month granularity has become too coarse. According to the Evolution docs, even if you change the spec, old data stays under the old spec and only new data is written with the new layout. Specs pile up in a list, and each manifest remembers the spec ID it was written with. When planning, each spec gets its own converted condition to select files, which the docs call split planning. The docs also state clearly that partition evolution is a metadata operation and does not hurry to rewrite files.

In Spark you use ALTER TABLE to add a field (ADD PARTITION FIELD), remove one (DROP PARTITION FIELD), and change one (REPLACE PARTITION FIELD … WITH …).

ALTER TABLE lake.demo.events REPLACE PARTITION FIELD months(event_ts) WITH days(event_ts);

What it looks like in the field

"We changed the partitioning, so we have to rewrite the old data too." This is the most common misunderstanding. Old month-granularity files are still read exactly as they are. Rewriting is optional and needs a reason — for example, when you often query the old period by day and want to reduce how much you read. Even then, you do it as a separate job such as compaction (rewrite_data_files).

A one-day condition reads a month. If you put a one-day condition on a month-granularity file, the partitions force you to select that month's files as a whole. Column statistics (lower and upper bounds) inside the file help, but if a file holds the whole month evenly, they do not. The difference between the rows of the files that planning opens and the rows that actually match is that cost.

We split by customer number and now there are tens of thousands of partitions. If you split a column with many distinct values by identity, small files explode. It is right to fix the number of slots with bucket.

What really matters in practice

What you will do in the next lab

You load March into a table split by months(order_ts) and use pyiceberg to plan a one-day condition and measure how many files and rows it opens. After changing the partitioning to days(order_ts), you load April and confirm that the March files are still spec 0 and only the April files are spec 1. You compare the number of rows read for a day in March and a day in April, then split the customer table with bucket(4, customer_id) and check that the slots were computed as the spec says.