Iceberg merge-on-read: VecRuntime vs Apache Spark
A merge-on-read Iceberg table carries delete files between compactions — the record of the UPDATEs, DELETEs and MERGEs applied since the data files were written. A reader must merge those deletes against the data on every scan. This page measures that merge on the JVM: VecRuntime — which runs the scan's filters, aggregates and joins on Arrow batches over Iceberg's own vectorized reader, turning the deletes into a selection — against plain Apache Spark, over both the v2 delete encodings (positional and equality delete files, #261) and the v3 deletion vectors (one roaring bitmap per data file in a Puffin blob, #262).
eq_10): Iceberg
evaluates the equality predicate row-at-a-time to build the mapping. Every checksum equalled Spark's, on all 90 cells.
CDC MERGE on the cluster (~20 GB store_sales, 432M live rows, 20% position + 5% equality deletes, a 10% change batch touching every data file): VecRuntime is not faster yet. The MERGE runs at 0.96× (v2) and 0.92× (v3) of Spark and the reads at 0.89–1.06×, although 7 of the MERGE's 8 operators and every read operator run columnar, through the columnar shuffle. Every v3 checksum equalled Spark's.
What was measured
BatchScanExec), which is the reader both engines use for v2 and v3. It is not the 1 TB
cluster campaign; the Iceberg legs of that campaign are prepared (they read S3 through S3FileIO with the Analytics
Accelerator stream) but not yet run. Comet was not measured here — no Comet jar on this host — so this comparison is
VecRuntime against OSS Spark only, as requested.| Component | Configuration |
|---|---|
| Dataset | lineitem, TPC-H SF1 (6,001,215 rows), decimals-as-doubles, generated by the harness (gen-iceberg-mor.sh); the delete shapes a lakehouse table carries between compactions |
| Table format | Apache Iceberg, merge-on-read; v2 positional and equality delete files (#261), v3 deletion vectors in Puffin (#262); read through Iceberg's JVM vectorized reader under both engines |
| Host | One 8-vCPU x86 EC2 instance, local[8], 8 GB heap, quiet machine; Corretto 25 |
| Protocol | 5 warm-up and 10 measured iterations per cell; median wall-clock in milliseconds; spark vs vector; all cells returned identical checksums to Spark |
| Queries | probe-sum / probe-group (a merge then a sum / a grouped sum — the merge in isolation), probe-count (the merge alone), TPC-H Q6 (a 2%-selective filter+sum), Q1 (a full scan, grouped aggregate) |
The delete variants: plain (no deletes) as the baseline, pos_N / dv_N (N% of rows deleted by position, v2 files or v3 vectors),
*_clustered (the deletes in whole 64-row blocks rather than scattered), *_upd_N (an UPDATE/MERGE shape: several snapshots and small files),
and eq_N (v2 equality deletes — a predicate, not positions). "live %" is the share of rows that survive the deletes.
v3 deletion vectors
The current encoding: one roaring bitmap per data file, deserialised once. Median ms; green = VecRuntime faster than Spark, the ratio in the last column.
v2 delete files
The older encoding: positional delete files (a sorted list of deleted row positions, indexed per task) and equality delete files (a predicate evaluated per row). Median ms.
What the numbers say
The margin is flat in the delete share, not growing
The hypothesis going in was that the more rows are deleted, the more VecRuntime's in-place kernels over a bitmap would
pull ahead of Spark's per-row mapping[i] indirection. The data refutes it: on probe-sum VecRuntime is
1.38× faster on a clean table and 1.0–1.2× on every positional variant whatever the percentage; on probe-group
1.61× becomes 1.2–1.4×; on Q1 4.36× becomes 3.1–4.0×. The reason is in the per-batch cost: the deletes cost a fixed price
per batch that is flat in the delete share — Spark pays 20–25 ms on probe-sum at 2%, 10% and 30% alike,
VecRuntime pays 30–50 ms — and a fixed price hurts the faster engine's ratio more. At SF1 the per-row indirection term the
hypothesis was about is invisible against the per-batch work.
Where the fixed price goes — Iceberg's reader, in both engines
Profiling puts Iceberg at the top of every list: ColumnarBatchUtil.buildRowIdMapping (the int[] the
reader builds per batch), and for v2 Deletes.toPositionIndexes plus the path hashing under it (a third of the Iceberg
samples). On v3 that index construction is gone — RoaringPositionBitmap.contains takes its place at a fraction
of the cost (one blob per file against thousands of delete rows to sort). VecRuntime's own delete-to-selection conversion
(selectionFromIndices) is ~1.5% of samples and the validity copy absent, so the merge cost is Iceberg's, not the
plugin's. A direct bitmap path on the plugin's side would remove that 1.5%, not the 19% the mapping construction costs — that
is an Iceberg-reader change, worth proposing upstream with these numbers.
Deletion vectors are the cheaper encoding
Against the v2 cells of the same mutation: the pure-merge probe-count costs Spark 68–79 ms on dv_* against
99–101 on pos_*, and VecRuntime 77–83 against 93–96. The delete price over a clean table is +15 ms (Spark) and
+17–22 ms (VecRuntime) on v3, against +21–23 and +32–40 on v2 position files. What disappears is the per-task index
construction; the Puffin read itself does not register (once per file). The update/merge shape (dv_upd_*) — the one
v3 is designed for — pays the same as its v2 twin within noise, because its extra cost is the small files, not the deletes.
Equality deletes are the expensive kind
eq_10 costs Spark 3× on probe-sum (104 → 330 ms) and VecRuntime 5× (76 → 372 ms) — the one variant where
VecRuntime loses the probe (0.89×). Iceberg's reader evaluates the equality predicate per row on a ColumnarBatchRow to
build the mapping, which is row-at-a-time work in front of a columnar reader and dominates everything. The lever that would move
it — taking the position-filtered batch and evaluating the equality-delete set as VecRuntime's own IN/anti-join over
the batch — needs the reader to expose the un-applied equality deletes, which the Iceberg 1.11 API does not.
Q6, the one shape VecRuntime trails Spark
Q6 keeps 2% of the rows. On a clean table the filter compacts once and the aggregate sees a tiny batch (VecRuntime even or ahead); on a deleted variant the filter's predicate runs over a batch that already carries a delete selection, and the compaction is of a selection-over-selection. Spark's whole-stage codegen fuses the delete mapping and the predicate into one row loop. This is the one shape where VecRuntime's per-row work over the physical rows shows — at most ~40 ms per query at this scale.
CDC MERGE on the cluster: ~20 GB store_sales, mixed deletes
The SF1 probes above isolate the delete merge on one host. This section is the workload a CDC pipeline actually runs: a
MERGE INTO of a change batch into a merge-on-read table that already carries both kinds of v2 delete (or deletion
vectors on v3), on a cluster, with the reads a consumer would run before and after it.
| Component | Configuration |
|---|---|
| Cluster | EKS, 8 × m5.4xlarge executors (16 vCPU, 64 GiB, 300 GB disk each); Spark 4.1.3, Iceberg 1.11.0; tables on S3 |
| Table | A 20% sample of TPC-DS SF1000 store_sales: 64 data files, 28.8 GiB, 432.0M live rows; 616 delete files holding 126.2M deleted rows |
v2 mix_20_5 | 20% of rows deleted by position (parquet position-delete files) + 5% by equality (parquet equality-delete files) |
v3 dvmix_20_5 | The same, with the position deletes as deletion vectors (Puffin) and the equality deletes unchanged |
| Change batch | 10% of rows: 42.0M updates, 11.0M deletes, 4.6M inserts (3.6 GiB parquet), touching 64 of 64 data files |
| Protocol | Reads: 2 warm-ups + 5 timed runs each. MERGE: 1 warm-up + 3 timed runs, each rolled back to the base snapshot, so every run starts from the same table. Median wall-clock; p90 in the JSON |
| Configs | spark (OSS Spark, no plugin) vs vector-shuffle (VecRuntime with its columnar Arrow shuffle); image built from the branch carrying #476, #477 and #478 |
v2: position + equality delete files
v3: deletion vectors + equality delete files
What the cluster numbers say
Parity on the reads, a small loss on the MERGE. Every read is within ±11% of Spark and the medians straddle 1.0×, which is
the run-to-run spread here. The MERGE is consistently slower: all three of VecRuntime's v3 runs (174.7–176.3 s) are slower than
Spark's slowest (163.9 s). The plans are columnar where they can be — the scan, filters and projections, the columnar shuffle and a
columnar shuffled hash join — with one row operator left: Iceberg's WriteDelta, which VecRuntime feeds through a
ColumnarToRow.
Where the time goes is not yet profiled. The working hypothesis is that each MERGE is dominated by work both engines share:
the ~70 s merge-on-read scan (every read costs about that, whatever runs above it), and Iceberg's row writer, which here writes
46.6M rows and 64 position-delete sets of 168M positions. The ColumnarToRow in front of that writer is extra work
VecRuntime pays and Spark does not, which would fit a loss of this size. A columnar delete writer (#20) and a faster
merge-on-read scan are the two changes that would move this benchmark; a profile of the MERGE will split the time before either is
built.
Deletion vectors barely change the MERGE. Spark's MERGE is 165.0 s on v2 and 161.7 s on v3, VecRuntime's 172.2 s and 175.4 s, and the reads are within noise of v2. The v2 MERGE rewrites every position-delete file; the v3 MERGE writes one deletion vector per data file instead. That rewrite is therefore not what the MERGE time is made of.
Correctness. On v3 all seven checksums — the three reads, the merged table, and the three reads after the merge — equalled Spark's, and both engines committed the same 46,641,979 added records and 64 deletion vectors (VecRuntime in 8 data files, Spark in 14). On v2 both engines left 425,665,996 rows after the merge; Spark's v2 checksums were not kept (the harness lost that job's result file), so v2 is checked by row count only.
Numbers from the #260 harness (#261 v2, #262 v3): TPC-H SF1 on one 8-vCPU host, median of 10 measured iterations, local[8]. Not a
cluster benchmark and not comparable to the 1 TB TPC-DS results. Checksums equalled Spark's on all 90 cells; every dv_*
checksum equalled its pos_* twin (same rows, different encoding). Two of sixteen JVMs showed a transient JIT-deopt spike
on Q1's first post-warm-up iterations that recovered within the same JVM — a compilation storm, not a delete-layout cost.