From 45187f62dde143e091ae244c54c86003eb15fa0f Mon Sep 17 00:00:00 2001 From: Angel Conde Date: Sat, 18 Jul 2026 17:17:11 +0200 Subject: [PATCH 1/2] Core, Spark 3.5, 4.0, 4.1: Add clustered DV writer for ordered position delta writes BaseDVFileWriter buffers one in-memory position index per data file the task deletes from and merges each file's previous deletes only on close, so its peak resident state is proportional to the number of data files the task touches. For hot-key MERGE workloads on tables with many files, that working set grows with the table. The data path has long offered a bounded alternative for ordered input (ClusteredDataWriter vs FanoutDataWriter); the DV path had no such variant. This adds ClusteredDVFileWriter, which requires the incoming deletes to be clustered by data file: it keeps a single live position index and, on the first delete of the next file, merges the previous file's existing deletes, flushes the blob into the shared Puffin file and releases the bitmap. Peak state becomes one position index plus small per-file blob metadata, independent of how many files the task deletes from. The writer fails fast if the input revisits a completed data file, as that would produce multiple DVs for one file in a commit. Spark position delta writes request ordering by spec, partition, file and position when fanout writers are disabled, so SparkPositionDeltaWrite now selects the clustered writer when the input is ordered and keeps PartitioningDVWriter for unordered (fanout) input. --- .../deletes/ClusteredDVFileWriter.java | 241 ++++++++++++++++++ .../apache/iceberg/io/ClusteredDVWriter.java | 68 +++++ .../org/apache/iceberg/io/TestDVWriters.java | 124 +++++++++ .../spark/source/SparkPositionDeltaWrite.java | 7 + .../spark/source/SparkPositionDeltaWrite.java | 7 + .../spark/source/SparkPositionDeltaWrite.java | 7 + 6 files changed, 454 insertions(+) create mode 100644 core/src/main/java/org/apache/iceberg/deletes/ClusteredDVFileWriter.java create mode 100644 core/src/main/java/org/apache/iceberg/io/ClusteredDVWriter.java diff --git a/core/src/main/java/org/apache/iceberg/deletes/ClusteredDVFileWriter.java b/core/src/main/java/org/apache/iceberg/deletes/ClusteredDVFileWriter.java new file mode 100644 index 000000000000..d76f4ee9263b --- /dev/null +++ b/core/src/main/java/org/apache/iceberg/deletes/ClusteredDVFileWriter.java @@ -0,0 +1,241 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.deletes; + +import java.io.IOException; +import java.util.List; +import java.util.Set; +import java.util.function.Function; +import java.util.function.Supplier; +import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.FileFormat; +import org.apache.iceberg.FileMetadata; +import org.apache.iceberg.IcebergBuild; +import org.apache.iceberg.MetadataColumns; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.StructLike; +import org.apache.iceberg.io.DeleteWriteResult; +import org.apache.iceberg.io.OutputFile; +import org.apache.iceberg.io.OutputFileFactory; +import org.apache.iceberg.puffin.Blob; +import org.apache.iceberg.puffin.BlobMetadata; +import org.apache.iceberg.puffin.Puffin; +import org.apache.iceberg.puffin.PuffinWriter; +import org.apache.iceberg.puffin.StandardBlobTypes; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.util.CharSequenceSet; +import org.apache.iceberg.util.ContentFileUtil; +import org.apache.iceberg.util.StructLikeUtil; + +/** + * A {@link DVFileWriter} for inputs that are clustered by data file. + * + *

{@link BaseDVFileWriter} buffers one in-memory position index per data file the task deletes + * from and merges each file's previous deletes only on close, so its peak working set is + * proportional to the number of data files touched by the task (which, for hot-key MERGE workloads, + * grows with the table). When the incoming deletes are clustered by data file - e.g. Spark position + * delta writes request ordering by spec, partition, file and position when fanout writers are + * disabled - that buffering is unnecessary: this writer keeps a single live position index. On the + * first delete of the next data file, the previous file's index is merged with that file's existing + * deletes, flushed as a blob into the shared Puffin file and released, so the resident state is one + * position index plus small per-file blob metadata needed to create the {@link DeleteFile} entries + * on close. + * + *

The writer fails fast if the input revisits a data file after moving on to another one, as + * that would produce multiple DVs for one data file in a single commit. + */ +public class ClusteredDVFileWriter implements DVFileWriter { + + private static final String REFERENCED_DATA_FILE_KEY = "referenced-data-file"; + private static final String CARDINALITY_KEY = "cardinality"; + + private final Supplier dvOutputFile; + private final Function loadPreviousDeletes; + private final List flushedDeletes = Lists.newArrayList(); + private final Set seenPaths = Sets.newHashSet(); + private final CharSequenceSet referencedDataFiles = CharSequenceSet.empty(); + private final List rewrittenDeleteFiles = Lists.newArrayList(); + + private PuffinWriter writer = null; + private String currentPath = null; + private PositionDeleteIndex currentPositions = null; + private PartitionSpec currentSpec = null; + private StructLike currentPartition = null; + private DeleteWriteResult result = null; + + public ClusteredDVFileWriter( + OutputFileFactory fileFactory, Function loadPreviousDeletes) { + this(() -> fileFactory.newOutputFile().encryptingOutputFile(), loadPreviousDeletes); + } + + public ClusteredDVFileWriter( + Supplier dvOutputFile, + Function loadPreviousDeletes) { + this.dvOutputFile = dvOutputFile; + this.loadPreviousDeletes = loadPreviousDeletes; + } + + @Override + public void delete(String path, long pos, PartitionSpec spec, StructLike partition) { + if (!path.equals(currentPath)) { + startNewFile(path, spec, partition); + } + + currentPositions.delete(pos); + } + + @Override + public DeleteWriteResult result() { + Preconditions.checkState(result != null, "Cannot get result from unclosed writer"); + return result; + } + + @Override + public void close() throws IOException { + if (result == null) { + flushCurrentFile(); + + List dvs = Lists.newArrayList(); + + if (writer != null) { + writer.close(); + + // DVs share the Puffin path and file size but have different offsets + String puffinPath = writer.location(); + long puffinFileSize = writer.fileSize(); + + for (FlushedDeletes deletes : flushedDeletes) { + dvs.add(createDV(puffinPath, puffinFileSize, deletes)); + } + } + + this.result = new DeleteWriteResult(dvs, referencedDataFiles, rewrittenDeleteFiles); + } + } + + private void startNewFile(String path, PartitionSpec spec, StructLike partition) { + Preconditions.checkState( + seenPaths.add(path), + "Incoming deletes are not clustered by data file: %s was already completed", + path); + + flushCurrentFile(); + + this.currentPath = path; + this.currentPositions = new BitmapPositionDeleteIndex(); + this.currentSpec = spec; + this.currentPartition = StructLikeUtil.copy(partition); + } + + private void flushCurrentFile() { + if (currentPath == null) { + return; + } + + PositionDeleteIndex previousPositions = loadPreviousDeletes.apply(currentPath); + if (previousPositions != null) { + currentPositions.merge(previousPositions); + for (DeleteFile previousDeleteFile : previousPositions.deleteFiles()) { + // only DVs and file-scoped deletes can be discarded from the table state + if (ContentFileUtil.isFileScoped(previousDeleteFile)) { + rewrittenDeleteFiles.add(previousDeleteFile); + } + } + } + + if (writer == null) { + this.writer = newWriter(); + } + + BlobMetadata blobMetadata = writer.write(toBlob(currentPositions, currentPath)); + flushedDeletes.add( + new FlushedDeletes( + currentPath, + currentSpec, + currentPartition, + blobMetadata, + currentPositions.cardinality())); + referencedDataFiles.add(currentPath); + + // release the bitmap - this is the whole point of the clustered writer + this.currentPath = null; + this.currentPositions = null; + this.currentSpec = null; + this.currentPartition = null; + } + + private DeleteFile createDV(String path, long size, FlushedDeletes deletes) { + return FileMetadata.deleteFileBuilder(deletes.spec) + .ofPositionDeletes() + .withFormat(FileFormat.PUFFIN) + .withPath(path) + .withPartition(deletes.partition) + .withFileSizeInBytes(size) + .withReferencedDataFile(deletes.path) + .withContentOffset(deletes.blobMetadata.offset()) + .withContentSizeInBytes(deletes.blobMetadata.length()) + .withRecordCount(deletes.cardinality) + .build(); + } + + private PuffinWriter newWriter() { + OutputFile outputFile = dvOutputFile.get(); + return Puffin.write(outputFile).createdBy(IcebergBuild.fullVersion()).build(); + } + + private Blob toBlob(PositionDeleteIndex positions, String path) { + return new Blob( + StandardBlobTypes.DV_V1, + ImmutableList.of(MetadataColumns.ROW_POSITION.fieldId()), + -1 /* snapshot ID is inherited */, + -1 /* sequence number is inherited */, + positions.serialize(), + null /* uncompressed */, + ImmutableMap.of( + REFERENCED_DATA_FILE_KEY, + path, + CARDINALITY_KEY, + String.valueOf(positions.cardinality()))); + } + + private static class FlushedDeletes { + private final String path; + private final PartitionSpec spec; + private final StructLike partition; + private final BlobMetadata blobMetadata; + private final long cardinality; + + private FlushedDeletes( + String path, + PartitionSpec spec, + StructLike partition, + BlobMetadata blobMetadata, + long cardinality) { + this.path = path; + this.spec = spec; + this.partition = partition; + this.blobMetadata = blobMetadata; + this.cardinality = cardinality; + } + } +} diff --git a/core/src/main/java/org/apache/iceberg/io/ClusteredDVWriter.java b/core/src/main/java/org/apache/iceberg/io/ClusteredDVWriter.java new file mode 100644 index 000000000000..d01f9541de7b --- /dev/null +++ b/core/src/main/java/org/apache/iceberg/io/ClusteredDVWriter.java @@ -0,0 +1,68 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.io; + +import java.io.IOException; +import java.util.function.Function; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.StructLike; +import org.apache.iceberg.deletes.ClusteredDVFileWriter; +import org.apache.iceberg.deletes.DVFileWriter; +import org.apache.iceberg.deletes.PositionDelete; +import org.apache.iceberg.deletes.PositionDeleteIndex; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; + +/** + * A {@link PartitioningWriter} for deletion vectors that requires the incoming deletes to be + * clustered by data file (e.g. sorted by spec, partition, file and position). Unlike {@link + * PartitioningDVWriter}, which buffers one position index per touched data file until close, this + * writer keeps a single live position index and flushes each data file's deletion vector as soon as + * the input moves on to the next file, bounding the task's resident state to one data file. + */ +public class ClusteredDVWriter + implements PartitioningWriter, DeleteWriteResult> { + + private final DVFileWriter writer; + private DeleteWriteResult result; + + public ClusteredDVWriter( + OutputFileFactory fileFactory, + Function loadPreviousDeletes) { + this.writer = new ClusteredDVFileWriter(fileFactory, loadPreviousDeletes::apply); + } + + @Override + public void write(PositionDelete row, PartitionSpec spec, StructLike partition) { + writer.delete(row.path().toString(), row.pos(), spec, partition); + } + + @Override + public DeleteWriteResult result() { + Preconditions.checkState(result != null, "Cannot get result from unclosed writer"); + return result; + } + + @Override + public void close() throws IOException { + if (result == null) { + writer.close(); + this.result = writer.result(); + } + } +} diff --git a/data/src/test/java/org/apache/iceberg/io/TestDVWriters.java b/data/src/test/java/org/apache/iceberg/io/TestDVWriters.java index b89f04a3c0e8..656f3d16e210 100644 --- a/data/src/test/java/org/apache/iceberg/io/TestDVWriters.java +++ b/data/src/test/java/org/apache/iceberg/io/TestDVWriters.java @@ -19,6 +19,7 @@ package org.apache.iceberg.io; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assumptions.assumeThat; import java.io.File; @@ -45,6 +46,7 @@ import org.apache.iceberg.data.BaseDeleteLoader; import org.apache.iceberg.data.DeleteLoader; import org.apache.iceberg.deletes.BaseDVFileWriter; +import org.apache.iceberg.deletes.ClusteredDVFileWriter; import org.apache.iceberg.deletes.DVFileWriter; import org.apache.iceberg.deletes.Deletes; import org.apache.iceberg.deletes.PositionDelete; @@ -170,6 +172,128 @@ public void testBasicDVs() throws IOException { assertRows(ImmutableList.of(toRow(11, "aaa"), toRow(12, "aaa"))); } + @TestTemplate + public void testClusteredDVsBasic() throws IOException { + assumeThat(formatVersion).isGreaterThanOrEqualTo(3); + + FileWriterFactory writerFactory = newWriterFactory(table.schema()); + + // add the first data file + List rows1 = ImmutableList.of(toRow(1, "aaa"), toRow(2, "aaa"), toRow(11, "aaa")); + DataFile dataFile1 = writeData(writerFactory, fileFactory, rows1, table.spec(), null); + table.newFastAppend().appendFile(dataFile1).commit(); + + // add the second data file + List rows2 = ImmutableList.of(toRow(3, "aaa"), toRow(4, "aaa"), toRow(12, "aaa")); + DataFile dataFile2 = writeData(writerFactory, fileFactory, rows2, table.spec(), null); + table.newFastAppend().appendFile(dataFile2).commit(); + + // init the clustered DV writer + DVFileWriter dvWriter = + new ClusteredDVFileWriter(fileFactory, new PreviousDeleteLoader(table, ImmutableMap.of())); + + // write deletes clustered by data file + dvWriter.delete(dataFile1.location(), 0L, table.spec(), null); + dvWriter.delete(dataFile1.location(), 1L, table.spec(), null); + dvWriter.delete(dataFile2.location(), 0L, table.spec(), null); + dvWriter.delete(dataFile2.location(), 1L, table.spec(), null); + dvWriter.close(); + + // verify the writer result + DeleteWriteResult result = dvWriter.result(); + assertThat(result.deleteFiles()).hasSize(2); + assertThat(result.referencedDataFiles()) + .hasSize(2) + .contains(dataFile1.location()) + .contains(dataFile2.location()); + assertThat(result.referencesDataFiles()).isTrue(); + + // commit the deletes + commit(result); + + // verify correctness + assertRows(ImmutableList.of(toRow(11, "aaa"), toRow(12, "aaa"))); + } + + @TestTemplate + public void testClusteredDVsMergePreviousDeletes() throws IOException { + assumeThat(formatVersion).isGreaterThanOrEqualTo(3); + + FileWriterFactory writerFactory = newWriterFactory(table.schema()); + + // add a data file with 3 data records + List rows = ImmutableList.of(toRow(1, "aaa"), toRow(2, "aaa"), toRow(3, "aaa")); + DataFile dataFile = writeData(writerFactory, parquetFileFactory, rows, table.spec(), null); + table.newFastAppend().appendFile(dataFile).commit(); + + // write the first DV + DVFileWriter dvWriter1 = + new ClusteredDVFileWriter(fileFactory, new PreviousDeleteLoader(table, ImmutableMap.of())); + dvWriter1.delete(dataFile.location(), 1L, table.spec(), null); + dvWriter1.close(); + + // validate the writer result + DeleteWriteResult result1 = dvWriter1.result(); + assertThat(result1.deleteFiles()).hasSize(1); + assertThat(result1.referencedDataFiles()).containsOnly(dataFile.location()); + assertThat(result1.rewrittenDeleteFiles()).isEmpty(); + + // commit the first DV + commit(result1); + + // write the second DV, merging with the first one + DeleteFile dv1 = Iterables.getOnlyElement(result1.deleteFiles()); + DVFileWriter dvWriter2 = + new ClusteredDVFileWriter( + fileFactory, + new PreviousDeleteLoader(table, ImmutableMap.of(dataFile.location(), dv1))); + dvWriter2.delete(dataFile.location(), 2L, table.spec(), null); + dvWriter2.close(); + + // validate the previous DV is merged and marked as rewritten + DeleteWriteResult result2 = dvWriter2.result(); + assertThat(result2.deleteFiles()).hasSize(1); + assertThat(result2.referencedDataFiles()).containsOnly(dataFile.location()); + assertThat(result2.rewrittenDeleteFiles()).hasSize(1); + + // replace DVs + commit(result2); + SnapshotChanges changes = SnapshotChanges.builderFor(table).build(); + assertThat(changes.addedDeleteFiles()).hasSize(1); + assertThat(changes.removedDeleteFiles()).hasSize(1); + + // verify correctness after replacing DVs + assertRows(ImmutableList.of(toRow(1, "aaa"))); + } + + @TestTemplate + public void testClusteredDVsRejectUnclusteredInput() throws IOException { + assumeThat(formatVersion).isGreaterThanOrEqualTo(3); + + FileWriterFactory writerFactory = newWriterFactory(table.schema()); + + // add two data files + List rows1 = ImmutableList.of(toRow(1, "aaa"), toRow(2, "aaa")); + DataFile dataFile1 = writeData(writerFactory, fileFactory, rows1, table.spec(), null); + table.newFastAppend().appendFile(dataFile1).commit(); + + List rows2 = ImmutableList.of(toRow(3, "aaa"), toRow(4, "aaa")); + DataFile dataFile2 = writeData(writerFactory, fileFactory, rows2, table.spec(), null); + table.newFastAppend().appendFile(dataFile2).commit(); + + DVFileWriter dvWriter = + new ClusteredDVFileWriter(fileFactory, new PreviousDeleteLoader(table, ImmutableMap.of())); + + // revisiting a completed data file must fail, as it would produce two DVs for one file + dvWriter.delete(dataFile1.location(), 0L, table.spec(), null); + dvWriter.delete(dataFile2.location(), 0L, table.spec(), null); + assertThatThrownBy(() -> dvWriter.delete(dataFile1.location(), 1L, table.spec(), null)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("not clustered by data file"); + + dvWriter.close(); + } + @TestTemplate public void testRewriteDVs() throws IOException { assumeThat(formatVersion).isGreaterThanOrEqualTo(3); diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java index f926bd96389a..26f7d2033ece 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java @@ -54,6 +54,7 @@ import org.apache.iceberg.expressions.Expression; import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.io.BasePositionDeltaWriter; +import org.apache.iceberg.io.ClusteredDVWriter; import org.apache.iceberg.io.ClusteredDataWriter; import org.apache.iceberg.io.ClusteredPositionDeleteWriter; import org.apache.iceberg.io.DataWriteResult; @@ -509,6 +510,12 @@ protected PartitioningWriter, DeleteWriteResult> new DeleteGranularity deleteGranularity = context.deleteGranularity(); if (context.useDVs()) { + if (inputOrdered) { + // clustered input allows flushing each data file's DV as soon as the input moves on, + // bounding the writer's resident state to one data file instead of all touched files + return new ClusteredDVWriter<>(files, previousDeleteLoader); + } + return new PartitioningDVWriter<>(files, previousDeleteLoader); } else if (inputOrdered && rewritableDeletes == null) { return new ClusteredPositionDeleteWriter<>( diff --git a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java index f0a58fc42107..c7b1e7ec19e3 100644 --- a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java +++ b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java @@ -55,6 +55,7 @@ import org.apache.iceberg.expressions.Expression; import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.io.BasePositionDeltaWriter; +import org.apache.iceberg.io.ClusteredDVWriter; import org.apache.iceberg.io.ClusteredDataWriter; import org.apache.iceberg.io.ClusteredPositionDeleteWriter; import org.apache.iceberg.io.DataWriteResult; @@ -544,6 +545,12 @@ protected PartitioningWriter, DeleteWriteResult> new DeleteGranularity deleteGranularity = context.deleteGranularity(); if (context.useDVs()) { + if (inputOrdered) { + // clustered input allows flushing each data file's DV as soon as the input moves on, + // bounding the writer's resident state to one data file instead of all touched files + return new ClusteredDVWriter<>(files, previousDeleteLoader); + } + return new PartitioningDVWriter<>(files, previousDeleteLoader); } else if (inputOrdered && rewritableDeletes == null) { return new ClusteredPositionDeleteWriter<>( diff --git a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java index 89f6e1c350fc..c666726e79b3 100644 --- a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java +++ b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkPositionDeltaWrite.java @@ -54,6 +54,7 @@ import org.apache.iceberg.exceptions.CleanableFailure; import org.apache.iceberg.expressions.Expression; import org.apache.iceberg.io.BasePositionDeltaWriter; +import org.apache.iceberg.io.ClusteredDVWriter; import org.apache.iceberg.io.ClusteredDataWriter; import org.apache.iceberg.io.ClusteredPositionDeleteWriter; import org.apache.iceberg.io.DataWriteResult; @@ -547,6 +548,12 @@ protected PartitioningWriter, DeleteWriteResult> new DeleteGranularity deleteGranularity = context.deleteGranularity(); if (context.useDVs()) { + if (inputOrdered) { + // clustered input allows flushing each data file's DV as soon as the input moves on, + // bounding the writer's resident state to one data file instead of all touched files + return new ClusteredDVWriter<>(files, previousDeleteLoader); + } + return new PartitioningDVWriter<>(files, previousDeleteLoader); } else if (inputOrdered && rewritableDeletes == null) { return new ClusteredPositionDeleteWriter<>( From 671c5b2795b69a376afb6effd5fdc329e67f181e Mon Sep 17 00:00:00 2001 From: Angel Conde Date: Sun, 19 Jul 2026 18:20:18 +0200 Subject: [PATCH 2/2] Core: Add DV writer buffered-state memory benchmark Measures retained heap just before close while writing deletes for N data files: BaseDVFileWriter holds one position index per file, the clustered writer only per-file blob metadata. At 2000 positions/file the buffered state is 258.5MB vs 27.9MB for 50k files (~9-10x across scales). --- .../deletes/TestDVWriterMemoryBenchmark.java | 128 ++++++++++++++++++ 1 file changed, 128 insertions(+) create mode 100644 core/src/test/java/org/apache/iceberg/deletes/TestDVWriterMemoryBenchmark.java diff --git a/core/src/test/java/org/apache/iceberg/deletes/TestDVWriterMemoryBenchmark.java b/core/src/test/java/org/apache/iceberg/deletes/TestDVWriterMemoryBenchmark.java new file mode 100644 index 000000000000..acbbca9142d4 --- /dev/null +++ b/core/src/test/java/org/apache/iceberg/deletes/TestDVWriterMemoryBenchmark.java @@ -0,0 +1,128 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.deletes; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.io.File; +import java.io.IOException; +import java.nio.file.Path; +import java.util.UUID; +import java.util.function.Function; +import java.util.function.Supplier; +import org.apache.iceberg.Files; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.io.OutputFile; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +/** + * Micro-benchmark comparing the retained heap of {@link BaseDVFileWriter} (one buffered position + * index per touched data file until close) against {@link ClusteredDVFileWriter} (one live position + * index, flushed per file) while writing deletes for N data files. + * + *

The measurement is taken just before {@code close()}: at that point the base writer holds + * every file's bitmap while the clustered writer has already flushed all but the last one. Heap is + * sampled after forced GC, so only strongly-reachable writer state is counted. + * + *

Results are printed and written to {@code build/dv-writer-memory-benchmark.txt}. + */ +public class TestDVWriterMemoryBenchmark { + + private static final PartitionSpec SPEC = PartitionSpec.unpartitioned(); + private static final int POSITIONS_PER_FILE = 2000; + private static final long POSITION_STRIDE = 37; + + @TempDir private Path temp; + + @Test + public void benchmarkBufferedStateBeforeClose() throws IOException { + StringBuilder report = new StringBuilder(); + report.append( + String.format( + "positionsPerFile=%d stride=%d%n%12s %14s %18s %8s%n", + POSITIONS_PER_FILE, + POSITION_STRIDE, + "files", + "stockPeakMB", + "clusteredPeakMB", + "ratio")); + + for (int files : new int[] {1_000, 10_000, 50_000}) { + long stockPeak = measurePeakHeap(files, out -> new BaseDVFileWriter(out, path -> null)); + long clusteredPeak = + measurePeakHeap(files, out -> new ClusteredDVFileWriter(out, path -> null)); + + report.append( + String.format( + "%12d %14.1f %18.1f %7.1fx%n", + files, + stockPeak / 1048576.0, + clusteredPeak / 1048576.0, + stockPeak / (double) Math.max(clusteredPeak, 1))); + + // the architectural claim: buffered state before close is bounded for the clustered + // writer and proportional to the touched file count for the base writer + assertThat(clusteredPeak).isLessThan(stockPeak); + } + + System.out.println(report); + java.nio.file.Files.write( + java.nio.file.Paths.get("build", "dv-writer-memory-benchmark.txt"), + report.toString().getBytes(java.nio.charset.StandardCharsets.UTF_8)); + } + + private long measurePeakHeap(int numFiles, Function, DVFileWriter> factory) + throws IOException { + Supplier outputFile = + () -> Files.localOutput(new File(temp.toFile(), "dv-" + UUID.randomUUID() + ".puffin")); + + long baseline = usedHeapAfterGc(); + + DVFileWriter writer = factory.apply(outputFile); + for (int file = 0; file < numFiles; file += 1) { + String path = "s3://bucket/warehouse/db/table/data/partition=0/file-" + file + ".parquet"; + for (int pos = 0; pos < POSITIONS_PER_FILE; pos += 1) { + writer.delete(path, pos * POSITION_STRIDE, SPEC, null); + } + } + + // peak buffered state: all deletes written, nothing closed yet + long peak = usedHeapAfterGc() - baseline; + + writer.close(); + assertThat(writer.result().deleteFiles()).hasSize(numFiles); + + return peak; + } + + private static long usedHeapAfterGc() { + for (int i = 0; i < 3; i += 1) { + System.gc(); + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + Runtime runtime = Runtime.getRuntime(); + return runtime.totalMemory() - runtime.freeMemory(); + } +}