13 Billion Rows into Apache Iceberg in 7 Minutes, for $3.69
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:
| Input | 12,959,823,744 rows · 1.13 TB zstd Parquet · 117 files of ~9.7 GB · 38 columns |
| Output | Apache 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 |
| Statement | one INSERT INTO … SELECT * — one atomic Iceberg commit |
| Compute | a Cazpian compute pool: 32 workers × 16 vCPU / 120 GB on AWS Graviton (arm64), Spark 4.1, Iceberg 1.11 |
| Clock | INSERT 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
| Run | Timed load | Throughput | Rows committed |
|---|---|---|---|
| 1 — first job on a freshly started pool | 431.7 s | 30.0 M rows/s | 12,959,823,744 ✅ |
| 2 | 424.5 s | 30.5 M rows/s | 12,959,823,744 ✅ |
| 3 | 423.0 s | 30.6 M rows/s | 12,959,823,744 ✅ |
| Median | 7 min 4.5 s | 30.5 M rows/s · 2.66 GB/s | exact 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 Graviton | Difference | |
|---|---|---|---|
| Rows written | 13.3 billion | 12.96 billion | same scale (97%) |
| Input read | 1.04 TB, as reported¹ | 1.13 TB, measured | ours is larger |
| Rows read to write them | 40.1 billion (3 passes) | 12.96 billion (1 pass) | 3.1× less work |
| Load time, clustered table | 14 min 5 s | 7 min 4 s | 2.0× faster |
| From zero, cluster start-up included | — (serverless) | 12 min 3 s | still faster |
| Write throughput | 15.7 M rows/s | 30.5 M rows/s | 1.9× |
| Read throughput | over 1.2 GB/s | 2.66 GB/s | 2.2× |
| Sub-10-minute target | missed | met, with clustering | |
| Compute cost per load | ~70 DBUs² ≈ $25–65 | ≈ $3.69 | ~7–18× cheaper |
| Tuning needed | 32 MB file splits to trigger autoscaling | none — 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:
| Phase | Graviton | x86 (same job) |
|---|---|---|
| Read 1.13 TB of Parquet, shuffle 2.18 TB by day | 3 min 29 s | 3 min 38 s |
| Read the shuffle, sort within each day, write 957 GB of Iceberg files | 3 min 32 s | 5 min 37 s |
| Iceberg commit (one snapshot, 8,331 files) | ~1 s | ~1 s |
| Total | 424.5 s | 559.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):
| Query | Time | Rows read | Data read |
|---|---|---|---|
| One day | 0.89 s | 3.2 M | 0.01 GB |
| One week, by day | 0.79 s | 27.8 M | 0.08 GB |
| One day, one warehouse | 0.79 s | 0.16 M | < 0.01 GB |
| One warehouse across all 1,837 days | 5.46 s | 645 M of 12.96 B | 8.3 of 957 GB |
| 30 days, grouped by warehouse | 1.24 s | 281.6 M | 0.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.