szehon-ho commented on code in PR #17622:
URL: https://github.com/apache/iceberg/pull/17622#discussion_r4138905664


##########
spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/RepairMetrics.java:
##########
@@ -0,0 +1,195 @@
+/*
+ * 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 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.Metrics;
+import org.apache.iceberg.MetricsConfig;
+import org.apache.iceberg.PartitionSpec;
+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;
+
+/**
+ * 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 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 the statistics of a file by reading it. */
+  static Metrics readMetrics(
+      InputFile input, ContentFile<?> file, MetricsConfig config, NameMapping 
mapping) {
+    switch (file.format()) {
+      case PARQUET:
+        return ParquetUtil.fileMetrics(input, config, mapping);

Review Comment:
   Resolve per-column metrics settings by field ID when reading files written 
before a column rename. For example, write a with default=none and a=full, then 
rename a to b. SchemaUpdate moves the override to b, but the footer still names 
the field a and MetricsUtil.metricsMode looks up that old name, falling back to 
none. With column repair enabled, this drops valid statistics despite the 
intended configuration staying the same. Please add a no-op test for this 
rename case.



##########
spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/RepairTableSparkAction.java:
##########
@@ -0,0 +1,856 @@
+/*
+ * 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.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.HasTableOperations;
+import org.apache.iceberg.ManifestContent;
+import org.apache.iceberg.ManifestFile;
+import org.apache.iceberg.ManifestFiles;
+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.RollingManifestWriter;
+import org.apache.iceberg.Snapshot;
+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.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.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, so 
the cost of the commit
+ * is proportional to the number of incorrect entries rather than to the size 
of the table.
+ */
+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() {
+    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() {
+    if (!repairFileMetrics) {
+      // no repair was selected through the configuration methods, so there is 
nothing to do
+      return EMPTY_RESULT;
+    }
+
+    Snapshot currentSnapshot = table.currentSnapshot();
+    if (currentSnapshot == null) {
+      return EMPTY_RESULT;
+    }
+
+    List<ManifestFile> repairedManifests = Lists.newArrayList();
+    List<ManifestFile> newManifests = Lists.newArrayList();
+    long repairedCount = 0L;
+
+    for (ManifestContent content : ManifestContent.values()) {
+      RepairedManifests repaired = repairTable(content, currentSnapshot);
+      repairedManifests.addAll(repaired.repairedManifests());
+      newManifests.addAll(repaired.newManifests());
+      repairedCount += repaired.repairedCount();
+    }
+
+    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;
+
+    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();
+    }
+
+    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 {
+            List<String> manifestsToRewrite =
+                
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()))

Review Comment:
   Use a Set for these path lookups. manifestsToRewrite is a List, so filtering 
N manifests requires roughly N squared / 2 string comparisons on the driver 
when all are affected.



##########
spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/RepairTableSparkAction.java:
##########
@@ -0,0 +1,856 @@
+/*
+ * 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.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.HasTableOperations;
+import org.apache.iceberg.ManifestContent;
+import org.apache.iceberg.ManifestFile;
+import org.apache.iceberg.ManifestFiles;
+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.RollingManifestWriter;
+import org.apache.iceberg.Snapshot;
+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.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.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, so 
the cost of the commit
+ * is proportional to the number of incorrect entries rather than to the size 
of the table.
+ */
+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() {
+    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() {
+    if (!repairFileMetrics) {
+      // no repair was selected through the configuration methods, so there is 
nothing to do
+      return EMPTY_RESULT;
+    }
+
+    Snapshot currentSnapshot = table.currentSnapshot();
+    if (currentSnapshot == null) {
+      return EMPTY_RESULT;
+    }
+
+    List<ManifestFile> repairedManifests = Lists.newArrayList();
+    List<ManifestFile> newManifests = Lists.newArrayList();
+    long repairedCount = 0L;
+
+    for (ManifestContent content : ManifestContent.values()) {
+      RepairedManifests repaired = repairTable(content, currentSnapshot);
+      repairedManifests.addAll(repaired.repairedManifests());
+      newManifests.addAll(repaired.newManifests());
+      repairedCount += repaired.repairedCount();
+    }
+
+    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;
+
+    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();
+    }
+
+    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 {
+            List<String> manifestsToRewrite =
+                
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(
+      Iterable<ManifestFile> deletedManifests, Iterable<ManifestFile> 
addedManifests) {
+    try {
+      org.apache.iceberg.RewriteManifests rewriteManifests = 
table.rewriteManifests();
+      deletedManifests.forEach(rewriteManifests::deleteManifest);
+      addedManifests.forEach(rewriteManifests::addManifest);
+      commit(rewriteManifests);

Review Comment:
   Recompute snapshot totals as part of the repair commit. RewriteManifests 
carries forward the previous totals, so a four-row file initially appended with 
record_count=104 is repaired to 4 here while total-records stays 104. Spark 
still uses that total for unfiltered partitioned-table scan estimates. The 
file-size and delete-count totals have the same issue. Please add a test that 
initially commits incorrect metrics and checks the repaired snapshot summary.



##########
spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRepairTableAction.java:
##########
@@ -0,0 +1,1085 @@
+/*
+ * 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 static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.assertj.core.api.Assumptions.assumeThat;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.when;
+
+import java.io.File;
+import java.io.IOException;
+import java.nio.file.Path;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.hadoop.conf.Configuration;
+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.Files;
+import org.apache.iceberg.ManifestFile;
+import org.apache.iceberg.ManifestFiles;
+import org.apache.iceberg.ManifestWriter;
+import org.apache.iceberg.Metrics;
+import org.apache.iceberg.Parameter;
+import org.apache.iceberg.ParameterizedTestExtension;
+import org.apache.iceberg.Parameters;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Snapshot;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.actions.RepairTable;
+import org.apache.iceberg.data.FileHelpers;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.deletes.BaseDVFileWriter;
+import org.apache.iceberg.deletes.DVFileWriter;
+import org.apache.iceberg.exceptions.CommitFailedException;
+import org.apache.iceberg.exceptions.CommitStateUnknownException;
+import org.apache.iceberg.exceptions.ValidationException;
+import org.apache.iceberg.hadoop.HadoopTables;
+import org.apache.iceberg.io.OutputFile;
+import org.apache.iceberg.io.OutputFileFactory;
+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.TestBase;
+import org.apache.iceberg.spark.source.ThreeColumnRecord;
+import org.apache.iceberg.types.Types;
+import org.apache.iceberg.util.Pair;
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.Row;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.junit.jupiter.api.io.TempDir;
+
+@ExtendWith(ParameterizedTestExtension.class)
+public class TestRepairTableAction extends TestBase {
+
+  private static final HadoopTables TABLES = new HadoopTables(new 
Configuration());
+  private static final Schema SCHEMA =
+      new Schema(
+          optional(1, "c1", Types.IntegerType.get()),
+          optional(2, "c2", Types.StringType.get()),
+          optional(3, "c3", Types.StringType.get()));
+
+  @Parameters(name = "formatVersion = {0}")
+  public static Object[] parameters() {
+    return new Object[][] {new Object[] {1}, new Object[] {2}, new Object[] 
{3}};
+  }
+
+  @Parameter private int formatVersion;
+
+  private String tableLocation = null;
+
+  @TempDir private Path temp;
+  @TempDir private File tableDir;
+
+  @BeforeEach
+  public void setupTableLocation() {
+    this.tableLocation = tableDir.toURI().toString();
+  }
+
+  @TestTemplate
+  public void testRepairEmptyTable() {
+    Table table = createTable(PartitionSpec.unpartitioned());
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+  }
+
+  @TestTemplate
+  public void testRepairTableWithCorrectStats() {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    Snapshot before = table.currentSnapshot();
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().snapshotId())
+        .as("should not commit a snapshot when nothing is repaired")
+        .isEqualTo(before.snapshotId());
+  }
+
+  @TestTemplate
+  public void testNoRepairSelectedIsNoOp() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    DataFile original = onlyDataFile(table);
+    replaceManifestWithCorruptStats(table, original);
+
+    table.refresh();
+    Snapshot before = table.currentSnapshot();
+    DataFile corrupt = onlyDataFile(table);
+
+    // no repair was selected, so execute() must do nothing even though the 
stats are incorrect
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().snapshotId())
+        .as("a repair with nothing selected must not commit")
+        .isEqualTo(before.snapshotId());
+    assertThat(onlyDataFile(table).recordCount())
+        .as("a repair with nothing selected must leave the incorrect stats in 
place")
+        .isEqualTo(corrupt.recordCount());
+  }
+
+  @TestTemplate
+  public void testRepairIncorrectRecordCountAndFileSize() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    List<Object[]> expectedRows = currentRows();
+    DataFile original = onlyDataFile(table);
+
+    // replace the manifest with one whose entry records a wrong record count 
and file size
+    replaceManifestWithCorruptStats(table, original);
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedEntryCount()).isEqualTo(1);
+    assertThat(result.repairedManifests()).hasSize(1);
+
+    table.refresh();
+    DataFile repaired = onlyDataFile(table);
+    assertThat(repaired.recordCount()).isEqualTo(original.recordCount());
+    
assertThat(repaired.fileSizeInBytes()).isEqualTo(original.fileSizeInBytes());
+    assertThat(repaired.location()).isEqualTo(original.location());
+
+    assertThat(currentRows())
+        .as("table contents must be unchanged by the repair")
+        .containsExactlyInAnyOrderElementsOf(expectedRows);
+  }
+
+  @TestTemplate
+  public void testRepairPreservesEntryLineage() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    DataFile original = onlyDataFile(table);
+    List<Row> lineageBefore = entryLineage();
+
+    replaceManifestWithCorruptStats(table, original);
+
+    SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    table.refresh();
+    assertThat(entryLineage())

Review Comment:
   Extend the v3 lineage test to multiple data files and assert that 
data_file.first_row_id survives the repair. entryLineage() currently checks 
only the snapshot ID and sequence numbers, so it would miss changed row IDs.



##########
spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRepairTableAction.java:
##########
@@ -0,0 +1,1085 @@
+/*
+ * 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 static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.assertj.core.api.Assumptions.assumeThat;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.when;
+
+import java.io.File;
+import java.io.IOException;
+import java.nio.file.Path;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.hadoop.conf.Configuration;
+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.Files;
+import org.apache.iceberg.ManifestFile;
+import org.apache.iceberg.ManifestFiles;
+import org.apache.iceberg.ManifestWriter;
+import org.apache.iceberg.Metrics;
+import org.apache.iceberg.Parameter;
+import org.apache.iceberg.ParameterizedTestExtension;
+import org.apache.iceberg.Parameters;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Snapshot;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.actions.RepairTable;
+import org.apache.iceberg.data.FileHelpers;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.deletes.BaseDVFileWriter;
+import org.apache.iceberg.deletes.DVFileWriter;
+import org.apache.iceberg.exceptions.CommitFailedException;
+import org.apache.iceberg.exceptions.CommitStateUnknownException;
+import org.apache.iceberg.exceptions.ValidationException;
+import org.apache.iceberg.hadoop.HadoopTables;
+import org.apache.iceberg.io.OutputFile;
+import org.apache.iceberg.io.OutputFileFactory;
+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.TestBase;
+import org.apache.iceberg.spark.source.ThreeColumnRecord;
+import org.apache.iceberg.types.Types;
+import org.apache.iceberg.util.Pair;
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.Row;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.junit.jupiter.api.io.TempDir;
+
+@ExtendWith(ParameterizedTestExtension.class)
+public class TestRepairTableAction extends TestBase {
+
+  private static final HadoopTables TABLES = new HadoopTables(new 
Configuration());
+  private static final Schema SCHEMA =
+      new Schema(
+          optional(1, "c1", Types.IntegerType.get()),
+          optional(2, "c2", Types.StringType.get()),
+          optional(3, "c3", Types.StringType.get()));
+
+  @Parameters(name = "formatVersion = {0}")
+  public static Object[] parameters() {
+    return new Object[][] {new Object[] {1}, new Object[] {2}, new Object[] 
{3}};
+  }
+
+  @Parameter private int formatVersion;
+
+  private String tableLocation = null;
+
+  @TempDir private Path temp;
+  @TempDir private File tableDir;
+
+  @BeforeEach
+  public void setupTableLocation() {
+    this.tableLocation = tableDir.toURI().toString();
+  }
+
+  @TestTemplate
+  public void testRepairEmptyTable() {
+    Table table = createTable(PartitionSpec.unpartitioned());
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+  }
+
+  @TestTemplate
+  public void testRepairTableWithCorrectStats() {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    Snapshot before = table.currentSnapshot();
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().snapshotId())
+        .as("should not commit a snapshot when nothing is repaired")
+        .isEqualTo(before.snapshotId());
+  }
+
+  @TestTemplate
+  public void testNoRepairSelectedIsNoOp() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    DataFile original = onlyDataFile(table);
+    replaceManifestWithCorruptStats(table, original);
+
+    table.refresh();
+    Snapshot before = table.currentSnapshot();
+    DataFile corrupt = onlyDataFile(table);
+
+    // no repair was selected, so execute() must do nothing even though the 
stats are incorrect
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().snapshotId())
+        .as("a repair with nothing selected must not commit")
+        .isEqualTo(before.snapshotId());
+    assertThat(onlyDataFile(table).recordCount())
+        .as("a repair with nothing selected must leave the incorrect stats in 
place")
+        .isEqualTo(corrupt.recordCount());
+  }
+
+  @TestTemplate
+  public void testRepairIncorrectRecordCountAndFileSize() throws IOException {

Review Comment:
   Parameterize the file-format coverage over Parquet, ORC, and Avro. The 
current tests only write Parquet, so the separate ORC and Avro metrics readers 
are untested by this action. Include a correct-file no-op and a 
record-count/file-size repair for each format.



##########
spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRepairTableAction.java:
##########
@@ -0,0 +1,1085 @@
+/*
+ * 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 static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.assertj.core.api.Assumptions.assumeThat;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.when;
+
+import java.io.File;
+import java.io.IOException;
+import java.nio.file.Path;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.hadoop.conf.Configuration;
+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.Files;
+import org.apache.iceberg.ManifestFile;
+import org.apache.iceberg.ManifestFiles;
+import org.apache.iceberg.ManifestWriter;
+import org.apache.iceberg.Metrics;
+import org.apache.iceberg.Parameter;
+import org.apache.iceberg.ParameterizedTestExtension;
+import org.apache.iceberg.Parameters;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Snapshot;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.actions.RepairTable;
+import org.apache.iceberg.data.FileHelpers;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.deletes.BaseDVFileWriter;
+import org.apache.iceberg.deletes.DVFileWriter;
+import org.apache.iceberg.exceptions.CommitFailedException;
+import org.apache.iceberg.exceptions.CommitStateUnknownException;
+import org.apache.iceberg.exceptions.ValidationException;
+import org.apache.iceberg.hadoop.HadoopTables;
+import org.apache.iceberg.io.OutputFile;
+import org.apache.iceberg.io.OutputFileFactory;
+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.TestBase;
+import org.apache.iceberg.spark.source.ThreeColumnRecord;
+import org.apache.iceberg.types.Types;
+import org.apache.iceberg.util.Pair;
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.Row;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.junit.jupiter.api.io.TempDir;
+
+@ExtendWith(ParameterizedTestExtension.class)
+public class TestRepairTableAction extends TestBase {
+
+  private static final HadoopTables TABLES = new HadoopTables(new 
Configuration());
+  private static final Schema SCHEMA =
+      new Schema(
+          optional(1, "c1", Types.IntegerType.get()),
+          optional(2, "c2", Types.StringType.get()),
+          optional(3, "c3", Types.StringType.get()));
+
+  @Parameters(name = "formatVersion = {0}")
+  public static Object[] parameters() {
+    return new Object[][] {new Object[] {1}, new Object[] {2}, new Object[] 
{3}};
+  }
+
+  @Parameter private int formatVersion;
+
+  private String tableLocation = null;
+
+  @TempDir private Path temp;
+  @TempDir private File tableDir;
+
+  @BeforeEach
+  public void setupTableLocation() {
+    this.tableLocation = tableDir.toURI().toString();
+  }
+
+  @TestTemplate
+  public void testRepairEmptyTable() {
+    Table table = createTable(PartitionSpec.unpartitioned());
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+  }
+
+  @TestTemplate
+  public void testRepairTableWithCorrectStats() {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    Snapshot before = table.currentSnapshot();
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().snapshotId())
+        .as("should not commit a snapshot when nothing is repaired")
+        .isEqualTo(before.snapshotId());
+  }
+
+  @TestTemplate
+  public void testNoRepairSelectedIsNoOp() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    DataFile original = onlyDataFile(table);
+    replaceManifestWithCorruptStats(table, original);
+
+    table.refresh();
+    Snapshot before = table.currentSnapshot();
+    DataFile corrupt = onlyDataFile(table);
+
+    // no repair was selected, so execute() must do nothing even though the 
stats are incorrect
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).execute();
+
+    assertThat(result.repairedManifests()).isEmpty();
+    assertThat(result.repairedEntryCount()).isEqualTo(0);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().snapshotId())
+        .as("a repair with nothing selected must not commit")
+        .isEqualTo(before.snapshotId());
+    assertThat(onlyDataFile(table).recordCount())
+        .as("a repair with nothing selected must leave the incorrect stats in 
place")
+        .isEqualTo(corrupt.recordCount());
+  }
+
+  @TestTemplate
+  public void testRepairIncorrectRecordCountAndFileSize() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    List<Object[]> expectedRows = currentRows();
+    DataFile original = onlyDataFile(table);
+
+    // replace the manifest with one whose entry records a wrong record count 
and file size
+    replaceManifestWithCorruptStats(table, original);
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedEntryCount()).isEqualTo(1);
+    assertThat(result.repairedManifests()).hasSize(1);
+
+    table.refresh();
+    DataFile repaired = onlyDataFile(table);
+    assertThat(repaired.recordCount()).isEqualTo(original.recordCount());
+    
assertThat(repaired.fileSizeInBytes()).isEqualTo(original.fileSizeInBytes());
+    assertThat(repaired.location()).isEqualTo(original.location());
+
+    assertThat(currentRows())
+        .as("table contents must be unchanged by the repair")
+        .containsExactlyInAnyOrderElementsOf(expectedRows);
+  }
+
+  @TestTemplate
+  public void testRepairPreservesEntryLineage() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    DataFile original = onlyDataFile(table);
+    List<Row> lineageBefore = entryLineage();
+
+    replaceManifestWithCorruptStats(table, original);
+
+    SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    table.refresh();
+    assertThat(entryLineage())
+        .as("snapshot id and sequence numbers must be carried through the 
repair")
+        .containsExactlyInAnyOrderElementsOf(lineageBefore);
+  }
+
+  @TestTemplate
+  public void testDryRunDoesNotCommit() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    DataFile original = onlyDataFile(table);
+    replaceManifestWithCorruptStats(table, original);
+
+    table.refresh();
+    Snapshot before = table.currentSnapshot();
+    DataFile corrupt = onlyDataFile(table);
+
+    RepairTable.Result result =
+        
SparkActions.get().repairTable(table).repairFileMetrics().dryRun().execute();
+
+    assertThat(result.repairedEntryCount()).isEqualTo(1);
+    assertThat(result.repairedManifests()).hasSize(1);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().snapshotId())
+        .as("dry run must not commit")
+        .isEqualTo(before.snapshotId());
+    assertThat(onlyDataFile(table).recordCount())
+        .as("dry run must leave the incorrect stats in place")
+        .isEqualTo(corrupt.recordCount());
+  }
+
+  @TestTemplate
+  public void testRepairOnlyRewritesAffectedManifests() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(2));
+    appendRecords(table, records(2));
+
+    table.refresh();
+    assertThat(table.currentSnapshot().dataManifests(table.io())).hasSize(2);
+
+    List<ManifestFile> manifests = 
table.currentSnapshot().dataManifests(table.io());
+    ManifestFile untouched = manifests.get(1);
+
+    // corrupt the entry of one manifest only
+    DataFile fileToCorrupt = readDataFiles(table, manifests.get(0)).get(0);
+    corruptStats(table, manifests.get(0), fileToCorrupt.location());
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedManifests()).hasSize(1);
+    assertThat(result.repairedEntryCount()).isEqualTo(1);
+
+    table.refresh();
+    assertThat(table.currentSnapshot().dataManifests(table.io()))
+        .as("the manifest without incorrect entries must be left in place")
+        .anyMatch(manifest -> manifest.path().equals(untouched.path()));
+  }
+
+  @TestTemplate
+  public void testRepairPartitionedTable() throws IOException {
+    Table table = 
createTable(PartitionSpec.builderFor(SCHEMA).identity("c1").build());
+
+    Dataset<Row> df =
+        spark
+            .createDataFrame(
+                Lists.newArrayList(
+                    new ThreeColumnRecord(1, "AAAA", "A"), new 
ThreeColumnRecord(2, "BBBB", "B")),
+                ThreeColumnRecord.class)
+            .coalesce(1);
+    df.select("c1", "c2", 
"c3").write().format("iceberg").mode("append").save(tableLocation);
+
+    table.refresh();
+    List<Object[]> expectedRows = currentRows();
+    ManifestFile manifest = 
table.currentSnapshot().dataManifests(table.io()).get(0);
+    List<DataFile> files = readDataFiles(table, manifest);
+
+    corruptStats(table, manifest, files.get(0).location());
+
+    RepairTable.Result result = 
SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(result.repairedEntryCount()).isEqualTo(1);
+    
assertThat(currentRows()).containsExactlyInAnyOrderElementsOf(expectedRows);
+  }
+
+  @TestTemplate
+  public void testRepairSkipsColumnMetricsByDefault() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    appendRecords(table, records(4));
+
+    DataFile original = onlyDataFile(table);
+
+    // only the column level statistics are wrong, the record count and the 
file size are correct
+    ManifestFile manifest = 
table.currentSnapshot().dataManifests(table.io()).get(0);
+    corruptStats(table, manifest, original.location(), false);
+
+    RepairTable.Result skipped =
+        SparkActions.get().repairTable(table).repairFileMetrics().execute();
+
+    assertThat(skipped.repairedEntryCount())
+        .as("column metrics must not be compared by default")
+        .isEqualTo(0);
+
+    RepairTable.Result repaired =
+        SparkActions.get()
+            .repairTable(table)
+            .repairFileMetrics()
+            .option(RepairTableSparkAction.REPAIR_COLUMN_METRICS, "true")
+            .execute();
+
+    assertThat(repaired.repairedEntryCount())
+        .as("column metrics are compared when enabled")
+        .isEqualTo(1);
+  }
+
+  @TestTemplate
+  void repairWithCustomColumnMetrics() throws IOException {
+    Table table = createTable(PartitionSpec.unpartitioned());
+    table
+        .updateProperties()
+        .set(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX + "c2", "none")
+        .commit();
+    appendRecords(table, records(4));
+
+    List<Object[]> expectedRows = currentRows();
+    DataFile original = onlyDataFile(table);
+    int columnId = table.schema().findField("c2").fieldId();
+    assertThat(original.valueCounts()).doesNotContainKey(columnId);
+    assertThat(original.lowerBounds()).doesNotContainKey(columnId);
+    assertThat(original.upperBounds()).doesNotContainKey(columnId);
+
+    replaceManifestWithCorruptStats(table, original);
+
+    RepairTable.Result result =
+        SparkActions.get()
+            .repairTable(table)
+            .repairFileMetrics()
+            .option(RepairTableSparkAction.REPAIR_COLUMN_METRICS, "true")
+            .execute();
+
+    assertThat(result.repairedEntryCount()).isEqualTo(1);
+    assertThat(result.repairedManifests()).hasSize(1);
+
+    DataFile repaired = onlyDataFile(table);
+    assertThat(repaired.recordCount()).isEqualTo(original.recordCount());
+    
assertThat(repaired.fileSizeInBytes()).isEqualTo(original.fileSizeInBytes());
+    assertThat(repaired.columnSizes()).isEqualTo(original.columnSizes());
+    assertThat(repaired.valueCounts()).isEqualTo(original.valueCounts());
+    
assertThat(repaired.nullValueCounts()).isEqualTo(original.nullValueCounts());
+    assertThat(repaired.lowerBounds()).isEqualTo(original.lowerBounds());
+    assertThat(repaired.upperBounds()).isEqualTo(original.upperBounds());
+    
assertThat(currentRows()).containsExactlyInAnyOrderElementsOf(expectedRows);

Review Comment:
   Run the repair a second time with the same column-metrics option and assert 
zero repaired entries, no repaired manifests, and an unchanged snapshot ID. The 
current tests stop after the first repair, so they do not verify that repeating 
the action is a no-op.



-- 
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]

Reply via email to