anoopj commented on code in PR #18200: URL: https://github.com/apache/iceberg/pull/18200#discussion_r4067176058
########## core/src/main/java/org/apache/iceberg/SnapshotValidator.java: ########## @@ -0,0 +1,589 @@ +/* + * 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; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import java.util.concurrent.ExecutorService; +import java.util.stream.Collectors; +import org.apache.iceberg.exceptions.ValidationException; +import org.apache.iceberg.expressions.Expression; +import org.apache.iceberg.expressions.Expressions; +import org.apache.iceberg.expressions.ManifestEvaluator; +import org.apache.iceberg.expressions.Projections; +import org.apache.iceberg.io.CloseableIterable; +import org.apache.iceberg.io.CloseableIterator; +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.Iterators; +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.Sets; +import org.apache.iceberg.util.CharSequenceSet; +import org.apache.iceberg.util.ContentFileUtil; +import org.apache.iceberg.util.Pair; +import org.apache.iceberg.util.PartitionSet; +import org.apache.iceberg.util.SnapshotUtil; +import org.apache.iceberg.util.Tasks; + +/** + * Checks a commit for conflicts against the table history between a starting snapshot and the + * current parent snapshot. + */ +class SnapshotValidator { + // data is only added in "append" and "overwrite" operations + private static final Set<String> VALIDATE_ADDED_FILES_OPERATIONS = + ImmutableSet.of(DataOperations.APPEND, DataOperations.OVERWRITE); + // data files are removed in "overwrite", "replace", and "delete" + private static final Set<String> VALIDATE_DATA_FILES_EXIST_OPERATIONS = + ImmutableSet.of(DataOperations.OVERWRITE, DataOperations.REPLACE, DataOperations.DELETE); + private static final Set<String> VALIDATE_DATA_FILES_EXIST_SKIP_DELETE_OPERATIONS = + ImmutableSet.of(DataOperations.OVERWRITE, DataOperations.REPLACE); + // delete files can be added in "overwrite" or "delete" operations + private static final Set<String> VALIDATE_ADDED_DELETE_FILES_OPERATIONS = + ImmutableSet.of(DataOperations.OVERWRITE, DataOperations.DELETE); + // DVs can be added in "overwrite", "delete", and "replace" operations + private static final Set<String> VALIDATE_ADDED_DVS_OPERATIONS = + ImmutableSet.of(DataOperations.OVERWRITE, DataOperations.DELETE, DataOperations.REPLACE); + + private final TableOperations ops; + private final boolean caseSensitive; + + SnapshotValidator(TableOperations ops, boolean caseSensitive) { + this.ops = ops; + this.caseSensitive = caseSensitive; + } + + void validateAddedDataFiles( + TableMetadata base, Long startingSnapshotId, PartitionSet partitionSet, Snapshot parent) { + CloseableIterable<ManifestEntry<DataFile>> conflictEntries = + addedDataFiles(base, startingSnapshotId, null, partitionSet, parent); + + try (CloseableIterator<ManifestEntry<DataFile>> conflicts = conflictEntries.iterator()) { + if (conflicts.hasNext()) { + throw new ValidationException( + "Found conflicting files that can contain records matching partitions %s: %s", + partitionSet, + Iterators.toString( + Iterators.transform(conflicts, entry -> entry.file().location().toString()))); + } + + } catch (IOException e) { + throw new UncheckedIOException( + String.format("Failed to validate no appends matching %s", partitionSet), e); + } + } + + void validateAddedDataFiles( + TableMetadata base, + Long startingSnapshotId, + Expression conflictDetectionFilter, + Snapshot parent) { + CloseableIterable<ManifestEntry<DataFile>> conflictEntries = + addedDataFiles(base, startingSnapshotId, conflictDetectionFilter, null, parent); + + try (CloseableIterator<ManifestEntry<DataFile>> conflicts = conflictEntries.iterator()) { + if (conflicts.hasNext()) { + throw new ValidationException( + "Found conflicting files that can contain records matching %s: %s", + conflictDetectionFilter, + Iterators.toString( + Iterators.transform(conflicts, entry -> entry.file().location().toString()))); + } + + } catch (IOException e) { + throw new UncheckedIOException( + String.format("Failed to validate no appends matching %s", conflictDetectionFilter), e); + } + } + + private CloseableIterable<ManifestEntry<DataFile>> addedDataFiles( + TableMetadata base, + Long startingSnapshotId, + Expression dataFilter, + PartitionSet partitionSet, + Snapshot parent) { + // if there is no current table state, no files have been added + if (parent == null) { + return CloseableIterable.empty(); + } + + Pair<List<ManifestFile>, Set<Long>> history = + validationHistory( + base, + startingSnapshotId, + VALIDATE_ADDED_FILES_OPERATIONS, + ManifestContent.DATA, + parent); + List<ManifestFile> manifests = history.first(); + Set<Long> newSnapshots = history.second(); + + ManifestGroup manifestGroup = + new ManifestGroup(ops.io(), manifests, ImmutableList.of()) + .caseSensitive(caseSensitive) + .filterManifestEntries(entry -> newSnapshots.contains(entry.snapshotId())) + .specsById(base.specsById()) + .ignoreDeleted() + .ignoreExisting(); + + if (dataFilter != null) { + manifestGroup = manifestGroup.filterData(dataFilter); + } + + if (partitionSet != null) { + manifestGroup = + manifestGroup.filterManifestEntries( + entry -> partitionSet.contains(entry.file().specId(), entry.file().partition())); + } + + return manifestGroup.entries(); + } + + void validateNoNewDeletesForDataFiles( + TableMetadata base, + Long startingSnapshotId, + Expression dataFilter, + Iterable<DataFile> dataFiles, + boolean ignoreEqualityDeletes, + Snapshot parent) { + // if there is no current table state, no files have been added + if (parent == null || base.formatVersion() < 2) { + return; + } + + List<DeleteFileIndex> deleteIndexes = + addedDeleteFilesIndexedPerSnapshot(base, startingSnapshotId, dataFilter, null, parent); + + long startingSequenceNumber = startingSequenceNumber(base, startingSnapshotId); + for (DataFile dataFile : dataFiles) { + for (DeleteFileIndex deletes : deleteIndexes) { + // if any delete is found that applies to files written in or before the starting snapshot, + // fail + DeleteFile[] deleteFiles = deletes.forDataFile(startingSequenceNumber, dataFile); + // rewrites can omit equality-delete checks: when added files keep the replaced files' data + // sequence number, higher-sequence equality deletes still apply to them, so there is no + // RewriteFiles/RowDelta conflict; only a new position delete signals a real conflict + if (ignoreEqualityDeletes) { + ValidationException.check( + !containsPositionDeletes(deleteFiles), + "Cannot commit, found new position delete for replaced data file: %s", + dataFile); + } else { + ValidationException.check( + deleteFiles.length == 0, + "Cannot commit, found new delete for replaced data file: %s", + dataFile); + } + } + } + } + + private static boolean containsPositionDeletes(DeleteFile[] deleteFiles) { + for (DeleteFile deleteFile : deleteFiles) { + if (deleteFile.content() == FileContent.POSITION_DELETES) { + return true; + } + } + + return false; + } + + void validateNoNewDeleteFiles( + TableMetadata base, Long startingSnapshotId, Expression dataFilter, Snapshot parent) { + Set<String> locations = + referencedDeleteFileLocations( + addedDeleteFilesIndexedPerSnapshot(base, startingSnapshotId, dataFilter, null, parent)); + ValidationException.check( + locations.isEmpty(), + "Found new conflicting delete files that can apply to records matching %s: %s", + dataFilter, + locations); + } + + void validateNoNewDeleteFiles( + TableMetadata base, Long startingSnapshotId, PartitionSet partitionSet, Snapshot parent) { + Set<String> locations = + referencedDeleteFileLocations( + addedDeleteFilesIndexedPerSnapshot( + base, startingSnapshotId, null, partitionSet, parent)); + ValidationException.check( + locations.isEmpty(), + "Found new conflicting delete files that can apply to records matching %s: %s", + partitionSet, + locations); + } + + private static Set<String> referencedDeleteFileLocations(List<DeleteFileIndex> deleteIndexes) { + Set<String> locations = Sets.newLinkedHashSet(); + for (DeleteFileIndex deletes : deleteIndexes) { + for (DeleteFile deleteFile : deletes.referencedDeleteFiles()) { + locations.add(deleteFile.location()); + } + } + + return locations; + } + + private List<DeleteFileIndex> addedDeleteFilesIndexedPerSnapshot( + TableMetadata base, + Long startingSnapshotId, + Expression dataFilter, + PartitionSet partitionSet, + Snapshot parent) { + // if there is no current table state, no delete files have been added + if (parent == null || base.formatVersion() < 2) { + return ImmutableList.of(); + } + + Pair<List<ManifestFile>, Set<Long>> history = + validationHistory( + base, + startingSnapshotId, + VALIDATE_ADDED_DELETE_FILES_OPERATIONS, + ManifestContent.DELETES, + parent); + + // the history collects a manifest only from the snapshot that added it, so grouping by Review Comment: Left all of these comments intact. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
