Skip to main content

13 Billion Rows into Apache Iceberg in 7 Minutes, for $3.69

· 10 min read
Cazpian Engineering
Platform Engineering Team

12.96 billion rows of raw Parquet into a query-ready Apache Iceberg table in 7 minutes 4 seconds on 32 Graviton workers

A post making the rounds on LinkedIn showed a serverless Spark platform processing "1 TB (40 billion rows)" in 14 minutes 5 seconds: 91 raw Parquet files ingested into a table clustered by date and warehouse. Read closely, the 40 billion is rows read — the post notes the engine read three times more rows than it wrote, 13.3 billion, because clustering on write re-processed the data. It is a good benchmark all the same: no toy data, no caching tricks, a real shuffle and a real write, and a table you can query the moment the job commits.

So we ran the same shape of job on Cazpian.

12,959,823,744 rows — 1.13 TB of raw Parquet — landed in an Apache Iceberg table partitioned by day and sorted within each day, in a median of 7 minutes 4 seconds on 32 AWS Graviton workers. Every row counted, in three runs out of three. About $3.69 of compute per load. Starting from nothing — no cluster running at all — clicking Create Pool to the committed table took 12 minutes 3 seconds.

The job​

We kept it deliberately plain, because the interesting part of this benchmark is the work the engine cannot avoid:

Input12,959,823,744 rows · 1.13 TB zstd Parquet · 117 files of ~9.7 GB · 38 columns
OutputApache Iceberg v2 table, partitioned by days(sales_date) (1,837 partitions), sorted by sales_date, warehouse_name within each partition, zstd, 512 MB target files
Statementone INSERT INTO … SELECT * — one atomic Iceberg commit
Computea Cazpian compute pool: 32 workers × 16 vCPU / 120 GB on AWS Graviton (arm64), Spark 4.1, Iceberg 1.11
ClockINSERT start → Iceberg commit, on a running pool (plus a separate cold-start measurement below)

The data is TPC-DS catalog_sales at scale factor 1000, joined once to date_dim and warehouse so the table carries a real date and a real warehouse name — the columns a query-ready layout is built around. Like the original post, the input is a fixed set of large Parquet files that already exist before the clock starts; we built one 1.44 billion-row replica with Spark and copied it nine times in S3 (more on that under Honest notes).

"Query-ready" is the expensive part. Iceberg does not just append files here: every one of the 12.96 billion rows has to be shuffled to the executor that owns its day, then sorted by warehouse inside that day, then written as right-sized files. That is a full shuffle of 1.13 TB of compressed data — 2.18 TB of shuffle blocks, because Spark's shuffle format is far less compact than zstd Parquet.

CREATE TABLE lake.bench.sales_clustered
USING iceberg
PARTITIONED BY (days(sales_date))
TBLPROPERTIES ('write.parquet.compression-codec' = 'zstd',
'write.target-file-size-bytes' = '536870912')
AS SELECT * FROM landing.raw_sales LIMIT 0;

ALTER TABLE lake.bench.sales_clustered
WRITE DISTRIBUTED BY PARTITION LOCALLY ORDERED BY sales_date, warehouse_name;

-- the timed statement
INSERT INTO lake.bench.sales_clustered
SELECT * FROM landing.raw_sales;

The result​

RunTimed loadThroughputRows committed
1 — first job on a freshly started pool431.7 s30.0 M rows/s12,959,823,744 ✅
2424.5 s30.5 M rows/s12,959,823,744 ✅
3423.0 s30.6 M rows/s12,959,823,744 ✅
Median7 min 4.5 s30.5 M rows/s · 2.66 GB/sexact every time

The three runs landed within 2% of each other. Each produced the same table: about 957 GB of zstd Parquet in 8,331 files (median 87 MB), one Iceberg snapshot, 1,837 day partitions. Rows in equal rows out, to the row: the engine's scan, shuffle and write counters and Iceberg's own commit summary all read 12,959,823,744. And the whole 1.13 TB was really read — during the read stage the 32 workers pulled 1.07–1.31 TB from S3 in every run, measured from their network counters.

Cost. The pool — 32 Graviton workers plus a driver — costs $31.28 an hour at AWS Fargate on-demand list price. A 424.5-second load is $3.69 of compute. The cold end-to-end run, including bringing the cluster up from zero, is about $6.30.

Cold start. From clicking Create Pool with nothing running: 32 workers were up in 4 minutes 16 seconds, the first job started 17 seconds later, and the table committed at 12 minutes 3 seconds. Even counting cluster creation, the whole thing finishes inside the reference run's 14 minutes 5 seconds.

Side by side​

Reference post (Databricks Serverless)Cazpian on GravitonDifference
Rows written13.3 billion12.96 billionsame scale (97%)
Input read1.04 TB, as reported¹1.13 TB, measuredours is larger
Rows read to write them40.1 billion (3 passes)12.96 billion (1 pass)3.1× less work
Load time, clustered table14 min 5 s7 min 4 s2.0× faster
From zero, cluster start-up included— (serverless)12 min 3 sstill faster
Write throughput15.7 M rows/s30.5 M rows/s1.9×
Read throughputover 1.2 GB/s2.66 GB/s2.2×
Sub-10-minute targetmissedmet, with clustering
Compute cost per load~70 DBUs² ≈ $25–65≈ $3.69~7–18× cheaper
Tuning needed32 MB file splits to trigger autoscalingnone — default 128 MB splits

¹ The post describes 91 files of 6–8 GB, which is 0.55–0.73 TB on disk, so our 1.13 TB input is at least as large and likely larger. ² The DBU figure comes from the post's discussion, priced at list serverless DBU rates (platform included); ours is AWS infrastructure at Fargate on-demand list price.

