From 4a093d2457a7d752b3dd774f2146e874940c668a Mon Sep 17 00:00:00 2001 From: Hongyue Zhang Date: Thu, 25 Jul 2024 22:50:25 -0700 Subject: [PATCH 1/5] Core, Spark: Remove dangling delete files Optionally can be enabled as part of RewriteDataFilesSparkAction Co-authored-by: Szehon Ho --- .../iceberg/actions/ActionsProvider.java | 6 + .../actions/RemoveDanglingDeleteFiles.java | 30 ++ .../iceberg/actions/RewriteDataFiles.java | 17 + .../iceberg/actions/BaseRewriteDataFiles.java | 6 + ...RemoveDanglingDeleteFilesActionResult.java | 36 ++ .../iceberg/spark/SparkContentFile.java | 7 +- .../RemoveDanglingDeletesSparkAction.java | 175 +++++++ .../actions/RewriteDataFilesSparkAction.java | 34 +- .../iceberg/spark/actions/SparkActions.java | 6 + .../TestRemoveDanglingDeleteAction.java | 443 ++++++++++++++++++ .../actions/TestRewriteDataFilesAction.java | 215 ++++++++- .../TestRewritePositionDeleteFilesAction.java | 3 +- 12 files changed, 966 insertions(+), 12 deletions(-) create mode 100644 api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java create mode 100644 core/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFilesActionResult.java create mode 100644 spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java create mode 100644 spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java diff --git a/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java b/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java index 2d6ff2679a17..6e744fa9c7e1 100644 --- a/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java +++ b/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java @@ -70,4 +70,10 @@ default RewritePositionDeleteFiles rewritePositionDeletes(Table table) { throw new UnsupportedOperationException( this.getClass().getName() + " does not implement rewritePositionDeletes"); } + + /** Instantiates an action to remove dangling delete files from current snapshot. */ + default RemoveDanglingDeleteFiles removeDanglingDeleteFiles(Table table) { + throw new UnsupportedOperationException( + this.getClass().getName() + " does not implement removeDanglingDeleteFiles"); + } } diff --git a/api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java b/api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java new file mode 100644 index 000000000000..879dc9f0c196 --- /dev/null +++ b/api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java @@ -0,0 +1,30 @@ +/* + * 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.actions; + +import java.util.List; +import org.apache.iceberg.DeleteFile; + +public interface RemoveDanglingDeleteFiles + extends Action { + + interface Result { + List removedDeleteFiles(); + } +} diff --git a/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java b/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java index f6ef40270852..066661b566c2 100644 --- a/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java +++ b/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java @@ -106,6 +106,19 @@ public interface RewriteDataFiles boolean USE_STARTING_SEQUENCE_NUMBER_DEFAULT = true; + /** + * Remove dangling delete files from the current snapshot after compaction. A delete file is + * considered dangling if it does not apply to any non-expired data file. + * + *

Dangling delete files will be pruned from iceberg metadata. Pruning apply to both position + * delete and equality delete based on data sequence number + * + *

Defaults to false. + */ + String REMOVE_DANGLING_DELETES = "remove-dangling-deletes"; + + boolean REMOVE_DANGLING_DELETES_DEFAULT = false; + /** * Forces the rewrite job order based on the value. * @@ -216,6 +229,10 @@ default long rewrittenBytesCount() { default int failedDataFilesCount() { return rewriteFailures().stream().mapToInt(FileGroupFailureResult::dataFilesCount).sum(); } + + default int removedDeleteFilesCount() { + return 0; + } } /** diff --git a/core/src/main/java/org/apache/iceberg/actions/BaseRewriteDataFiles.java b/core/src/main/java/org/apache/iceberg/actions/BaseRewriteDataFiles.java index 953439484a15..2faa1f1b756c 100644 --- a/core/src/main/java/org/apache/iceberg/actions/BaseRewriteDataFiles.java +++ b/core/src/main/java/org/apache/iceberg/actions/BaseRewriteDataFiles.java @@ -55,6 +55,12 @@ default long rewrittenBytesCount() { return RewriteDataFiles.Result.super.rewrittenBytesCount(); } + @Override + @Value.Default + default int removedDeleteFilesCount() { + return RewriteDataFiles.Result.super.removedDeleteFilesCount(); + } + @Override @Value.Default default int failedDataFilesCount() { diff --git a/core/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFilesActionResult.java b/core/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFilesActionResult.java new file mode 100644 index 000000000000..0559d14fa535 --- /dev/null +++ b/core/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFilesActionResult.java @@ -0,0 +1,36 @@ +/* + * 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.actions; + +import java.util.List; +import org.apache.iceberg.DeleteFile; + +public class RemoveDanglingDeleteFilesActionResult implements RemoveDanglingDeleteFiles.Result { + + private final List removedDeleteFiles; + + public RemoveDanglingDeleteFilesActionResult(List removeDeleteFiles) { + this.removedDeleteFiles = removeDeleteFiles; + } + + @Override + public List removedDeleteFiles() { + return removedDeleteFiles; + } +} diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkContentFile.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkContentFile.java index f756c4cde015..99586f2503c2 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkContentFile.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkContentFile.java @@ -52,6 +52,7 @@ public abstract class SparkContentFile implements ContentFile { private final int keyMetadataPosition; private final int splitOffsetsPosition; private final int sortOrderIdPosition; + private final int fileSpecIdPosition; private final int equalityIdsPosition; private final Type lowerBoundsType; private final Type upperBoundsType; @@ -100,6 +101,7 @@ public abstract class SparkContentFile implements ContentFile { this.keyMetadataPosition = positions.get(DataFile.KEY_METADATA.name()); this.splitOffsetsPosition = positions.get(DataFile.SPLIT_OFFSETS.name()); this.sortOrderIdPosition = positions.get(DataFile.SORT_ORDER_ID.name()); + this.fileSpecIdPosition = positions.get(DataFile.SPEC_ID.name()); this.equalityIdsPosition = positions.get(DataFile.EQUALITY_IDS.name()); } @@ -120,7 +122,10 @@ public Long pos() { @Override public int specId() { - return -1; + if (wrapped.isNullAt(fileSpecIdPosition)) { + return -1; + } + return wrapped.getAs(fileSpecIdPosition); } @Override diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java new file mode 100644 index 000000000000..6b4e156a6559 --- /dev/null +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java @@ -0,0 +1,175 @@ +/* + * 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.spark.actions; + +import static org.apache.spark.sql.functions.col; +import static org.apache.spark.sql.functions.min; + +import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; +import org.apache.iceberg.DataFile; +import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.MetadataTableType; +import org.apache.iceberg.Partitioning; +import org.apache.iceberg.RewriteFiles; +import org.apache.iceberg.Table; +import org.apache.iceberg.actions.RemoveDanglingDeleteFiles; +import org.apache.iceberg.actions.RemoveDanglingDeleteFilesActionResult; +import org.apache.iceberg.spark.JobGroupInfo; +import org.apache.iceberg.spark.SparkDeleteFile; +import org.apache.iceberg.types.Types; +import org.apache.spark.sql.Column; +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.types.StructType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * An action that removes dangling delete files from the current snapshot. A delete file is dangling + * if its deletes no longer applies to any non-expired data file. + * + *

The following dangling delete files are removed: + * + *

    + *
  • Position delete files with a data sequence number less than that of any data file in the + * same partition + *
  • Equality delete files with a data sequence number less than or equal to that of any data + * file in the same partition + *
+ */ +class RemoveDanglingDeletesSparkAction + extends BaseSnapshotUpdateSparkAction + implements RemoveDanglingDeleteFiles { + + private static final Logger LOG = LoggerFactory.getLogger(RemoveDanglingDeletesSparkAction.class); + private final Table table; + + protected RemoveDanglingDeletesSparkAction(SparkSession spark, Table table) { + super(spark.cloneSession()); + this.table = table; + } + + @Override + protected RemoveDanglingDeletesSparkAction self() { + return this; + } + + public Result execute() { + if (table.specs().size() == 1 && table.spec().isUnpartitioned()) { + // ManifestFilterManager already performs this table-wide delete on each commit + return new RemoveDanglingDeleteFilesActionResult(Collections.emptyList()); + } + + String desc = String.format("Removing dangling delete in %s", table.name()); + JobGroupInfo info = newJobGroupInfo("REMOVE-DELETES", desc); + return withJobGroupInfo(info, this::doExecute); + } + + Result doExecute() { + RewriteFiles rewriteFiles = table.newRewrite(); + List danglingDeletes = findDanglingDeletes(); + for (DeleteFile deleteFile : danglingDeletes) { + LOG.debug("Removing dangling delete file {}", deleteFile.path()); + rewriteFiles.deleteFile(deleteFile); + } + + if (!danglingDeletes.isEmpty()) { + commit(rewriteFiles); + } + + return new RemoveDanglingDeleteFilesActionResult(danglingDeletes); + } + + /** + * Dangling delete files can be identified with following steps + * + *

1. Query live data entries table to group by partition spec ID and partition value aggregate + * to compute min data sequence number per group + * + *

2. Left join live delete entries table on grouped partition spec ID and partition value to + * account for partition evolution + * + *

3. Filter to identify dangling deletes that can be discarded by comparing its data sequence + * number having single predicate to account for both position and equality deletes + * + *

4. Collect results row to driver and use {@link SparkDeleteFile SparkDeleteFile} to wrap + * rows to valid delete files + */ + private List findDanglingDeletes() { + Dataset minSequenceNumberByPartition = + loadMetadataTable(table, MetadataTableType.ENTRIES) + .filter(" data_file.content == 0 AND status < 2") + .selectExpr( + "data_file.partition as partition", + "data_file.spec_id as spec_id", + "sequence_number") + .groupBy("partition", "spec_id") + .agg(min("sequence_number")) + .toDF("grouped_partition", "grouped_spec_id", "min_data_sequence_number"); + + Dataset deleteEntries = + loadMetadataTable(table, MetadataTableType.ENTRIES) + .filter(" data_file.content != 0 AND status < 2"); + + Column joinCond = + deleteEntries + .col("data_file.spec_id") + .equalTo(minSequenceNumberByPartition.col("grouped_spec_id")) + .and( + deleteEntries + .col("data_file.partition") + .equalTo(minSequenceNumberByPartition.col("grouped_partition"))); + + Column filterCondition = + col("min_data_sequence_number") + .isNull() + // dangling position delete files + .or( + col("data_file.content") + .equalTo("1") + .and(col("sequence_number").$less(col("min_data_sequence_number")))) + // dangling equality delete files + .or( + col("data_file.content") + .equalTo("2") + .and(col("sequence_number").$less$eq(col("min_data_sequence_number")))); + + Dataset danglingDeletes = + deleteEntries + .join(minSequenceNumberByPartition, joinCond, "left") + .filter(filterCondition) + .select("data_file.*"); + return danglingDeletes.collectAsList().stream() + .map( + row -> + deleteFileWrapper(danglingDeletes.schema(), row.getInt(row.fieldIndex("spec_id"))) + .wrap(row)) + .collect(Collectors.toList()); + } + + private SparkDeleteFile deleteFileWrapper(StructType sparkFileType, int specId) { + Types.StructType combinedFileType = DataFile.getType(Partitioning.partitionType(table)); + // deleteFile need to use the same spec for which manifest was written in order for deletion + Types.StructType projection = DataFile.getType(table.specs().get(specId).partitionType()); + return new SparkDeleteFile(combinedFileType, projection, sparkFileType); + } +} diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java index d33e5e540893..9a3e62da826c 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java @@ -33,6 +33,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import org.apache.iceberg.DataFile; +import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.RewriteJobOrder; import org.apache.iceberg.SortOrder; @@ -82,8 +83,9 @@ public class RewriteDataFilesSparkAction PARTIAL_PROGRESS_MAX_FAILED_COMMITS, TARGET_FILE_SIZE_BYTES, USE_STARTING_SEQUENCE_NUMBER, - REWRITE_JOB_ORDER, - OUTPUT_SPEC_ID); + OUTPUT_SPEC_ID, + REMOVE_DANGLING_DELETES, + REWRITE_JOB_ORDER); private static final RewriteDataFilesSparkAction.Result EMPTY_RESULT = ImmutableRewriteDataFiles.Result.builder().rewriteResults(ImmutableList.of()).build(); @@ -95,6 +97,7 @@ public class RewriteDataFilesSparkAction private int maxCommits; private int maxFailedCommits; private boolean partialProgressEnabled; + private boolean removeDanglingDeletes; private boolean useStartingSequenceNumber; private RewriteJobOrder rewriteJobOrder; private FileRewriter rewriter = null; @@ -175,11 +178,21 @@ public RewriteDataFiles.Result execute() { Stream groupStream = toGroupStream(ctx, fileGroupsByPartition); + ImmutableRewriteDataFiles.Result.Builder resultBuilder; if (partialProgressEnabled) { - return doExecuteWithPartialProgress(ctx, groupStream, commitManager(startingSnapshotId)); + resultBuilder = + doExecuteWithPartialProgress(ctx, groupStream, commitManager(startingSnapshotId)); } else { - return doExecute(ctx, groupStream, commitManager(startingSnapshotId)); + resultBuilder = doExecute(ctx, groupStream, commitManager(startingSnapshotId)); } + + if (removeDanglingDeletes) { + RemoveDanglingDeletesSparkAction action = + new RemoveDanglingDeletesSparkAction(spark(), table); + List removed = action.execute().removedDeleteFiles(); + resultBuilder.removedDeleteFilesCount(removed.size()); + } + return resultBuilder.build(); } StructLikeMap>> planFileGroups(long startingSnapshotId) { @@ -264,7 +277,7 @@ RewriteDataFilesCommitManager commitManager(long startingSnapshotId) { table, startingSnapshotId, useStartingSequenceNumber, commitSummary()); } - private Result doExecute( + private ImmutableRewriteDataFiles.Result.Builder doExecute( RewriteExecutionContext ctx, Stream groupStream, RewriteDataFilesCommitManager commitManager) { @@ -326,10 +339,10 @@ private Result doExecute( List rewriteResults = rewrittenGroups.stream().map(RewriteFileGroup::asResult).collect(Collectors.toList()); - return ImmutableRewriteDataFiles.Result.builder().rewriteResults(rewriteResults).build(); + return ImmutableRewriteDataFiles.Result.builder().rewriteResults(rewriteResults); } - private Result doExecuteWithPartialProgress( + private ImmutableRewriteDataFiles.Result.Builder doExecuteWithPartialProgress( RewriteExecutionContext ctx, Stream groupStream, RewriteDataFilesCommitManager commitManager) { @@ -386,8 +399,7 @@ private Result doExecuteWithPartialProgress( return ImmutableRewriteDataFiles.Result.builder() .rewriteResults(toRewriteResults(commitService.results())) - .rewriteFailures(rewriteFailures) - .build(); + .rewriteFailures(rewriteFailures); } Stream toGroupStream( @@ -456,6 +468,10 @@ void validateAndInitOptions() { PropertyUtil.propertyAsBoolean( options(), USE_STARTING_SEQUENCE_NUMBER, USE_STARTING_SEQUENCE_NUMBER_DEFAULT); + removeDanglingDeletes = + PropertyUtil.propertyAsBoolean( + options(), REMOVE_DANGLING_DELETES, REMOVE_DANGLING_DELETES_DEFAULT); + rewriteJobOrder = RewriteJobOrder.fromName( PropertyUtil.propertyAsString(options(), REWRITE_JOB_ORDER, REWRITE_JOB_ORDER_DEFAULT)); diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java index fb67ded96e35..c30859994ec7 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java @@ -20,6 +20,7 @@ import org.apache.iceberg.Table; import org.apache.iceberg.actions.ActionsProvider; +import org.apache.iceberg.actions.RemoveDanglingDeleteFiles; import org.apache.iceberg.spark.Spark3Util; import org.apache.iceberg.spark.Spark3Util.CatalogAndIdentifier; import org.apache.spark.sql.SparkSession; @@ -96,4 +97,9 @@ public DeleteReachableFilesSparkAction deleteReachableFiles(String metadataLocat public RewritePositionDeleteFilesSparkAction rewritePositionDeletes(Table table) { return new RewritePositionDeleteFilesSparkAction(spark, table); } + + @Override + public RemoveDanglingDeleteFiles removeDanglingDeleteFiles(Table table) { + return new RemoveDanglingDeletesSparkAction(spark, table); + } } diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java new file mode 100644 index 000000000000..533356bf4e86 --- /dev/null +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java @@ -0,0 +1,443 @@ +/* + * 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.spark.actions; + +import static org.apache.iceberg.types.Types.NestedField.optional; + +import java.io.File; +import java.nio.file.Path; +import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; +import org.apache.hadoop.conf.Configuration; +import org.apache.iceberg.DataFile; +import org.apache.iceberg.DataFiles; +import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.FileMetadata; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.Schema; +import org.apache.iceberg.Table; +import org.apache.iceberg.TableProperties; +import org.apache.iceberg.actions.RemoveDanglingDeleteFiles; +import org.apache.iceberg.hadoop.HadoopTables; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.spark.TestBase; +import org.apache.iceberg.types.Types; +import org.apache.spark.sql.Encoders; +import org.assertj.core.api.Assertions; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import scala.Tuple2; + +public class TestRemoveDanglingDeleteAction extends TestBase { + + private static final HadoopTables TABLES = new HadoopTables(new Configuration()); + private static final Schema SCHEMA = + new Schema( + optional(1, "c1", Types.StringType.get()), + optional(2, "c2", Types.StringType.get()), + optional(3, "c3", Types.StringType.get())); + + private static final PartitionSpec SPEC = PartitionSpec.builderFor(SCHEMA).identity("c1").build(); + + static final DataFile FILE_A = + DataFiles.builder(SPEC) + .withPath("/path/to/data-a.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=a") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DataFile FILE_A2 = + DataFiles.builder(SPEC) + .withPath("/path/to/data-a.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=a") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DataFile FILE_B = + DataFiles.builder(SPEC) + .withPath("/path/to/data-b.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=b") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DataFile FILE_B2 = + DataFiles.builder(SPEC) + .withPath("/path/to/data-b.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=b") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DataFile FILE_C = + DataFiles.builder(SPEC) + .withPath("/path/to/data-c.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=c") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DataFile FILE_C2 = + DataFiles.builder(SPEC) + .withPath("/path/to/data-c.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=c") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DataFile FILE_D = + DataFiles.builder(SPEC) + .withPath("/path/to/data-d.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=d") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DataFile FILE_D2 = + DataFiles.builder(SPEC) + .withPath("/path/to/data-d.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=d") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_A_POS_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofPositionDeletes() + .withPath("/path/to/data-a-pos-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=a") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_A2_POS_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofPositionDeletes() + .withPath("/path/to/data-a2-pos-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=a") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_A_EQ_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofEqualityDeletes() + .withPath("/path/to/data-a-eq-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=a") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_A2_EQ_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofEqualityDeletes() + .withPath("/path/to/data-a2-eq-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=a") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_B_POS_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofPositionDeletes() + .withPath("/path/to/data-b-pos-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=b") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_B2_POS_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofPositionDeletes() + .withPath("/path/to/data-b2-pos-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=b") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_B_EQ_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofEqualityDeletes() + .withPath("/path/to/data-b-eq-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=b") // easy way to set partition data for now + .withRecordCount(1) + .build(); + static final DeleteFile FILE_B2_EQ_DELETES = + FileMetadata.deleteFileBuilder(SPEC) + .ofEqualityDeletes() + .withPath("/path/to/data-b2-eq-deletes.parquet") + .withFileSizeInBytes(10) + .withPartitionPath("c1=b") // easy way to set partition data for now + .withRecordCount(1) + .build(); + + static final DataFile FILE_UNPARTITIONED = + DataFiles.builder(PartitionSpec.unpartitioned()) + .withPath("/path/to/data-unpartitioned.parquet") + .withFileSizeInBytes(10) + .withRecordCount(1) + .build(); + static final DeleteFile FILE_UNPARTITIONED_POS_DELETE = + FileMetadata.deleteFileBuilder(PartitionSpec.unpartitioned()) + .ofEqualityDeletes() + .withPath("/path/to/data-unpartitioned-pos-deletes.parquet") + .withFileSizeInBytes(10) + .withRecordCount(1) + .build(); + static final DeleteFile FILE_UNPARTITIONED_EQ_DELETE = + FileMetadata.deleteFileBuilder(PartitionSpec.unpartitioned()) + .ofEqualityDeletes() + .withPath("/path/to/data-unpartitioned-eq-deletes.parquet") + .withFileSizeInBytes(10) + .withRecordCount(1) + .build(); + + @TempDir private Path temp; + + private String tableLocation = null; + private Table table; + + @BeforeEach + public void before() throws Exception { + File tableDir = temp.resolve("junit").toFile(); + this.tableLocation = tableDir.toURI().toString(); + } + + @AfterEach + public void after() { + TABLES.dropTable(tableLocation); + } + + private void setupPartitionedTable() { + this.table = + TABLES.create( + SCHEMA, SPEC, ImmutableMap.of(TableProperties.FORMAT_VERSION, "2"), tableLocation); + } + + private void setupUnpartitionedTable() { + this.table = + TABLES.create( + SCHEMA, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2"), + tableLocation); + } + + @Test + public void testPartitionedDeletesWithLesserSeqNo() { + setupPartitionedTable(); + + // Add Data Files + table.newAppend().appendFile(FILE_B).appendFile(FILE_C).appendFile(FILE_D).commit(); + + // Add Delete Files + table + .newRowDelta() + .addDeletes(FILE_A_POS_DELETES) + .addDeletes(FILE_A2_POS_DELETES) + .addDeletes(FILE_B_POS_DELETES) + .addDeletes(FILE_B2_POS_DELETES) + .addDeletes(FILE_A_EQ_DELETES) + .addDeletes(FILE_A2_EQ_DELETES) + .addDeletes(FILE_B_EQ_DELETES) + .addDeletes(FILE_B2_EQ_DELETES) + .commit(); + + // Add More Data Files + table + .newAppend() + .appendFile(FILE_A2) + .appendFile(FILE_B2) + .appendFile(FILE_C2) + .appendFile(FILE_D2) + .commit(); + + List> actual = + spark + .read() + .format("iceberg") + .load(tableLocation + "#entries") + .select("sequence_number", "data_file.file_path") + .sort("sequence_number", "data_file.file_path") + .as(Encoders.tuple(Encoders.LONG(), Encoders.STRING())) + .collectAsList(); + List> expected = + ImmutableList.of( + Tuple2.apply(1L, FILE_B.path().toString()), + Tuple2.apply(1L, FILE_C.path().toString()), + Tuple2.apply(1L, FILE_D.path().toString()), + Tuple2.apply(2L, FILE_A_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A2_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A2_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B2_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B2_POS_DELETES.path().toString()), + Tuple2.apply(3L, FILE_A2.path().toString()), + Tuple2.apply(3L, FILE_B2.path().toString()), + Tuple2.apply(3L, FILE_C2.path().toString()), + Tuple2.apply(3L, FILE_D2.path().toString())); + Assertions.assertThat(actual).isEqualTo(expected); + + RemoveDanglingDeleteFiles.Result result = + SparkActions.get().removeDanglingDeleteFiles(table).execute(); + + // All Delete files of the FILE A partition should be removed + // because there are no data files in partition with a lesser sequence number + Set removedDeleteFiles = + result.removedDeleteFiles().stream().map(DeleteFile::path).collect(Collectors.toSet()); + Assertions.assertThat(removedDeleteFiles) + .as("Expected 4 delete files removed") + .hasSize(4) + .containsExactlyInAnyOrder( + FILE_A_POS_DELETES.path(), + FILE_A2_POS_DELETES.path(), + FILE_A_EQ_DELETES.path(), + FILE_A2_EQ_DELETES.path()); + + List> actualAfter = + spark + .read() + .format("iceberg") + .load(tableLocation + "#entries") + .filter("status < 2") // live files + .select("sequence_number", "data_file.file_path") + .sort("sequence_number", "data_file.file_path") + .as(Encoders.tuple(Encoders.LONG(), Encoders.STRING())) + .collectAsList(); + List> expectedAfter = + ImmutableList.of( + Tuple2.apply(1L, FILE_B.path().toString()), + Tuple2.apply(1L, FILE_C.path().toString()), + Tuple2.apply(1L, FILE_D.path().toString()), + Tuple2.apply(2L, FILE_B_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B2_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B2_POS_DELETES.path().toString()), + Tuple2.apply(3L, FILE_A2.path().toString()), + Tuple2.apply(3L, FILE_B2.path().toString()), + Tuple2.apply(3L, FILE_C2.path().toString()), + Tuple2.apply(3L, FILE_D2.path().toString())); + Assertions.assertThat(actualAfter).isEqualTo(expectedAfter); + } + + @Test + public void testPartitionedDeletesWithEqSeqNo() { + setupPartitionedTable(); + + // Add Data Files + table.newAppend().appendFile(FILE_A).appendFile(FILE_C).appendFile(FILE_D).commit(); + + // Add Data Files with EQ and POS deletes + table + .newRowDelta() + .addRows(FILE_A2) + .addRows(FILE_B2) + .addRows(FILE_C2) + .addRows(FILE_D2) + .addDeletes(FILE_A_POS_DELETES) + .addDeletes(FILE_A2_POS_DELETES) + .addDeletes(FILE_A_EQ_DELETES) + .addDeletes(FILE_A2_EQ_DELETES) + .addDeletes(FILE_B_POS_DELETES) + .addDeletes(FILE_B2_POS_DELETES) + .addDeletes(FILE_B_EQ_DELETES) + .addDeletes(FILE_B2_EQ_DELETES) + .commit(); + + List> actual = + spark + .read() + .format("iceberg") + .load(tableLocation + "#entries") + .select("sequence_number", "data_file.file_path") + .sort("sequence_number", "data_file.file_path") + .as(Encoders.tuple(Encoders.LONG(), Encoders.STRING())) + .collectAsList(); + List> expected = + ImmutableList.of( + Tuple2.apply(1L, FILE_A.path().toString()), + Tuple2.apply(1L, FILE_C.path().toString()), + Tuple2.apply(1L, FILE_D.path().toString()), + Tuple2.apply(2L, FILE_A_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A2.path().toString()), + Tuple2.apply(2L, FILE_A2_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A2_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B2.path().toString()), + Tuple2.apply(2L, FILE_B2_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B2_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_C2.path().toString()), + Tuple2.apply(2L, FILE_D2.path().toString())); + Assertions.assertThat(actual).isEqualTo(expected); + + RemoveDanglingDeleteFiles.Result result = + SparkActions.get().removeDanglingDeleteFiles(table).execute(); + + // Eq Delete files of the FILE B partition should be removed + // because there are no data files in partition with a lesser sequence number + Set removedDeleteFiles = + result.removedDeleteFiles().stream().map(DeleteFile::path).collect(Collectors.toSet()); + Assertions.assertThat(removedDeleteFiles) + .as("Expected two delete files removed") + .hasSize(2) + .containsExactlyInAnyOrder(FILE_B_EQ_DELETES.path(), FILE_B2_EQ_DELETES.path()); + + List> actualAfter = + spark + .read() + .format("iceberg") + .load(tableLocation + "#entries") + .filter("status < 2") // live files + .select("sequence_number", "data_file.file_path") + .sort("sequence_number", "data_file.file_path") + .as(Encoders.tuple(Encoders.LONG(), Encoders.STRING())) + .collectAsList(); + List> expectedAfter = + ImmutableList.of( + Tuple2.apply(1L, FILE_A.path().toString()), + Tuple2.apply(1L, FILE_C.path().toString()), + Tuple2.apply(1L, FILE_D.path().toString()), + Tuple2.apply(2L, FILE_A_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A2.path().toString()), + Tuple2.apply(2L, FILE_A2_EQ_DELETES.path().toString()), + Tuple2.apply(2L, FILE_A2_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_B2.path().toString()), + Tuple2.apply(2L, FILE_B2_POS_DELETES.path().toString()), + Tuple2.apply(2L, FILE_C2.path().toString()), + Tuple2.apply(2L, FILE_D2.path().toString())); + Assertions.assertThat(actualAfter).isEqualTo(expectedAfter); + } + + @Test + public void testUnpartitionedTable() { + setupUnpartitionedTable(); + + table + .newRowDelta() + .addDeletes(FILE_UNPARTITIONED_POS_DELETE) + .addDeletes(FILE_UNPARTITIONED_EQ_DELETE) + .commit(); + table.newAppend().appendFile(FILE_UNPARTITIONED).commit(); + + RemoveDanglingDeleteFiles.Result result = + SparkActions.get().removeDanglingDeleteFiles(table).execute(); + Assertions.assertThat(result.removedDeleteFiles()) + .as("No-op for unpartitioned tables") + .isEmpty(); + } +} diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java index b67ee87c7d3e..2de83f8b355c 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewriteDataFilesAction.java @@ -24,6 +24,7 @@ import static org.apache.spark.sql.functions.current_date; import static org.apache.spark.sql.functions.date_add; import static org.apache.spark.sql.functions.expr; +import static org.apache.spark.sql.functions.min; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; @@ -56,6 +57,7 @@ import org.apache.iceberg.FileScanTask; import org.apache.iceberg.MetadataTableType; import org.apache.iceberg.PartitionData; +import org.apache.iceberg.PartitionKey; import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.RewriteJobOrder; import org.apache.iceberg.RowDelta; @@ -73,7 +75,9 @@ import org.apache.iceberg.actions.SizeBasedDataRewriter; import org.apache.iceberg.actions.SizeBasedFileRewriter; import org.apache.iceberg.data.GenericAppenderFactory; +import org.apache.iceberg.data.GenericRecord; import org.apache.iceberg.data.Record; +import org.apache.iceberg.deletes.EqualityDeleteWriter; import org.apache.iceberg.deletes.PositionDelete; import org.apache.iceberg.deletes.PositionDeleteWriter; import org.apache.iceberg.encryption.EncryptedFiles; @@ -86,6 +90,7 @@ import org.apache.iceberg.hadoop.HadoopTables; import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.io.OutputFile; +import org.apache.iceberg.io.OutputFileFactory; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet; @@ -105,9 +110,11 @@ import org.apache.iceberg.types.Conversions; import org.apache.iceberg.types.Types; import org.apache.iceberg.types.Types.NestedField; +import org.apache.iceberg.util.ArrayUtil; import org.apache.iceberg.util.Pair; import org.apache.iceberg.util.StructLikeMap; import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Encoders; import org.apache.spark.sql.Row; import org.apache.spark.sql.internal.SQLConf; import org.junit.jupiter.api.BeforeAll; @@ -128,6 +135,8 @@ public class TestRewriteDataFilesAction extends TestBase { optional(2, "c2", Types.StringType.get()), optional(3, "c3", Types.StringType.get())); + private static final PartitionSpec SPEC = PartitionSpec.builderFor(SCHEMA).identity("c1").build(); + @TempDir private Path temp; private final FileRewriteCoordinator coordinator = FileRewriteCoordinator.get(); @@ -336,6 +345,125 @@ public void testBinPackWithDeletes() { assertThat(actualRecords).as("7 rows are removed").hasSize(total - 7); } + @Test + public void testRemoveDangledEqualityDeletesPartitionEvolution() { + Table table = + TABLES.create( + SCHEMA, + SPEC, + Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"), + tableLocation); + + // data seq = 1, write 4 files in 2 partitions + List records1 = + Lists.newArrayList( + new ThreeColumnRecord(1, null, "AAAA"), new ThreeColumnRecord(1, "BBBBBBBBBB", "BBBB")); + writeRecords(records1); + List records2 = + Lists.newArrayList( + new ThreeColumnRecord(0, "CCCCCCCCCC", "CCCC"), + new ThreeColumnRecord(0, "DDDDDDDDDD", "DDDD")); + writeRecords(records2); + table.refresh(); + shouldHaveFiles(table, 4); + + // data seq = 2 & 3, write 2 equality deletes in both partitions + writeEqDeleteRecord(table, "c1", 1, "c3", "AAAA"); + writeEqDeleteRecord(table, "c1", 2, "c3", "CCCC"); + table.refresh(); + Set existingDeletes = TestHelpers.deleteFiles(table); + assertThat(existingDeletes) + .as("Only one equality delete c1=1 is used in query planning") + .hasSize(1); + + // partition evolution + table.refresh(); + table.updateSpec().addField(Expressions.ref("c3")).commit(); + + // data seq = 4, write 2 new data files in both partitions for evolved spec + List records3 = + Lists.newArrayList( + new ThreeColumnRecord(1, "A", "CCCC"), new ThreeColumnRecord(2, "D", "DDDD")); + writeRecords(records3); + + List originalData = currentData(); + + RewriteDataFiles.Result result = + basicRewrite(table) + .option(SizeBasedFileRewriter.REWRITE_ALL, "true") + .filter(Expressions.equal("c1", 1)) + .option(RewriteDataFiles.REMOVE_DANGLING_DELETES, "true") + .execute(); + + existingDeletes = TestHelpers.deleteFiles(table); + assertThat(existingDeletes).as("Shall pruned dangling deletes after rewrite").hasSize(0); + + assertThat(result) + .extracting( + Result::addedDataFilesCount, + Result::rewrittenDataFilesCount, + Result::removedDeleteFilesCount) + .as("Should compact 3 data files into 2 and remove both dangled equality delete file") + .containsExactly(2, 3, 2); + shouldHaveMinSequenceNumberInPartition(table, "data_file.partition.c1 == 1", 5); + + List postRewriteData = currentData(); + assertEquals("We shouldn't have changed the data", originalData, postRewriteData); + + shouldHaveSnapshots(table, 7); + shouldHaveFiles(table, 5); + } + + @Test + public void testRemoveDangledPositionDeletesPartitionEvolution() { + Table table = + TABLES.create( + SCHEMA, + SPEC, + Collections.singletonMap(TableProperties.FORMAT_VERSION, "2"), + tableLocation); + + // data seq = 1, write 4 files in 2 partitions + writeRecords(2, 2, 2); + List dataFilesBefore = TestHelpers.dataFiles(table, null); + shouldHaveFiles(table, 4); + + // data seq = 2, write 1 position deletes in c1=1 + table + .newRowDelta() + .addDeletes(writePosDeletesToFile(table, dataFilesBefore.get(3), 1).get(0)) + .commit(); + + // partition evolution + table.updateSpec().addField(Expressions.ref("c3")).commit(); + + // data seq = 3, write 1 new data files in c1=1 for evolved spec + writeRecords(1, 1, 1); + shouldHaveFiles(table, 5); + List expectedRecords = currentData(); + + Result result = + actions() + .rewriteDataFiles(table) + .filter(Expressions.equal("c1", 1)) + .option(SizeBasedFileRewriter.REWRITE_ALL, "true") + .option(RewriteDataFiles.REMOVE_DANGLING_DELETES, "true") + .execute(); + + assertThat(result) + .extracting( + Result::addedDataFilesCount, + Result::rewrittenDataFilesCount, + Result::removedDeleteFilesCount) + .as("Should rewrite 2 data files into 1 and remove 1 dangled position delete file") + .containsExactly(1, 2, 1); + shouldHaveMinSequenceNumberInPartition(table, "data_file.partition.c1 == 1", 3); + + shouldHaveSnapshots(table, 5); + assertThat(table.currentSnapshot().summary().get("total-position-deletes")).isEqualTo("0"); + assertEquals("Rows must match", expectedRecords, currentData()); + } + @Test public void testBinPackWithDeleteAllData() { Map options = Maps.newHashMap(); @@ -1697,6 +1825,21 @@ protected void shouldHaveFiles(Table table, int numExpected) { assertThat(numFiles).as("Did not have the expected number of files").isEqualTo(numExpected); } + protected long shouldHaveMinSequenceNumberInPartition( + Table table, String partitionFilter, long expected) { + long actual = + SparkTableUtil.loadMetadataTable(spark, table, MetadataTableType.ENTRIES) + .filter("status != 2") + .filter(partitionFilter) + .select("sequence_number") + .agg(min("sequence_number")) + .as(Encoders.LONG()) + .collectAsList() + .get(0); + assertThat(actual).as("Did not have the expected min sequence number").isEqualTo(expected); + return actual; + } + protected void shouldHaveSnapshots(Table table, int expectedSnapshots) { table.refresh(); int actualSnapshots = Iterables.size(table.snapshots()); @@ -1893,6 +2036,11 @@ protected int averageFileSize(Table table) { .getAsDouble(); } + private void writeRecords(List records) { + Dataset df = spark.createDataFrame(records, ThreeColumnRecord.class); + writeDF(df); + } + private void writeRecords(int files, int numRecords) { writeRecords(files, numRecords, 0); } @@ -1946,7 +2094,10 @@ private List writePosDeletes( table .io() .newOutputFile( - table.locationProvider().newDataLocation(UUID.randomUUID().toString())); + table + .locationProvider() + .newDataLocation( + FileFormat.PARQUET.addExtension(UUID.randomUUID().toString()))); EncryptedOutputFile encryptedOutputFile = EncryptedFiles.encryptedOutput(outputFile, EncryptionKeyMetadata.EMPTY); @@ -1972,6 +2123,68 @@ private List writePosDeletes( return results; } + private void writeEqDeleteRecord( + Table table, String partCol, Object partVal, String delCol, Object delVal) { + List equalityFieldIds = Lists.newArrayList(table.schema().findField(delCol).fieldId()); + Schema eqDeleteRowSchema = table.schema().select(delCol); + Record partitionRecord = + GenericRecord.create(table.schema().select(partCol)) + .copy(ImmutableMap.of(partCol, partVal)); + Record record = GenericRecord.create(eqDeleteRowSchema).copy(ImmutableMap.of(delCol, delVal)); + writeEqDeleteRecord(table, equalityFieldIds, partitionRecord, eqDeleteRowSchema, record); + } + + private void writeEqDeleteRecord( + Table table, + List equalityFieldIds, + Record partitionRecord, + Schema eqDeleteRowSchema, + Record deleteRecord) { + OutputFileFactory fileFactory = + OutputFileFactory.builderFor(table, 1, 1).format(FileFormat.PARQUET).build(); + GenericAppenderFactory appenderFactory = + new GenericAppenderFactory( + table.schema(), + table.spec(), + ArrayUtil.toIntArray(equalityFieldIds), + eqDeleteRowSchema, + null); + + EncryptedOutputFile file = + createEncryptedOutputFile(createPartitionKey(table, partitionRecord), fileFactory); + + EqualityDeleteWriter eqDeleteWriter = + appenderFactory.newEqDeleteWriter( + file, FileFormat.PARQUET, createPartitionKey(table, partitionRecord)); + + try (EqualityDeleteWriter clsEqDeleteWriter = eqDeleteWriter) { + clsEqDeleteWriter.write(deleteRecord); + } catch (Exception e) { + throw new RuntimeException(e); + } + table.newRowDelta().addDeletes(eqDeleteWriter.toDeleteFile()).commit(); + } + + private PartitionKey createPartitionKey(Table table, Record record) { + if (table.spec().isUnpartitioned()) { + return null; + } + + PartitionKey partitionKey = new PartitionKey(table.spec(), table.schema()); + partitionKey.partition(record); + + return partitionKey; + } + + private EncryptedOutputFile createEncryptedOutputFile( + PartitionKey partition, OutputFileFactory fileFactory) { + if (partition == null) { + return fileFactory.newOutputFile(); + } else { + return fileFactory.newOutputFile(partition); + } + } + private SparkActions actions() { return SparkActions.get(); } diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewritePositionDeleteFilesAction.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewritePositionDeleteFilesAction.java index 37b6cd86fb92..8547f9753f5e 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewritePositionDeleteFilesAction.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRewritePositionDeleteFilesAction.java @@ -862,6 +862,7 @@ private void writePosDeletesForFiles( files.stream().collect(Collectors.groupingBy(ContentFile::partition)); List deleteFiles = Lists.newArrayListWithCapacity(deleteFilesPerPartition * filesByPartition.size()); + String suffix = String.format(".%s", FileFormat.PARQUET.name().toLowerCase()); for (Map.Entry> filesByPartitionEntry : filesByPartition.entrySet()) { @@ -886,7 +887,7 @@ private void writePosDeletesForFiles( if (counter == deleteFileSize) { // Dump to file and reset variables OutputFile output = - Files.localOutput(File.createTempFile("junit", null, temp.toFile())); + Files.localOutput(File.createTempFile("junit", suffix, temp.toFile())); deleteFiles.add(FileHelpers.writeDeleteFile(table, output, partition, deletes).first()); counter = 0; deletes.clear(); From 82b68890007d4879db90ca7565023146adcc7c14 Mon Sep 17 00:00:00 2001 From: Hongyue Zhang Date: Thu, 25 Jul 2024 23:40:38 -0700 Subject: [PATCH 2/5] Fix assertJ for static import --- .../TestRemoveDanglingDeleteAction.java | 18 ++++++++---------- 1 file changed, 8 insertions(+), 10 deletions(-) diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java index 533356bf4e86..287b91b5007d 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java @@ -19,6 +19,7 @@ package org.apache.iceberg.spark.actions; import static org.apache.iceberg.types.Types.NestedField.optional; +import static org.assertj.core.api.Assertions.assertThat; import java.io.File; import java.nio.file.Path; @@ -41,7 +42,6 @@ import org.apache.iceberg.spark.TestBase; import org.apache.iceberg.types.Types; import org.apache.spark.sql.Encoders; -import org.assertj.core.api.Assertions; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -287,7 +287,7 @@ public void testPartitionedDeletesWithLesserSeqNo() { Tuple2.apply(3L, FILE_B2.path().toString()), Tuple2.apply(3L, FILE_C2.path().toString()), Tuple2.apply(3L, FILE_D2.path().toString())); - Assertions.assertThat(actual).isEqualTo(expected); + assertThat(actual).isEqualTo(expected); RemoveDanglingDeleteFiles.Result result = SparkActions.get().removeDanglingDeleteFiles(table).execute(); @@ -296,7 +296,7 @@ public void testPartitionedDeletesWithLesserSeqNo() { // because there are no data files in partition with a lesser sequence number Set removedDeleteFiles = result.removedDeleteFiles().stream().map(DeleteFile::path).collect(Collectors.toSet()); - Assertions.assertThat(removedDeleteFiles) + assertThat(removedDeleteFiles) .as("Expected 4 delete files removed") .hasSize(4) .containsExactlyInAnyOrder( @@ -328,7 +328,7 @@ public void testPartitionedDeletesWithLesserSeqNo() { Tuple2.apply(3L, FILE_B2.path().toString()), Tuple2.apply(3L, FILE_C2.path().toString()), Tuple2.apply(3L, FILE_D2.path().toString())); - Assertions.assertThat(actualAfter).isEqualTo(expectedAfter); + assertThat(actualAfter).isEqualTo(expectedAfter); } @Test @@ -381,7 +381,7 @@ public void testPartitionedDeletesWithEqSeqNo() { Tuple2.apply(2L, FILE_B2_POS_DELETES.path().toString()), Tuple2.apply(2L, FILE_C2.path().toString()), Tuple2.apply(2L, FILE_D2.path().toString())); - Assertions.assertThat(actual).isEqualTo(expected); + assertThat(actual).isEqualTo(expected); RemoveDanglingDeleteFiles.Result result = SparkActions.get().removeDanglingDeleteFiles(table).execute(); @@ -390,7 +390,7 @@ public void testPartitionedDeletesWithEqSeqNo() { // because there are no data files in partition with a lesser sequence number Set removedDeleteFiles = result.removedDeleteFiles().stream().map(DeleteFile::path).collect(Collectors.toSet()); - Assertions.assertThat(removedDeleteFiles) + assertThat(removedDeleteFiles) .as("Expected two delete files removed") .hasSize(2) .containsExactlyInAnyOrder(FILE_B_EQ_DELETES.path(), FILE_B2_EQ_DELETES.path()); @@ -420,7 +420,7 @@ public void testPartitionedDeletesWithEqSeqNo() { Tuple2.apply(2L, FILE_B2_POS_DELETES.path().toString()), Tuple2.apply(2L, FILE_C2.path().toString()), Tuple2.apply(2L, FILE_D2.path().toString())); - Assertions.assertThat(actualAfter).isEqualTo(expectedAfter); + assertThat(actualAfter).isEqualTo(expectedAfter); } @Test @@ -436,8 +436,6 @@ public void testUnpartitionedTable() { RemoveDanglingDeleteFiles.Result result = SparkActions.get().removeDanglingDeleteFiles(table).execute(); - Assertions.assertThat(result.removedDeleteFiles()) - .as("No-op for unpartitioned tables") - .isEmpty(); + assertThat(result.removedDeleteFiles()).as("No-op for unpartitioned tables").isEmpty(); } } From f6b0be37270a6248b9df6ed80bfa1f601adbf40b Mon Sep 17 00:00:00 2001 From: Hongyue Zhang Date: Tue, 13 Aug 2024 17:10:25 -0700 Subject: [PATCH 3/5] Address szehon feedback - Rewording documentation and add more comments - Changed removed deletes in results from List to Iterable to save on memory - Added BaseRemoveDanglingDeleteFiles and generate Immutable implementation --- .../actions/RemoveDanglingDeleteFiles.java | 9 +++- .../iceberg/actions/RewriteDataFiles.java | 5 +- ...ava => BaseRemoveDanglingDeleteFiles.java} | 23 ++++----- .../RemoveDanglingDeletesSparkAction.java | 48 ++++++++++--------- .../actions/RewriteDataFilesSparkAction.java | 17 +++---- .../TestRemoveDanglingDeleteAction.java | 10 +++- 6 files changed, 62 insertions(+), 50 deletions(-) rename core/src/main/java/org/apache/iceberg/actions/{RemoveDanglingDeleteFilesActionResult.java => BaseRemoveDanglingDeleteFiles.java} (65%) diff --git a/api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java b/api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java index 879dc9f0c196..b0ef0d5e35f8 100644 --- a/api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java +++ b/api/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFiles.java @@ -18,13 +18,18 @@ */ package org.apache.iceberg.actions; -import java.util.List; import org.apache.iceberg.DeleteFile; +/** + * An action that removes dangling delete files from the current snapshot. A delete file is dangling + * if its deletes no longer applies to any live data files. + */ public interface RemoveDanglingDeleteFiles extends Action { + /** An action that remove dangling deletes. */ interface Result { - List removedDeleteFiles(); + /** Return removed deletes. */ + Iterable removedDeleteFiles(); } } diff --git a/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java b/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java index 066661b566c2..589b9017741e 100644 --- a/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java +++ b/api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java @@ -108,10 +108,9 @@ public interface RewriteDataFiles /** * Remove dangling delete files from the current snapshot after compaction. A delete file is - * considered dangling if it does not apply to any non-expired data file. + * considered dangling if it does not apply to any live data files. * - *

Dangling delete files will be pruned from iceberg metadata. Pruning apply to both position - * delete and equality delete based on data sequence number + *

Both equality and position dangling delete files will be removed. * *

Defaults to false. */ diff --git a/core/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFilesActionResult.java b/core/src/main/java/org/apache/iceberg/actions/BaseRemoveDanglingDeleteFiles.java similarity index 65% rename from core/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFilesActionResult.java rename to core/src/main/java/org/apache/iceberg/actions/BaseRemoveDanglingDeleteFiles.java index 0559d14fa535..3b5ce9e79a43 100644 --- a/core/src/main/java/org/apache/iceberg/actions/RemoveDanglingDeleteFilesActionResult.java +++ b/core/src/main/java/org/apache/iceberg/actions/BaseRemoveDanglingDeleteFiles.java @@ -18,19 +18,16 @@ */ package org.apache.iceberg.actions; -import java.util.List; -import org.apache.iceberg.DeleteFile; +import org.immutables.value.Value; -public class RemoveDanglingDeleteFilesActionResult implements RemoveDanglingDeleteFiles.Result { +@Value.Enclosing +@SuppressWarnings("ImmutablesStyle") +@Value.Style( + typeImmutableEnclosing = "ImmutableRemoveDanglingDeleteFiles", + visibilityString = "PUBLIC", + builderVisibilityString = "PUBLIC") +interface BaseRemoveDanglingDeleteFiles extends RemoveDanglingDeleteFiles { - private final List removedDeleteFiles; - - public RemoveDanglingDeleteFilesActionResult(List removeDeleteFiles) { - this.removedDeleteFiles = removeDeleteFiles; - } - - @Override - public List removedDeleteFiles() { - return removedDeleteFiles; - } + @Value.Immutable + interface Result extends RemoveDanglingDeleteFiles.Result {} } diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java index 6b4e156a6559..eff7fd13e8bf 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java @@ -30,8 +30,8 @@ import org.apache.iceberg.Partitioning; import org.apache.iceberg.RewriteFiles; import org.apache.iceberg.Table; +import org.apache.iceberg.actions.ImmutableRemoveDanglingDeleteFiles; import org.apache.iceberg.actions.RemoveDanglingDeleteFiles; -import org.apache.iceberg.actions.RemoveDanglingDeleteFilesActionResult; import org.apache.iceberg.spark.JobGroupInfo; import org.apache.iceberg.spark.SparkDeleteFile; import org.apache.iceberg.types.Types; @@ -45,7 +45,7 @@ /** * An action that removes dangling delete files from the current snapshot. A delete file is dangling - * if its deletes no longer applies to any non-expired data file. + * if its deletes no longer applies to any live data files. * *

The following dangling delete files are removed: * @@ -76,7 +76,9 @@ protected RemoveDanglingDeletesSparkAction self() { public Result execute() { if (table.specs().size() == 1 && table.spec().isUnpartitioned()) { // ManifestFilterManager already performs this table-wide delete on each commit - return new RemoveDanglingDeleteFilesActionResult(Collections.emptyList()); + return ImmutableRemoveDanglingDeleteFiles.Result.builder() + .removedDeleteFiles(Collections.emptyList()) + .build(); } String desc = String.format("Removing dangling delete in %s", table.name()); @@ -96,20 +98,21 @@ Result doExecute() { commit(rewriteFiles); } - return new RemoveDanglingDeleteFilesActionResult(danglingDeletes); + return ImmutableRemoveDanglingDeleteFiles.Result.builder() + .removedDeleteFiles(danglingDeletes) + .build(); } /** * Dangling delete files can be identified with following steps * - *

1. Query live data entries table to group by partition spec ID and partition value aggregate - * to compute min data sequence number per group + *

1.Group data files by partition keys and find the minimum data sequence number in each + * group. * - *

2. Left join live delete entries table on grouped partition spec ID and partition value to - * account for partition evolution + *

2. Left outer join delete files with partition-grouped data files on partition keys. * - *

3. Filter to identify dangling deletes that can be discarded by comparing its data sequence - * number having single predicate to account for both position and equality deletes + *

3. Filter results to find dangling delete files by comparing delete file sequence_number to + * its partitions' minimum data sequence number. * *

4. Collect results row to driver and use {@link SparkDeleteFile SparkDeleteFile} to wrap * rows to valid delete files @@ -117,7 +120,8 @@ Result doExecute() { private List findDanglingDeletes() { Dataset minSequenceNumberByPartition = loadMetadataTable(table, MetadataTableType.ENTRIES) - .filter(" data_file.content == 0 AND status < 2") + // find live entry for data files + .filter("data_file.content == 0 AND status < 2") .selectExpr( "data_file.partition as partition", "data_file.spec_id as spec_id", @@ -128,9 +132,10 @@ private List findDanglingDeletes() { Dataset deleteEntries = loadMetadataTable(table, MetadataTableType.ENTRIES) + // find live entry for delete files .filter(" data_file.content != 0 AND status < 2"); - Column joinCond = + Column joinOnPartition = deleteEntries .col("data_file.spec_id") .equalTo(minSequenceNumberByPartition.col("grouped_spec_id")) @@ -139,8 +144,9 @@ private List findDanglingDeletes() { .col("data_file.partition") .equalTo(minSequenceNumberByPartition.col("grouped_partition"))); - Column filterCondition = + Column filterOnDanglingDeletes = col("min_data_sequence_number") + // when all data files are rewritten with new partition spec .isNull() // dangling position delete files .or( @@ -155,21 +161,19 @@ private List findDanglingDeletes() { Dataset danglingDeletes = deleteEntries - .join(minSequenceNumberByPartition, joinCond, "left") - .filter(filterCondition) + .join(minSequenceNumberByPartition, joinOnPartition, "left") + .filter(filterOnDanglingDeletes) .select("data_file.*"); return danglingDeletes.collectAsList().stream() - .map( - row -> - deleteFileWrapper(danglingDeletes.schema(), row.getInt(row.fieldIndex("spec_id"))) - .wrap(row)) + .map(row -> deleteFileWrapper(danglingDeletes.schema(), row)) .collect(Collectors.toList()); } - private SparkDeleteFile deleteFileWrapper(StructType sparkFileType, int specId) { + private DeleteFile deleteFileWrapper(StructType sparkFileType, Row row) { + int specId = row.getInt(row.fieldIndex("spec_id")); Types.StructType combinedFileType = DataFile.getType(Partitioning.partitionType(table)); - // deleteFile need to use the same spec for which manifest was written in order for deletion + // Set correct spec id Types.StructType projection = DataFile.getType(table.specs().get(specId).partitionType()); - return new SparkDeleteFile(combinedFileType, projection, sparkFileType); + return new SparkDeleteFile(combinedFileType, projection, sparkFileType).wrap(row); } } diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java index 9a3e62da826c..4c0406a76489 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java @@ -33,7 +33,6 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import org.apache.iceberg.DataFile; -import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.RewriteJobOrder; import org.apache.iceberg.SortOrder; @@ -41,6 +40,7 @@ import org.apache.iceberg.Table; import org.apache.iceberg.actions.FileRewriter; import org.apache.iceberg.actions.ImmutableRewriteDataFiles; +import org.apache.iceberg.actions.ImmutableRewriteDataFiles.Result.Builder; import org.apache.iceberg.actions.RewriteDataFiles; import org.apache.iceberg.actions.RewriteDataFilesCommitManager; import org.apache.iceberg.actions.RewriteFileGroup; @@ -54,6 +54,7 @@ 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.ImmutableSet; +import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.relocated.com.google.common.collect.Queues; @@ -84,8 +85,8 @@ public class RewriteDataFilesSparkAction TARGET_FILE_SIZE_BYTES, USE_STARTING_SEQUENCE_NUMBER, OUTPUT_SPEC_ID, - REMOVE_DANGLING_DELETES, - REWRITE_JOB_ORDER); + REWRITE_JOB_ORDER, + REMOVE_DANGLING_DELETES); private static final RewriteDataFilesSparkAction.Result EMPTY_RESULT = ImmutableRewriteDataFiles.Result.builder().rewriteResults(ImmutableList.of()).build(); @@ -178,7 +179,7 @@ public RewriteDataFiles.Result execute() { Stream groupStream = toGroupStream(ctx, fileGroupsByPartition); - ImmutableRewriteDataFiles.Result.Builder resultBuilder; + Builder resultBuilder; if (partialProgressEnabled) { resultBuilder = doExecuteWithPartialProgress(ctx, groupStream, commitManager(startingSnapshotId)); @@ -189,8 +190,8 @@ public RewriteDataFiles.Result execute() { if (removeDanglingDeletes) { RemoveDanglingDeletesSparkAction action = new RemoveDanglingDeletesSparkAction(spark(), table); - List removed = action.execute().removedDeleteFiles(); - resultBuilder.removedDeleteFilesCount(removed.size()); + int removedCount = Iterables.size(action.execute().removedDeleteFiles()); + resultBuilder.removedDeleteFilesCount(removedCount); } return resultBuilder.build(); } @@ -277,7 +278,7 @@ RewriteDataFilesCommitManager commitManager(long startingSnapshotId) { table, startingSnapshotId, useStartingSequenceNumber, commitSummary()); } - private ImmutableRewriteDataFiles.Result.Builder doExecute( + private Builder doExecute( RewriteExecutionContext ctx, Stream groupStream, RewriteDataFilesCommitManager commitManager) { @@ -342,7 +343,7 @@ private ImmutableRewriteDataFiles.Result.Builder doExecute( return ImmutableRewriteDataFiles.Result.builder().rewriteResults(rewriteResults); } - private ImmutableRewriteDataFiles.Result.Builder doExecuteWithPartialProgress( + private Builder doExecuteWithPartialProgress( RewriteExecutionContext ctx, Stream groupStream, RewriteDataFilesCommitManager commitManager) { diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java index 287b91b5007d..e15b2fb2174a 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveDanglingDeleteAction.java @@ -26,6 +26,7 @@ import java.util.List; import java.util.Set; import java.util.stream.Collectors; +import java.util.stream.StreamSupport; import org.apache.hadoop.conf.Configuration; import org.apache.iceberg.DataFile; import org.apache.iceberg.DataFiles; @@ -294,8 +295,11 @@ public void testPartitionedDeletesWithLesserSeqNo() { // All Delete files of the FILE A partition should be removed // because there are no data files in partition with a lesser sequence number + Set removedDeleteFiles = - result.removedDeleteFiles().stream().map(DeleteFile::path).collect(Collectors.toSet()); + StreamSupport.stream(result.removedDeleteFiles().spliterator(), false) + .map(DeleteFile::path) + .collect(Collectors.toSet()); assertThat(removedDeleteFiles) .as("Expected 4 delete files removed") .hasSize(4) @@ -389,7 +393,9 @@ public void testPartitionedDeletesWithEqSeqNo() { // Eq Delete files of the FILE B partition should be removed // because there are no data files in partition with a lesser sequence number Set removedDeleteFiles = - result.removedDeleteFiles().stream().map(DeleteFile::path).collect(Collectors.toSet()); + StreamSupport.stream(result.removedDeleteFiles().spliterator(), false) + .map(DeleteFile::path) + .collect(Collectors.toSet()); assertThat(removedDeleteFiles) .as("Expected two delete files removed") .hasSize(2) From 18e2f798ad4c60aa3970236cf6ef25dc5130ec36 Mon Sep 17 00:00:00 2001 From: Hongyue Zhang Date: Tue, 20 Aug 2024 16:38:59 -0700 Subject: [PATCH 4/5] Address review feedback and improve docs - Instantiate spark action without clone session - Update javadoc to use html order list - Inline resultBuilder in RewriteDataFilesSparkAction --- .../RemoveDanglingDeletesSparkAction.java | 36 +++++++++---------- .../actions/RewriteDataFilesSparkAction.java | 13 +++---- 2 files changed, 23 insertions(+), 26 deletions(-) diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java index eff7fd13e8bf..bbf65f58e19c 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RemoveDanglingDeletesSparkAction.java @@ -64,7 +64,7 @@ class RemoveDanglingDeletesSparkAction private final Table table; protected RemoveDanglingDeletesSparkAction(SparkSession spark, Table table) { - super(spark.cloneSession()); + super(spark); this.table = table; } @@ -81,7 +81,7 @@ public Result execute() { .build(); } - String desc = String.format("Removing dangling delete in %s", table.name()); + String desc = String.format("Removing dangling delete files in %s", table.name()); JobGroupInfo info = newJobGroupInfo("REMOVE-DELETES", desc); return withJobGroupInfo(info, this::doExecute); } @@ -106,21 +106,20 @@ Result doExecute() { /** * Dangling delete files can be identified with following steps * - *

1.Group data files by partition keys and find the minimum data sequence number in each - * group. - * - *

2. Left outer join delete files with partition-grouped data files on partition keys. - * - *

3. Filter results to find dangling delete files by comparing delete file sequence_number to - * its partitions' minimum data sequence number. - * - *

4. Collect results row to driver and use {@link SparkDeleteFile SparkDeleteFile} to wrap - * rows to valid delete files + *

    + *
  1. Group data files by partition keys and find the minimum data sequence number in each + * group. + *
  2. Left outer join delete files with partition-grouped data files on partition keys. + *
  3. Find dangling deletes by comparing each delete file's sequence number to its partition's + * minimum data sequence number. + *
  4. Collect results row to driver and use {@link SparkDeleteFile SparkDeleteFile} to wrap + * rows to valid delete files + *
*/ private List findDanglingDeletes() { Dataset minSequenceNumberByPartition = loadMetadataTable(table, MetadataTableType.ENTRIES) - // find live entry for data files + // find live data files .filter("data_file.content == 0 AND status < 2") .selectExpr( "data_file.partition as partition", @@ -132,8 +131,8 @@ private List findDanglingDeletes() { Dataset deleteEntries = loadMetadataTable(table, MetadataTableType.ENTRIES) - // find live entry for delete files - .filter(" data_file.content != 0 AND status < 2"); + // find live delete files + .filter("data_file.content != 0 AND status < 2"); Column joinOnPartition = deleteEntries @@ -146,14 +145,14 @@ private List findDanglingDeletes() { Column filterOnDanglingDeletes = col("min_data_sequence_number") - // when all data files are rewritten with new partition spec + // delete fies without any data files in partition .isNull() - // dangling position delete files + // position delete files without any applicable data files in partition .or( col("data_file.content") .equalTo("1") .and(col("sequence_number").$less(col("min_data_sequence_number")))) - // dangling equality delete files + // equality delete files without any applicable data files in the partition .or( col("data_file.content") .equalTo("2") @@ -165,6 +164,7 @@ private List findDanglingDeletes() { .filter(filterOnDanglingDeletes) .select("data_file.*"); return danglingDeletes.collectAsList().stream() + // map on driver because SparkDeleteFile is not serializable .map(row -> deleteFileWrapper(danglingDeletes.schema(), row)) .collect(Collectors.toList()); } diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java index 4c0406a76489..4e381a7bd362 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/RewriteDataFilesSparkAction.java @@ -84,8 +84,8 @@ public class RewriteDataFilesSparkAction PARTIAL_PROGRESS_MAX_FAILED_COMMITS, TARGET_FILE_SIZE_BYTES, USE_STARTING_SEQUENCE_NUMBER, - OUTPUT_SPEC_ID, REWRITE_JOB_ORDER, + OUTPUT_SPEC_ID, REMOVE_DANGLING_DELETES); private static final RewriteDataFilesSparkAction.Result EMPTY_RESULT = @@ -179,13 +179,10 @@ public RewriteDataFiles.Result execute() { Stream groupStream = toGroupStream(ctx, fileGroupsByPartition); - Builder resultBuilder; - if (partialProgressEnabled) { - resultBuilder = - doExecuteWithPartialProgress(ctx, groupStream, commitManager(startingSnapshotId)); - } else { - resultBuilder = doExecute(ctx, groupStream, commitManager(startingSnapshotId)); - } + Builder resultBuilder = + partialProgressEnabled + ? doExecuteWithPartialProgress(ctx, groupStream, commitManager(startingSnapshotId)) + : doExecute(ctx, groupStream, commitManager(startingSnapshotId)); if (removeDanglingDeletes) { RemoveDanglingDeletesSparkAction action = From a84cf88eccfc914a48d12d897b508f3d251b4c53 Mon Sep 17 00:00:00 2001 From: Hongyue Zhang Date: Mon, 14 Oct 2024 15:18:59 -0700 Subject: [PATCH 5/5] spotlessApply to reformat spaces --- .../java/org/apache/iceberg/actions/ActionsProvider.java | 2 +- .../org/apache/iceberg/spark/actions/SparkActions.java | 7 +++---- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java b/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java index 23ed92baf875..e033c632e555 100644 --- a/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java +++ b/api/src/main/java/org/apache/iceberg/actions/ActionsProvider.java @@ -76,7 +76,7 @@ default ComputeTableStats computeTableStats(Table table) { throw new UnsupportedOperationException( this.getClass().getName() + " does not implement computeTableStats"); } - + /** Instantiates an action to remove dangling delete files from current snapshot. */ default RemoveDanglingDeleteFiles removeDanglingDeleteFiles(Table table) { throw new UnsupportedOperationException( diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java index 1ba53105228c..ba9fa2e7b4db 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/actions/SparkActions.java @@ -98,15 +98,14 @@ public DeleteReachableFilesSparkAction deleteReachableFiles(String metadataLocat public RewritePositionDeleteFilesSparkAction rewritePositionDeletes(Table table) { return new RewritePositionDeleteFilesSparkAction(spark, table); } - + @Override public ComputeTableStats computeTableStats(Table table) { return new ComputeTableStatsSparkAction(spark, table); } - @Override + @Override public RemoveDanglingDeleteFiles removeDanglingDeleteFiles(Table table) { - return new RemoveDanglingDeletesSparkAction(spark, table); + return new RemoveDanglingDeletesSparkAction(spark, table); } - }