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).

TL;DR. On the pure merge-and-aggregate probes, VecRuntime runs 1.0–1.4× faster than Spark across every delete shape and share, and Q1 (a full scan-aggregate) 3.0–4.6× faster. The margin does not grow with the delete share — it shrinks, because the delete merge costs a fixed price per batch (paid inside Iceberg's reader, in both engines) and a fixed price hurts the faster engine's ratio more. Deletion vectors (v3) are the cheaper encoding for both engines: the pure-merge probe costs Spark +15 ms and VecRuntime +17–22 ms over a clean table, against +21–40 ms for v2 position files. Equality deletes (v2) are the expensive kind and the one shape VecRuntime can lose (0.89× on 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

Scope. This is the merge-on-read merge cost at TPC-H scale factor 1 on a single host, over Iceberg's JVM vectorized reader (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.
ComponentConfiguration
Datasetlineitem, 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 formatApache 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
HostOne 8-vCPU x86 EC2 instance, local[8], 8 GB heap, quiet machine; Corretto 25
Protocol5 warm-up and 10 measured iterations per cell; median wall-clock in milliseconds; spark vs vector; all cells returned identical checksums to Spark
Queriesprobe-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.

ComponentConfiguration
ClusterEKS, 8 × m5.4xlarge executors (16 vCPU, 64 GiB, 300 GB disk each); Spark 4.1.3, Iceberg 1.11.0; tables on S3
TableA 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_520% of rows deleted by position (parquet position-delete files) + 5% by equality (parquet equality-delete files)
v3 dvmix_20_5The same, with the position deletes as deletion vectors (Puffin) and the equality deletes unchanged
Change batch10% of rows: 42.0M updates, 11.0M deletes, 4.6M inserts (3.6 GiB parquet), touching 64 of 64 data files
ProtocolReads: 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
Configsspark (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.