Why read every row three times? Clustering a table on write needs the data grouped by the clustering keys before files are written, and the reference engine got there in several passes over the input. Here the same layout — day partitions, sorted by warehouse inside each — comes out of a single hash shuffle followed by a local sort, so each of the 12.96 billion rows is read once, shuffled once and written once.

We did not run the other platform ourselves; its numbers are the ones published with the post, and the dollar range depends on which serverless rate applies. The two cost figures also price different things: a serverless DBU bundles the platform into the compute price, while Cazpian is licensed as a flat annual fee per deployment, so the per-job number above is the whole marginal cost of the load.

Where the 7 minutes go​

We read the engine's own counters from the live driver after every run. The load is two stages and a commit:

PhaseGravitonx86 (same job)
Read 1.13 TB of Parquet, shuffle 2.18 TB by day3 min 29 s3 min 38 s
Read the shuffle, sort within each day, write 957 GB of Iceberg files3 min 32 s5 min 37 s
Iceberg commit (one snapshot, 8,331 files)~1 s~1 s
Total424.5 s559.3 s

The first stage is bound by reading and repartitioning, and both architectures do it in about the same time. The second stage is where Graviton wins — sorting billions of wide rows and encoding them as compressed Parquet is dense CPU work, and the Graviton workers did it 1.6× faster. Overall the same job was 1.32× faster on Graviton and about 39% cheaper per load ($3.69 against $6.05). It matches what we saw on TPC-DS on Graviton, and it is why we are moving Cazpian compute pools to Graviton by default.

Per worker during the load, the containers averaged 55–61% of 16 vCPUs (peaking at 100%), 30–34 GB of their 120 GB memory with a peak of 51 GB, and never more than 75 GiB of their 200 GiB disk. Nothing was starved; the remaining headroom is in the tail — the slowest write task ran about 100 seconds against a median of 31, most likely the rows with no sale date, which all share one partition.

Query-ready, measured​

A table is only "query-ready" if queries are actually fast the moment it commits. Straight after run 1 we ran five typical queries against the new table (warm timings):

QueryTimeRows readData read
One day0.89 s3.2 M0.01 GB
One week, by day0.79 s27.8 M0.08 GB
One day, one warehouse0.79 s0.16 M< 0.01 GB
One warehouse across all 1,837 days5.46 s645 M of 12.96 B8.3 of 957 GB
30 days, grouped by warehouse1.24 s281.6 M0.89 GB

The five queries, exactly as run. The day is the latest date in the data and the warehouse is one of TPC-DS's twenty; cs_net_paid is the sale's net amount.

-- One day
SELECT count(*), sum(cs_net_paid)
FROM lake.bench.sales_clustered
WHERE sales_date = DATE '2003-01-14';

-- One week, by day
SELECT sales_date, count(*), sum(cs_net_paid)
FROM lake.bench.sales_clustered
WHERE sales_date BETWEEN DATE '2003-01-14' - INTERVAL 6 DAYS AND DATE '2003-01-14'
GROUP BY sales_date;

-- One day, one warehouse
SELECT count(*), sum(cs_net_paid)
FROM lake.bench.sales_clustered
WHERE sales_date = DATE '2003-01-14'
AND warehouse_name = 'Bad cards must make.';

-- One warehouse across all 1,837 days (no partition can be pruned)
SELECT count(*), sum(cs_net_paid)
FROM lake.bench.sales_clustered
WHERE warehouse_name = 'Bad cards must make.';

-- 30 days, grouped by warehouse
SELECT warehouse_name, count(*), sum(cs_net_paid)
FROM lake.bench.sales_clustered
WHERE sales_date BETWEEN DATE '2003-01-14' - INTERVAL 29 DAYS AND DATE '2003-01-14'
GROUP BY warehouse_name;

The layout is doing the work. Day partitions let Iceberg skip whole files, and sorting by warehouse inside each day gives every Parquet row group a tight min/max range — so a warehouse filter reads about 5% of the table even when no partition can be pruned. These queries also run on Cazpian's native query acceleration.

Honest notes​

The load itself does not use native query acceleration — on purpose. Native acceleration is what makes our queries fast, so we expected it to help here too. It did not. It cannot yet read raw files through our governed catalog path, and it declines to run the shuffle when the target is partitioned by an Iceberg transform such as days(). We forced it into the plan three different ways, and every variant was slower — the conversions between Spark's row format and the accelerator's columnar format cost more than the acceleration saved. The fastest plan for this job is plain Spark on Graviton, and that is what we measured. We are sharing the details — plans and numbers — with the upstream open-source project.

The input is nine copies of one replica. We built 1.44 billion rows with Spark and copied those 13 Parquet files nine times in S3 — 117 distinct objects, the way the original post used a fixed set of existing files rather than generating data inside the clock. Rows repeat across copies; the load still reads, shuffles, sorts and writes every one of them, and the row count is checked exactly before and after.

Warm and cold are different clocks, and we report both. The 7-minute figure is a running pool, statement start to commit. The 12-minute figure starts before the cluster exists. Serverless platforms hide the second clock inside the first; we would rather show it.

The x86 run was single. Graviton was run three times; x86 once, as the comparison.

Reproduce it​

The benchmark tooling — building the input, the timed load with its integrity checks, collecting the Spark metrics from the live driver, and the query suite — is part of our internal benchmarking lab, along with every physical plan and metric behind this post. If you want to run the same job on your own data in a Cazpian pool, the three SQL statements above are the whole thing.

Twelve minutes from nothing to 13 billion query-ready rows, for about six dollars. Seven minutes and under four dollars if the pool is already up.