Create an Apache Iceberg Table with Spark SQL
This procedure creates a partitioned Apache Iceberg table from the spark-sql shell, writes a few rows, and confirms the table is healthy by reading its metadata tables. It uses a local, path-based catalog so you can run it on a laptop before pointing the same SQL at a production catalog.
Prerequisites
Section titled “Prerequisites”- Apache Spark 4.1 (Scala 2.13) with Java 17 or 21. Spark 3.5 works too if you swap the runtime package shown below.
- The Iceberg Spark runtime for Iceberg 1.11.0:
org.apache.iceberg:iceberg-spark-runtime-4.1_2.13:1.11.0. For Spark 3.5 useorg.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.11.0. - Write access to a local directory for the warehouse.
Start the shell with the runtime package, the Iceberg SQL extensions, and a catalog named local:
spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-4.1_2.13:1.11.0 \ --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \ --conf spark.sql.catalog.local=org.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.local.type=hadoop \ --conf spark.sql.catalog.local.warehouse=$PWD/warehouseTo use a Hive Metastore or REST catalog instead, change the type and add a uri, for example --conf spark.sql.catalog.local.type=rest --conf spark.sql.catalog.local.uri=http://localhost:8181. Every step below stays the same.
-
Create a namespace to hold the table.
CREATE NAMESPACE IF NOT EXISTS local.db; -
Create the table.
USING icebergmakes it an Iceberg table, and thePARTITIONED BYclause uses transforms, so readers filter onevent_tsanduser_idand never need to know how the data is laid out on disk (hidden partitioning).CREATE TABLE local.db.events (event_id bigint,user_id bigint,event_type string,event_ts timestamp)USING icebergPARTITIONED BY (day(event_ts), bucket(16, user_id));The other transforms are
year(ts),month(ts),hour(ts), andtruncate(L, col). New tables use format version 2 by default. -
Write rows in a single statement, which produces exactly one commit.
INSERT INTO local.db.events VALUES(1, 101, 'login', TIMESTAMP '2026-09-28 08:15:00'),(2, 102, 'purchase', TIMESTAMP '2026-09-28 09:30:00'),(3, 101, 'logout', TIMESTAMP '2026-09-29 17:45:00'),(4, 103, 'login', TIMESTAMP '2026-09-29 18:05:00'); -
Query the table with a filter on the source column. Iceberg maps the predicate to the
day(event_ts)partition for you.SELECT event_type, count(*) AS eventsFROM local.db.eventsWHERE event_ts >= TIMESTAMP '2026-09-29 00:00:00'GROUP BY event_type;
Done when
Section titled “Done when”Run each check. If all four match, the table exists, is partitioned the way you declared, and has one clean commit.
-
One snapshot, an
append, with four records:SELECT snapshot_id, operation,summary['added-records'] AS added_records,summary['total-records'] AS total_recordsFROM local.db.events.snapshots;Expected: one row,
operation=append,added_records=4,total_records=4. -
The history has that snapshot as the current state:
SELECT made_current_at, snapshot_id, parent_id, is_current_ancestorFROM local.db.events.history;Expected: one row,
parent_idisNULL,is_current_ancestoristrue, andsnapshot_idmatches check 1. -
The data files add up to the rows you wrote:
SELECT count(*) AS data_files, sum(record_count) AS recordsFROM local.db.events.files;Expected:
records=4anddata_filesis at least1. -
The partitions reflect the transforms:
SELECT partition, record_count, file_countFROM local.db.events.partitions;Expected: each
partitionvalue is a struct with a day field and a bucket field, andrecord_countsums to4across rows. The two days in the sample mean you see at least two distinct day values.
Next steps
Section titled “Next steps”- Compact data files once many small writes pile up.
- Expire snapshots to cap metadata growth and storage cost.
- Migrate a Hive table instead of starting from an empty table.
Sources
Section titled “Sources”Verified against the Apache Iceberg 1.11.0 documentation:
Work with Alex