szehon-ho commented on code in PR #17622: URL: https://github.com/apache/iceberg/pull/17622#discussion_r4150782798
########## spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/RepairMetrics.java: ########## @@ -0,0 +1,262 @@ +/* + * 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.TableProperties.DEFAULT_NAME_MAPPING; + +import java.nio.ByteBuffer; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import org.apache.iceberg.ContentFile; +import org.apache.iceberg.DataFile; +import org.apache.iceberg.DataFiles; +import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.FileContent; +import org.apache.iceberg.FileFormat; +import org.apache.iceberg.FileMetadata; +import org.apache.iceberg.MetadataColumns; +import org.apache.iceberg.Metrics; +import org.apache.iceberg.MetricsConfig; +import org.apache.iceberg.MetricsUtil; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.Schema; +import org.apache.iceberg.Table; +import org.apache.iceberg.avro.Avro; +import org.apache.iceberg.io.InputFile; +import org.apache.iceberg.mapping.NameMapping; +import org.apache.iceberg.mapping.NameMappingParser; +import org.apache.iceberg.orc.OrcMetrics; +import org.apache.iceberg.parquet.ParquetUtil; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.relocated.com.google.common.collect.Maps; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.TypeUtil; + +/** + * Reads the statistics of data and delete files and compares them against the statistics recorded + * in manifest entries. + * + * <p>Recomputed statistics always respect the metrics config of the table so that they are + * comparable with the stored statistics. + */ +class RepairMetrics { + + private static final Set<Integer> POSITION_DELETE_FIELD_IDS = + ImmutableSet.of( + MetadataColumns.DELETE_FILE_PATH.fieldId(), MetadataColumns.DELETE_FILE_POS.fieldId()); + + private RepairMetrics() {} + + /** Returns the name mapping of the table, or null if the table does not define one. */ + static NameMapping nameMapping(Table table) { + String mapping = table.properties().get(DEFAULT_NAME_MAPPING); + return mapping != null ? NameMappingParser.fromJson(mapping) : null; + } + + /** + * Returns the metrics config to use when recomputing the statistics of the given file. + * + * <p>Position delete files record statistics for the path and position columns only, which is a + * fixed config rather than the config of the table. + */ + static MetricsConfig metricsConfig(Table table, FileContent content) { + return content == FileContent.POSITION_DELETES + ? MetricsConfig.forPositionDelete() + : MetricsConfig.forTable(table); + } + + /** + * Returns true if the statistics of the file can be recomputed by reading it. + * + * <p>Deletion vectors are stored as blobs inside a Puffin file, so their statistics cannot be + * derived by reading the file they are stored in. + */ + static boolean supportsMetrics(ContentFile<?> file) { + FileFormat format = file.format(); + return format == FileFormat.PARQUET || format == FileFormat.ORC || format == FileFormat.AVRO; + } + + /** Recomputes recoverable statistics, preserving metrics that require writer-side tracking. */ + static Metrics readMetrics( + InputFile input, + ContentFile<?> file, + MetricsConfig config, + NameMapping mapping, + Schema schema) { + Metrics metrics = + switch (file.format()) { + case PARQUET -> ParquetUtil.fileMetrics(input, config, mapping); + case ORC -> OrcMetrics.fromInputFile(input, config, mapping); + case AVRO -> new Metrics(Avro.rowCount(input)); + default -> + throw new UnsupportedOperationException( + "Cannot read metrics of format: " + file.format()); + }; + if (file.format() == FileFormat.AVRO) { + return recordCountOnly(file, metrics); + } + + if (file.content() == FileContent.POSITION_DELETES) { + // Match PositionDeleteWriter: counts are omitted, and bounds are kept only for a single path. + int pathId = MetadataColumns.DELETE_FILE_PATH.fieldId(); + ByteBuffer lowerPath = normalize(metrics.lowerBounds()).get(pathId); + ByteBuffer upperPath = normalize(metrics.upperBounds()).get(pathId); + metrics = + lowerPath != null && lowerPath.equals(upperPath) + ? MetricsUtil.copyWithoutFieldCounts(metrics, POSITION_DELETE_FIELD_IDS) + : MetricsUtil.copyWithoutFieldCountsAndBounds(metrics, POSITION_DELETE_FIELD_IDS); + } + + // Footers cannot recover NaN counts or the NaN-safe bounds tracked by writers. Include + // stored NaN-count IDs so that bounds for dropped columns are preserved as well. + Set<Integer> floatingPointIds = Sets.newHashSet(normalize(file.nanValueCounts()).keySet()); + TypeUtil.indexById(schema.asStruct()) + .forEach( + (id, field) -> { + Type.TypeID typeId = field.type().typeId(); + if (typeId == Type.TypeID.FLOAT || typeId == Type.TypeID.DOUBLE) { Review Comment: Preserve geometry bounds alongside floating-point bounds. `GenericParquetWriter` collects geometry bounding boxes, but `ParquetMetrics.metricsFromFooter` intentionally returns only counts for geometry. With `repair-column-metrics=true`, a correct file therefore appears corrupt, and the replacement entry loses its valid geometry bounds. Add a no-op test using a geometry file written by `GenericParquetWriter`. ########## spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/RepairTableSparkAction.java: ########## @@ -0,0 +1,990 @@ +/* + * 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.MetadataTableType.ENTRIES; + +import java.io.Serializable; +import java.util.EnumMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.function.Function; +import java.util.stream.Collectors; +import org.apache.hadoop.fs.Path; +import org.apache.iceberg.ContentFile; +import org.apache.iceberg.DataFile; +import org.apache.iceberg.DataFiles; +import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.FileContent; +import org.apache.iceberg.FileFormat; +import org.apache.iceberg.FileMetadata; +import org.apache.iceberg.GenericManifestFile; +import org.apache.iceberg.HasTableOperations; +import org.apache.iceberg.ManifestContent; +import org.apache.iceberg.ManifestFile; +import org.apache.iceberg.ManifestFiles; +import org.apache.iceberg.ManifestReader; +import org.apache.iceberg.ManifestWriter; +import org.apache.iceberg.Metrics; +import org.apache.iceberg.MetricsConfig; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.Partitioning; +import org.apache.iceberg.RewriteManifests; +import org.apache.iceberg.RollingManifestWriter; +import org.apache.iceberg.Snapshot; +import org.apache.iceberg.SnapshotSummary; +import org.apache.iceberg.Table; +import org.apache.iceberg.TableOperations; +import org.apache.iceberg.TableProperties; +import org.apache.iceberg.actions.ImmutableRepairTable; +import org.apache.iceberg.actions.RepairTable; +import org.apache.iceberg.encryption.EncryptedOutputFile; +import org.apache.iceberg.encryption.EncryptingFileIO; +import org.apache.iceberg.exceptions.CleanableFailure; +import org.apache.iceberg.exceptions.CommitStateUnknownException; +import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.InputFile; +import org.apache.iceberg.io.OutputFile; +import org.apache.iceberg.io.SupportsBulkOperations; +import org.apache.iceberg.mapping.NameMapping; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +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.Sets; +import org.apache.iceberg.spark.JobGroupInfo; +import org.apache.iceberg.spark.SparkContentFile; +import org.apache.iceberg.spark.SparkDataFile; +import org.apache.iceberg.spark.SparkDeleteFile; +import org.apache.iceberg.spark.source.SerializableTableWithSize; +import org.apache.iceberg.types.Types; +import org.apache.iceberg.util.PropertyUtil; +import org.apache.iceberg.util.ScanTaskUtil; +import org.apache.iceberg.util.ThreadPools; +import org.apache.spark.api.java.function.MapPartitionsFunction; +import org.apache.spark.broadcast.Broadcast; +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Encoder; +import org.apache.spark.sql.Encoders; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.functions; +import org.apache.spark.sql.types.StructType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import scala.Tuple2; + +/** + * An action that repairs incorrect statistics in the manifests of a table. + * + * <p>The statistics of every live manifest entry are compared against the file the entry refers to. + * Only manifests that contain at least one incorrect entry are rewritten. Snapshot totals are + * recomputed from the live entries when committing a repair. + * + * <p>Deletion vectors (delete blobs stored in Puffin files) are not verified or repaired. NaN + * counts and floating-point bounds are preserved because footers cannot reconstruct them reliably. + */ +public class RepairTableSparkAction extends BaseSnapshotUpdateSparkAction<RepairTableSparkAction> + implements RepairTable { + + public static final String USE_CACHING = "use-caching"; + public static final boolean USE_CACHING_DEFAULT = false; + + /** + * Whether to compare and repair column level statistics. When disabled, only record counts and + * file sizes are compared and repaired. + * + * <p>This is disabled by default. Recomputed column statistics reflect the current metrics config + * of the table, but the config a file was written under is not recorded, so a table whose config + * changed reports column statistics that legitimately differ from the recomputed ones. Repairing + * them in that case would overwrite correct statistics. Reading the footer of every candidate + * file happens regardless of this option; it only controls whether column statistics are + * compared. + */ + public static final String REPAIR_COLUMN_METRICS = "repair-column-metrics"; + + public static final boolean REPAIR_COLUMN_METRICS_DEFAULT = false; + + private static final Logger LOG = LoggerFactory.getLogger(RepairTableSparkAction.class); + + private static final RepairTable.Result EMPTY_RESULT = + ImmutableRepairTable.Result.builder() + .repairedManifests(ImmutableList.of()) + .repairedEntryCount(0L) + .build(); + + private static final String NEW_MANIFEST_PREFIX = "repaired-m-"; + + private final Table table; + private final int formatVersion; + private final long targetManifestSizeBytes; + private final boolean shouldStageManifests; + private final String outputLocation; + + private boolean repairFileMetrics = false; + private boolean dryRun = false; + + RepairTableSparkAction(SparkSession spark, Table table) { + super(spark); + this.table = table; + this.targetManifestSizeBytes = + PropertyUtil.propertyAsLong( + table.properties(), + TableProperties.MANIFEST_TARGET_SIZE_BYTES, + TableProperties.MANIFEST_TARGET_SIZE_BYTES_DEFAULT); + + TableOperations ops = ((HasTableOperations) table).operations(); + Path metadataFilePath = new Path(ops.metadataFileLocation("file")); + this.outputLocation = metadataFilePath.getParent().toString(); + this.formatVersion = ops.current().formatVersion(); + + boolean snapshotIdInheritanceEnabled = + PropertyUtil.propertyAsBoolean( + table.properties(), + TableProperties.SNAPSHOT_ID_INHERITANCE_ENABLED, + TableProperties.SNAPSHOT_ID_INHERITANCE_ENABLED_DEFAULT); + this.shouldStageManifests = formatVersion == 1 && !snapshotIdInheritanceEnabled; + } + + @Override + protected RepairTableSparkAction self() { + return this; + } + + @Override + public RepairTableSparkAction repairFileMetrics() { + this.repairFileMetrics = true; + return this; + } + + @Override + public RepairTableSparkAction dryRun() { + this.dryRun = true; + return this; + } + + @Override + public RepairTable.Result execute() { + if (!repairFileMetrics) { + return EMPTY_RESULT; + } + + String desc = String.format("Repairing manifests in %s (dryRun=%s)", table.name(), dryRun); + JobGroupInfo info = newJobGroupInfo("REPAIR-TABLE", desc); + return withJobGroupInfo(info, this::doExecute); + } + + private RepairTable.Result doExecute() { + Snapshot currentSnapshot = table.currentSnapshot(); + if (currentSnapshot == null) { + return EMPTY_RESULT; + } + + List<ManifestFile> repairedManifests = Lists.newArrayList(); + List<ManifestFile> newManifests = Lists.newArrayList(); + long repairedCount = 0L; + + try { + for (ManifestContent content : ManifestContent.values()) { + RepairedManifests repaired = repairTable(content, currentSnapshot); + repairedManifests.addAll(repaired.repairedManifests()); + newManifests.addAll(repaired.newManifests()); + repairedCount += repaired.repairedCount(); + } + } catch (Exception e) { + // If a later content group fails before the commit, delete the manifests already written by + // earlier groups so a partial repair does not leave orphan files in the metadata directory. + deleteFiles(Iterables.transform(newManifests, ManifestFile::path)); + throw e; + } + + if (repairedManifests.isEmpty()) { + return EMPTY_RESULT; + } + + // a dry run writes no manifests, so there is nothing to commit or clean up + if (!dryRun) { + replaceManifests(repairedManifests, newManifests); + } + + LOG.info( + "Repaired the stats of {} manifest entries, rewriting {} manifests as {} (dryRun={})", + repairedCount, + repairedManifests.size(), + newManifests.size(), + dryRun); + + return ImmutableRepairTable.Result.builder() + .repairedManifests(repairedManifests) + .repairedEntryCount(repairedCount) + .build(); + } + + private RepairedManifests repairTable(ManifestContent content, Snapshot snapshot) { + List<ManifestFile> manifests = loadManifests(content, snapshot); + if (manifests.isEmpty()) { + return RepairedManifests.empty(); + } + + // A manifest is rewritten with the spec it was written under, so manifests are grouped by + // spec and each group is repaired separately. Rewriting a manifest of an older spec with the + // current spec of the table would change the partition data of its entries. + Map<Integer, List<ManifestFile>> manifestsBySpecId = + manifests.stream().collect(Collectors.groupingBy(ManifestFile::partitionSpecId)); + + List<ManifestFile> repairedManifests = Lists.newArrayList(); + List<ManifestFile> newManifests = Lists.newArrayList(); + long repairedCount = 0L; + + try { + for (Map.Entry<Integer, List<ManifestFile>> group : manifestsBySpecId.entrySet()) { + RepairedManifests repaired = repairManifests(content, group.getKey(), group.getValue()); + repairedManifests.addAll(repaired.repairedManifests()); + newManifests.addAll(repaired.newManifests()); + repairedCount += repaired.repairedCount(); + } + } catch (Exception e) { + // If a later spec group fails before the commit, delete the manifests already written by + // earlier groups so a partial repair does not leave orphan files behind. + deleteFiles(Iterables.transform(newManifests, ManifestFile::path)); + throw e; + } + + return RepairedManifests.of(repairedManifests, newManifests, repairedCount); + } + + private RepairedManifests repairManifests( + ManifestContent content, int specId, List<ManifestFile> manifests) { + Dataset<Row> entryDF = buildManifestEntryDF(manifests); + + return withReusableDS( + entryDF, + df -> { + // the entries whose stats disagree with the files they refer to, as (manifest, path). + // cached because it is small, one row per incorrect entry, and is read by several actions + // below, whereas recomputing it would re-read every file. + Dataset<Row> verdicts = + df.mapPartitions( + newCheckStatsFunc(content, specId), + Encoders.tuple(Encoders.STRING(), Encoders.STRING())) + .toDF("manifest", "path") + .cache(); + + try { + Set<String> manifestsToRewrite = + Sets.newHashSet( + verdicts.select("manifest").distinct().as(Encoders.STRING()).collectAsList()); + + if (manifestsToRewrite.isEmpty()) { + return RepairedManifests.empty(); + } + + long repairedCount = verdicts.count(); + List<ManifestFile> rewritten = + manifests.stream() + .filter(manifest -> manifestsToRewrite.contains(manifest.path())) + .collect(Collectors.toList()); + + // a dry run reports what would be repaired without writing any manifests + if (dryRun) { + return RepairedManifests.of(rewritten, ImmutableList.of(), repairedCount); + } + + // mark every entry of the affected manifests with whether its stats need repair by + // joining on the file path, so the unbounded per file set stays distributed rather than + // being collected to the driver and broadcast back out + Dataset<Row> entriesToRewrite = + df.filter(df.col("manifest").isin(manifestsToRewrite.toArray())); + Dataset<Row> markedEntries = markEntriesToRepair(entriesToRewrite, verdicts); + List<ManifestFile> written = + writeManifests(content, specId, markedEntries, rewritten.size()); + + return RepairedManifests.of(rewritten, written, repairedCount); + } finally { + verdicts.unpersist(false); + } + }); + } + + /** + * Marks every entry with a boolean {@code repair} column that is true when the entry's file has + * incorrect statistics, by left joining the entries against the verdicts on the file path. + */ + private Dataset<Row> markEntriesToRepair(Dataset<Row> entries, Dataset<Row> verdicts) { + Dataset<Row> repairedPaths = + verdicts.select(verdicts.col("path").as("repaired_path")).distinct(); + return entries + .join( + repairedPaths, + entries.col("data_file.file_path").equalTo(repairedPaths.col("repaired_path")), + "left") + .withColumn("repair", functions.col("repaired_path").isNotNull()) + .drop("repaired_path") + .select( + "manifest", + "snapshot_id", + "sequence_number", + "file_sequence_number", + "data_file", + "repair"); + } + + /** + * Loads the live entries of the given manifests, keeping the manifest each entry was read from so + * that only the manifests containing an incorrect entry are rewritten. + */ + private Dataset<Row> buildManifestEntryDF(List<ManifestFile> manifests) { + Dataset<Row> manifestDF = + spark() + .createDataset(Lists.transform(manifests, ManifestFile::path), Encoders.STRING()) + .toDF("manifest"); + + Dataset<Row> entryDF = + loadMetadataTable(table, ENTRIES) + .filter("status < 2") // select only live entries + .selectExpr( + "input_file_name() as manifest", + "snapshot_id", + "sequence_number", + "file_sequence_number", + "data_file"); + + return entryDF.join( + manifestDF, manifestDF.col("manifest").equalTo(entryDF.col("manifest")), "left_semi"); + } + + private List<ManifestFile> writeManifests( + ManifestContent content, int specId, Dataset<Row> entryDF, int numManifests) { + StructType sparkType = (StructType) entryDF.schema().apply("data_file").dataType(); + Types.StructType combinedFileType = DataFile.getType(Partitioning.partitionType(table)); + Types.StructType fileType = DataFile.getType(table.specs().get(specId).partitionType()); + ManifestWriterFactory writers = manifestWriters(specId); + RepairContext context = newRepairContext(content, specId); + + WriteManifests<?> writeFunc = + content == ManifestContent.DATA + ? new WriteDataManifests(writers, combinedFileType, fileType, sparkType, context) + : new WriteDeleteManifests(writers, combinedFileType, fileType, sparkType, context); + + // repartition by manifest so the entries of each manifest are written together and the layout + // of the table is preserved, rather than scattered round robin as a plain repartition(n) would. + // this produces about as many manifests as are being replaced. + return writeFunc + .apply(entryDF.repartition(numManifests, entryDF.col("manifest"))) + .collectAsList(); + } + + private CheckStats newCheckStatsFunc(ManifestContent content, int specId) { + return new CheckStats(newRepairContext(content, specId)); + } + + private RepairContext newRepairContext(ManifestContent content, int specId) { + boolean repairColumnMetrics = + PropertyUtil.propertyAsBoolean( + options(), REPAIR_COLUMN_METRICS, REPAIR_COLUMN_METRICS_DEFAULT); + return new RepairContext( + sparkContext().broadcast(SerializableTableWithSize.copyOf(table)), + content, + specId, + repairColumnMetrics); + } + + private List<ManifestFile> loadManifests(ManifestContent content, Snapshot snapshot) { + switch (content) { + case DATA: + return snapshot.dataManifests(table.io()); + case DELETES: + return snapshot.deleteManifests(table.io()); + default: + throw new IllegalArgumentException("Unknown manifest content: " + content); + } + } + + private void replaceManifests( + List<ManifestFile> deletedManifests, List<ManifestFile> addedManifests) { + try { + RewriteManifests rewriteManifests = table.rewriteManifests(); + deletedManifests.forEach(rewriteManifests::deleteManifest); + addedManifests.forEach(rewriteManifests::addManifest); + Set<String> deletedPaths = + deletedManifests.stream().map(ManifestFile::path).collect(Collectors.toSet()); + // The validator receives the refreshed parent on every commit attempt, including retries. + // Recompute totals against that parent so concurrent changes are included in the summary. + rewriteManifests.validateWith( + snapshots -> { + Snapshot parent = Iterables.getFirst(snapshots, null); + if (parent == null) { + return false; + } + + List<ManifestFile> manifests = + parent.allManifests(table.io()).stream() + .filter(manifest -> !deletedPaths.contains(manifest.path())) + .collect(Collectors.toList()); + manifests.addAll(addedManifests); + updateSnapshotTotals(rewriteManifests, manifests); Review Comment: Handle failures from `updateSnapshotTotals` as cleanable before propagating them. Its Spark job runs during validation, before the metadata commit, but a `SparkException` bypasses the `CleanableFailure`-only cleanup below. `BaseRewriteManifests` also deliberately leaves caller-supplied manifests untouched, so all replacement manifests already written become orphans. Add coverage that fails the totals job after manifest writing completes. -- 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]
