This is an automated email from the ASF dual-hosted git repository.

claudevdm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new bffb545c2b3 Addfiles partition bounds fixes (#40252)
bffb545c2b3 is described below

commit bffb545c2b347eec83a52acd6bd7722c8a3a2335
Author: claudevdm <[email protected]>
AuthorDate: Tue Sep 29 12:35:22 2026 -0400

    Addfiles partition bounds fixes (#40252)
    
    * AddFiles: regression tests for column bounds and metrics-based partition 
inference
    
    These tests describe the correct behavior and fail on master; the commits
    that follow fix them one group at a time. Failures at this commit:
    
    Bounds (getFileMetrics):
    - millis and nanos timestamps, nanos time: bounds stored in the file's
      unit under a micros type, off by 1000x.
    - millis time and uint32: ClassCastException (Integer to Long) while
      collecting bounds.
    
    Partition inference without a location prefix:
    - all-null boolean partition column is registered under flag=false, an
      all-null string column under the string "null", and a file lacking
      the column under flag=false: the null partition goes through a path
      string and is parsed back.
    - all-null timestamp under day(ts) is refused: "Text 'null' could not
      be parsed".
    - a column with statistics switched off is registered under flag=false
      although every value is true; an INT96 column is refused with the
      same parse error instead of an unknown partition.
    - an Avro file on a partitioned table fails the pipeline with a
      NullPointerException (null bound maps).
    - a millis timestamp file lands in day 19 (1970) instead of 2024-01-01.
    - with write.metadata.metrics.default=none, an unrelated uint32 column
      fails the pipeline with a ClassCastException during the partition
      recollection.
    
    testMicrosTimestampBoundsAreUnchanged passes and stays as the control.
    
    * AddFiles: never guess a partition for a column without bounds
    
    Metrics-based partition inference treated a partition source column
    with no bounds as all null. That holds for a column whose values are
    all null, but bounds are also missing when statistics are switched off,
    for INT96 timestamps (Iceberg collects nothing for them) and for Avro
    files, whose Metrics carry null bound maps.
    
    - An Avro file on a partitioned table failed the bundle with a
      NullPointerException on metrics.lowerBounds().keySet(), so one Avro
      file failed the job. Null bound maps are now read as empty.
    - A column without bounds is an unknown partition (error row naming the
      column and pointing at the location prefix) unless the null partition
      is proven: the file is empty, the column's null count equals its
      value count, or the Parquet file does not contain the column at all
      (a file written before the column existed reads as null everywhere).
    
    Turns green: testAvroFileOnAPartitionedTableIsAnUnknownPartition,
    testPartitionColumnWithoutStatisticsIsAnUnknownPartition,
    testInt96PartitionColumnIsAnUnknownPartition.
    
    * AddFiles: set the metrics-inferred partition as values, not as a path 
string
    
    getPartitionFromMetrics returned PartitionKey.toPath() and the DataFile
    was built with withPartitionPath, which parses each value back from
    text. A null partition value is written as the text "null", so:
    
    - boolean identity partition: Boolean.valueOf("null") files the data
      under flag=false, and queries pruning on the partition return wrong
      rows;
    - string identity partition: the file lands under the string "null";
    - int, date and timestamp partitions: the parse fails and the file is
      refused ("Text 'null' could not be parsed at index 0").
    
    getPartitionFromMetrics now returns the PartitionKey and the DataFile is
    built with withPartition. The location-prefix option still goes through
    the path, where Iceberg maps __HIVE_DEFAULT_PARTITION__ to null.
    
    Turns green: the all-null boolean, string and day(timestamp) tests and
    testFileLackingThePartitionColumnRegistersUnderTheNullPartition.
    
    * AddFiles: recollect partition metrics for the partition columns only
    
    When the file's own metrics lack bounds for a partition source column
    (e.g. write.metadata.metrics.default=none), the partition is inferred
    from metrics recollected with "full" for the partition columns. Every
    other column fell back to Iceberg's default mode, so bounds were computed
    for the whole file, and a column Iceberg cannot collect bounds for
    (uint32, millis time: ClassCastException) failed the bundle from a code
    path that only catches UnknownPartitionException.
    
    The recollection now sets the default mode to none, so only the
    partition columns are read. It is also less work on wide files.
    
    Turns green:
    testUnrelatedColumnDoesNotBreakInferenceUnderMetricsModeNone.
    
    * AddFiles: store column bounds in the unit of the Iceberg type
    
    Files without field ids resolve through the name mapping, and their
    column bounds are what partition inference and query pruning read, so
    they must be encoded in the Iceberg type's unit. Iceberg's converter
    maps TIMESTAMP/TIME in MILLIS or NANOS to micros types and uint32 to
    long, but ParquetUtil.footerMetrics stores the file's raw statistic:
    
    - TIMESTAMP(MILLIS|NANOS), TIME(NANOS): bounds off by 1000x, silently.
      A 2024 millis timestamp gets a bound in 1970; day(ts) puts the file
      in the wrong partition and range predicates prune it away, while the
      data reads correctly.
    - TIME(MILLIS) int32 and UINT32: ClassCastException (Integer to Long)
      while collecting bounds, so the file is routed to errors.
    
    Fix: BoundAdjustment, applied in getFileMetrics after the id mapping.
    Affected int32 columns are presented to Iceberg without their
    annotation so it computes plain int bounds instead of throwing; every
    affected lower/upper bound is then rewritten as an 8-byte long in the
    Iceberg unit (millis x1000; nanos /1000 with lower rounded down and
    upper up so bounds stay conservative; uint32 via toUnsignedLong).
    Counts are untouched. Files whose columns need no adjustment take the
    old path unchanged.
    
    Turns green: the six bounds tests' remaining failures and
    testMillisTimestampFileLandsInItsDayPartition.
    
    * AddFiles: refuse a file whose partition column holds both nulls and values
    
    Found while fixing the partition inference; not in the fuzz findings.
    Column bounds ignore nulls, so a partition source column holding true
    and null has bounds true..true and the file was registered under
    flag=true. Its null rows belong to the null partition, so a query
    pruning on the partition silently loses them.
    
    Such a file spans two partitions and is now an unknown partition (error
    row naming the column). A transform that maps every value to null (void)
    is unaffected: nulls and values land in the same partition.
    
    Test and fix together: on the previous commit
    testPartitionColumnWithNullsAndValuesIsAnUnknownPartition fails with the
    file registered and no error row.
    
    * fix boundadjustment
    
    * comments
    
    * missing file
    
    * comments
    
    * trigger tests
    
    * move utils, tests
    
    * fixes
    
    * traverse from newest schema
---
 .../IO_Iceberg_Integration_Tests.json              |   2 +-
 .../org/apache/beam/sdk/io/iceberg/AddFiles.java   | 241 ++++++--
 .../beam/sdk/io/iceberg/BoundAdjustment.java       | 395 ++++++++++++
 .../apache/beam/sdk/io/iceberg/DryRunReport.java   |   2 +-
 .../beam/sdk/io/iceberg/NameMappingUtils.java      |  12 +
 .../beam/sdk/io/iceberg/ParquetFieldIds.java       | 185 ++++++
 .../beam/sdk/io/iceberg/AddFilesMetricsTest.java   | 680 +++++++++++++++++++++
 .../apache/beam/sdk/io/iceberg/AddFilesTest.java   | 118 +++-
 .../beam/sdk/io/iceberg/BoundAdjustmentTest.java   | 391 ++++++++++++
 .../beam/sdk/io/iceberg/NameMappingUtilsTest.java  |  31 +
 .../beam/sdk/io/iceberg/ParquetFieldIdsTest.java   | 412 +++++++++++++
 .../beam/sdk/io/iceberg/ParquetTestFiles.java      | 170 ++++++
 12 files changed, 2569 insertions(+), 70 deletions(-)

diff --git a/.github/trigger_files/IO_Iceberg_Integration_Tests.json 
b/.github/trigger_files/IO_Iceberg_Integration_Tests.json
index 7ab7bcd9a9c..7392be3b11c 100644
--- a/.github/trigger_files/IO_Iceberg_Integration_Tests.json
+++ b/.github/trigger_files/IO_Iceberg_Integration_Tests.json
@@ -1,4 +1,4 @@
 {
     "comment": "Modify this file in a trivial way to cause this test suite to 
run.",
-    "modification": 2
+    "modification": 10
 }
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
index 4093b2fb6e5..cb54183e0a0 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java
@@ -40,7 +40,6 @@ import java.util.concurrent.Callable;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Collectors;
-import java.util.stream.Stream;
 import org.apache.beam.sdk.coders.KvCoder;
 import org.apache.beam.sdk.coders.RowCoder;
 import org.apache.beam.sdk.coders.VarIntCoder;
@@ -86,6 +85,7 @@ import org.apache.iceberg.ManifestFiles;
 import org.apache.iceberg.ManifestWriter;
 import org.apache.iceberg.Metrics;
 import org.apache.iceberg.MetricsConfig;
+import org.apache.iceberg.MetricsModes;
 import org.apache.iceberg.PartitionField;
 import org.apache.iceberg.PartitionKey;
 import org.apache.iceberg.PartitionSpec;
@@ -104,11 +104,10 @@ import org.apache.iceberg.mapping.NameMapping;
 import org.apache.iceberg.mapping.NameMappingParser;
 import org.apache.iceberg.orc.OrcMetrics;
 import org.apache.iceberg.parquet.ParquetSchemaUtil;
-import org.apache.iceberg.parquet.ParquetUtil;
 import org.apache.iceberg.transforms.Transform;
 import org.apache.iceberg.types.Conversions;
 import org.apache.iceberg.types.Type;
-import org.apache.parquet.hadoop.metadata.FileMetaData;
+import org.apache.parquet.column.ColumnDescriptor;
 import org.apache.parquet.hadoop.metadata.ParquetMetadata;
 import org.apache.parquet.schema.MessageType;
 import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
@@ -472,6 +471,9 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
     static final String UNKNOWN_PARTITION_ERROR = "Could not determine the 
file's partition: ";
     static final String UNREADABLE_SCHEMA_ERROR = "Could not read the file's 
schema: ";
     static final String UNCOVERED_ERROR = "Table schema does not cover the 
file after refresh: ";
+    static final String FIELD_ID_ERROR =
+        "The file's Parquet field ids do not match the table's, and readers 
resolve its columns"
+            + " by those ids: ";
     static final String PINNED_COLUMN_ERROR = "Pinned required column ";
     static final String MISSING_TABLE_ERROR =
         "Table does not exist and the schema pre-pass could not create it (no 
readable Parquet"
@@ -684,6 +686,23 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
           return errorResult(filePath, verdict.error, timestamp, window, 
paneInfo);
         }
 
+        // What readers resolve a file without field ids with, so its metrics 
and partition agree
+        // with what they read.
+        NameMapping mapping =
+            NameMappingUtils.forReaders(
+                table.schema(),
+                NameMappingUtils.parseOrNull(
+                    
table.properties().get(TableProperties.DEFAULT_NAME_MAPPING)));
+        ParquetFieldIds.@Nullable Resolved resolvedFooter = null;
+        if (parquetFooter != null) {
+          try {
+            resolvedFooter = ParquetFieldIds.resolve(parquetFooter, table, 
mapping);
+          } catch (ParquetFieldIds.ConflictException e) {
+            return errorResult(
+                filePath, FIELD_ID_ERROR + e.getMessage(), timestamp, window, 
paneInfo);
+          }
+        }
+
         InputFile inputFile = table.io().newInputFile(filePath);
 
         Metrics metrics;
@@ -693,39 +712,51 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
                   inputFile,
                   format,
                   MetricsConfig.forTable(table),
-                  MappingUtil.create(table.schema()),
-                  parquetFooter);
+                  mapping,
+                  table.schema(),
+                  resolvedFooter);
         } catch (Exception e) {
           return errorResult(filePath, errorMessage(e), timestamp, window, 
paneInfo);
         }
 
         // Figure out which partition this DataFile should go to
-        String partitionPath;
-        if (table.spec().isUnpartitioned()) {
-          partitionPath = "";
-        } else if (!Strings.isNullOrEmpty(prefix)) {
+        String partitionPath = "";
+        @Nullable PartitionKey partitionFromMetrics = null;
+        boolean partitioned = table.spec().isPartitioned();
+        if (partitioned && !Strings.isNullOrEmpty(prefix)) {
           // option 1: use directory structure to determine partition
           // Note: we don't validate the DataFile content here
           partitionPath = getPartitionFromFilePath(filePath);
-        } else {
+        } else if (partitioned) {
           try {
             // option 2: examine DataFile min/max statistics to determine 
partition
-            partitionPath = getPartitionFromMetrics(metrics, inputFile, table, 
parquetFooter);
+            partitionFromMetrics =
+                getPartitionFromMetrics(metrics, inputFile, table, mapping, 
resolvedFooter);
           } catch (UnknownPartitionException e) {
             return errorResult(
                 filePath, UNKNOWN_PARTITION_ERROR + e.getMessage(), timestamp, 
window, paneInfo);
+          } catch (RuntimeException e) {
+            // e.g. a bound that does not decode as the table's type: one 
error row, not a failed
+            // bundle
+            return errorResult(
+                filePath, UNKNOWN_PARTITION_ERROR + errorMessage(e), 
timestamp, window, paneInfo);
           }
         }
 
         try {
-          DataFile df =
+          DataFiles.Builder builder =
               DataFiles.builder(table.spec())
                   .withPath(filePath)
                   .withFormat(format)
                   .withMetrics(metrics)
-                  .withFileSizeInBytes(inputFile.getLength())
-                  .withPartitionPath(partitionPath)
-                  .build();
+                  .withFileSizeInBytes(inputFile.getLength());
+          if (partitionFromMetrics != null) {
+            // Set as values: a path string cannot carry a null ("flag=null" 
parses as false).
+            builder = builder.withPartition(partitionFromMetrics);
+          } else {
+            builder = builder.withPartitionPath(partitionPath);
+          }
+          DataFile df = builder.build();
           return new ProcessResult(
               SerializableDataFile.from(df, table.spec()),
               null,
@@ -942,37 +973,49 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
      * <p>In these cases, we output the DataFile to the DLQ, because assigning 
an incorrect
      * partition may lead to it being incorrectly ignored by downstream 
queries.
      */
-    static String getPartitionFromMetrics(
-        Metrics metrics, InputFile inputFile, Table table, @Nullable 
ParquetMetadata preReadFooter)
+    static PartitionKey getPartitionFromMetrics(
+        Metrics metrics,
+        InputFile inputFile,
+        Table table,
+        NameMapping mapping,
+        ParquetFieldIds.@Nullable Resolved footer)
         throws UnknownPartitionException {
       List<PartitionField> fields = table.spec().fields();
       List<Integer> sourceIds =
           
fields.stream().map(PartitionField::sourceId).collect(Collectors.toList());
+      // Iceberg looks metrics modes up by the file's column names, which 
differ from the table's
+      // after a rename.
+      List<String> metricsNames = new ArrayList<>();
+      if (footer != null) {
+        metricsNames = footer.columnNames(sourceIds);
+      } else {
+        for (int sourceId : sourceIds) {
+          metricsNames.add(table.schema().findColumnName(sourceId));
+        }
+      }
+      // The table's metrics are reused only when they hold every partition 
column's exact bounds.
+      // Truncated bounds (under Iceberg's default, truncate(16), a 
23-character value is stored as
+      // two different 16-character bounds) and nanos bounds rounded outward 
work for pruning but
+      // cannot name the partition. Recollected metrics are not attached to 
the DataFile.
       Metrics partitionMetrics;
-      // Check if metrics already includes partition columns (configured by 
table properties):
-      if (metrics.lowerBounds().keySet().containsAll(sourceIds)
-          && metrics.upperBounds().keySet().containsAll(sourceIds)) {
+      if (orEmpty(metrics.lowerBounds()).keySet().containsAll(sourceIds)
+          && orEmpty(metrics.upperBounds()).keySet().containsAll(sourceIds)
+          && fullMetrics(table, metricsNames)
+          && (footer == null || !BoundAdjustment.roundsBounds(footer, 
table.schema(), sourceIds))) {
         partitionMetrics = metrics;
+      } else if (footer != null) {
+        partitionMetrics =
+            BoundAdjustment.partitionMetrics(
+                footer, table.schema(), fullMetricsFor(metricsNames), mapping);
       } else {
-        // Otherwise, recollect metrics and ensure it includes all partition 
fields.
-        // Note: we don't attach these additional metrics to the DataFile 
because we can't assume
-        // that's in the user's best interest.
-        // Some tables are very wide and users may not want to store excessive 
metadata.
-        List<String> sourceNames =
-            fields.stream()
-                .map(pf -> table.schema().findColumnName(pf.sourceId()))
-                .collect(Collectors.toList());
-        Map<String, String> configProps =
-            sourceNames.stream()
-                .collect(Collectors.toMap(s -> 
"write.metadata.metrics.column." + s, s -> "full"));
-        MetricsConfig configWithPartitionFields = 
MetricsConfig.fromProperties(configProps);
         partitionMetrics =
             getFileMetrics(
                 inputFile,
                 inferFormat(inputFile.location()),
-                configWithPartitionFields,
-                MappingUtil.create(table.schema()),
-                preReadFooter);
+                fullMetricsFor(metricsNames),
+                mapping,
+                table.schema(),
+                null);
       }
 
       PartitionKey pk = new PartitionKey(table.spec(), table.schema());
@@ -986,10 +1029,20 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
         // Make a best effort estimate by comparing the lower and upper 
transformed values.
         // If the transformed values are equal, assume that the DataFile's 
data safely
         // aligns with the same partition.
-        ByteBuffer lowerBytes = 
partitionMetrics.lowerBounds().get(field.sourceId());
-        ByteBuffer upperBytes = 
partitionMetrics.upperBounds().get(field.sourceId());
+        ByteBuffer lowerBytes = 
orEmpty(partitionMetrics.lowerBounds()).get(field.sourceId());
+        ByteBuffer upperBytes = 
orEmpty(partitionMetrics.upperBounds()).get(field.sourceId());
         if (lowerBytes == null && upperBytes == null) {
-          continue;
+          // No bounds. The null partition is right only when every value is 
known to be null;
+          // otherwise the partition is unknowable and must not be guessed.
+          if (allValuesNull(partitionMetrics, field.sourceId())
+              || lacksColumn(footer, field.sourceId())) {
+            continue;
+          }
+          throw new UnknownPartitionException(
+              "No column bounds for partition source column "
+                  + table.schema().findColumnName(field.sourceId())
+                  + " (statistics are missing, or are not collected for its 
type, e.g. INT96, or"
+                  + " for this file format); set a location prefix to 
partition by path instead");
         } else if (lowerBytes == null || upperBytes == null) {
           throw new UnknownPartitionException(
               "Only one of the min/max was was null, for field "
@@ -1004,11 +1057,91 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
           throw new UnknownPartitionException(
               "Min and max transformed values were not equal, for column: " + 
field.name());
         }
+        // Bounds ignore nulls, and a null row belongs to the null partition.
+        if (lowerTransformedValue != null && hasNulls(partitionMetrics, 
field.sourceId())) {
+          throw new UnknownPartitionException(
+              "Column has both null and non-null values, which belong to 
different partitions: "
+                  + table.schema().findColumnName(field.sourceId()));
+        }
 
         pk.set(i, lowerTransformedValue);
       }
 
-      return pk.toPath();
+      return pk;
+    }
+
+    /**
+     * Full metrics for the named columns and none for the rest: another 
column whose bounds cannot
+     * be collected must not fail the inference.
+     */
+    private static MetricsConfig fullMetricsFor(List<String> columnNames) {
+      Map<String, String> props = new HashMap<>();
+      props.put(TableProperties.DEFAULT_WRITE_METRICS_MODE, "none");
+      for (String columnName : columnNames) {
+        props.put(TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX + 
columnName, "full");
+      }
+      return MetricsConfig.fromProperties(props);
+    }
+
+    private static boolean fullMetrics(Table table, List<String> metricsNames) 
{
+      MetricsConfig config = MetricsConfig.forTable(table);
+      for (String name : metricsNames) {
+        if (!(config.columnMode(name) instanceof MetricsModes.Full)) {
+          return false;
+        }
+      }
+      return true;
+    }
+
+    /** Avro metrics carry null bound maps. */
+    private static Map<Integer, ByteBuffer> orEmpty(@Nullable Map<Integer, 
ByteBuffer> bounds) {
+      if (bounds == null) {
+        return Collections.emptyMap();
+      }
+      return bounds;
+    }
+
+    /** True when the file is empty or the column's null count equals its 
value count. */
+    private static boolean allValuesNull(Metrics metrics, int fieldId) {
+      Long records = metrics.recordCount();
+      if (records != null && records == 0) {
+        return true;
+      }
+      Map<Integer, Long> valueCounts = metrics.valueCounts();
+      Map<Integer, Long> nullCounts = metrics.nullValueCounts();
+      if (valueCounts == null || nullCounts == null) {
+        return false;
+      }
+      Long valueCount = valueCounts.get(fieldId);
+      Long nullCount = nullCounts.get(fieldId);
+      return valueCount != null && nullCount != null && 
valueCount.equals(nullCount);
+    }
+
+    private static boolean hasNulls(Metrics metrics, int fieldId) {
+      Map<Integer, Long> nullCounts = metrics.nullValueCounts();
+      if (nullCounts == null) {
+        return false;
+      }
+      Long nullCount = nullCounts.get(fieldId);
+      return nullCount != null && nullCount > 0;
+    }
+
+    /**
+     * True when a Parquet file does not contain the column at all, e.g. it 
was written before the
+     * column existed. Every row then reads as null. Unknown for other formats.
+     */
+    private static boolean lacksColumn(ParquetFieldIds.@Nullable Resolved 
footer, int fieldId) {
+      if (footer == null) {
+        return false;
+      }
+      // A partition source is a leaf column, so the file's leaf columns are 
enough.
+      for (ColumnDescriptor column : 
footer.footer().getFileMetaData().getSchema().getColumns()) {
+        org.apache.parquet.schema.Type.ID id = 
column.getPrimitiveType().getId();
+        if (id != null && id.intValue() == fieldId) {
+          return false;
+        }
+      }
+      return true;
     }
   }
 
@@ -1106,7 +1239,8 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
     private static void ensureNameMappingPresent(Table table) {
       // Forces name-based resolution: zero-copy files typically don't carry
       // field ids, so any schema column missing from the mapping is unreadable
-      // in registered files.
+      // in registered files. Registration resolves files with
+      // NameMappingUtils.forReaders, which must agree with what this stores.
       @Nullable NameMapping existing =
           NameMappingUtils.parseOrNull(
               table.properties().get(TableProperties.DEFAULT_NAME_MAPPING));
@@ -1204,21 +1338,20 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
   }
 
   @SuppressWarnings("argument")
-  public static Metrics getFileMetrics(
+  static Metrics getFileMetrics(
       InputFile file,
       FileFormat format,
       MetricsConfig config,
       NameMapping mapping,
-      @Nullable ParquetMetadata preReadFooter) {
+      org.apache.iceberg.Schema tableSchema,
+      ParquetFieldIds.@Nullable Resolved parquetFooter) {
     switch (format) {
       case PARQUET:
-        ParquetMetadata footer =
-            checkStateNotNull(preReadFooter, "Parquet metrics require the 
pre-read footer");
-        MessageType originalMessageType = footer.getFileMetaData().getSchema();
-        if (!ParquetSchemaUtil.hasIds(originalMessageType)) {
-          footer = getFooterWithTypeIds(originalMessageType, footer, mapping);
-        }
-        return ParquetUtil.footerMetrics(footer, Stream.empty(), config, 
mapping);
+        return BoundAdjustment.footerMetrics(
+            checkStateNotNull(parquetFooter, "Parquet metrics require the 
resolved footer"),
+            tableSchema,
+            config,
+            mapping);
       case ORC:
         return OrcMetrics.fromInputFile(file, config, mapping);
       case AVRO:
@@ -1251,16 +1384,6 @@ public class AddFiles extends 
PTransform<PCollection<String>, PCollectionRowTupl
     }
   }
 
-  static ParquetMetadata getFooterWithTypeIds(
-      MessageType originalMessageType, ParquetMetadata footer, NameMapping 
mapping) {
-    originalMessageType = 
ParquetSchemaUtil.applyNameMapping(originalMessageType, mapping);
-    FileMetaData oldFileMeta = footer.getFileMetaData();
-    FileMetaData newFileMeta =
-        new FileMetaData(
-            originalMessageType, oldFileMeta.getKeyValueMetaData(), 
oldFileMeta.getCreatedBy());
-    return new ParquetMetadata(newFileMeta, footer.getBlocks());
-  }
-
   static class UnknownFormatException extends IllegalArgumentException {}
 
   static class UnknownPartitionException extends IllegalStateException {
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundAdjustment.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundAdjustment.java
new file mode 100644
index 00000000000..5807aeb88e1
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundAdjustment.java
@@ -0,0 +1,395 @@
+/*
+ * 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.beam.sdk.io.iceberg;
+
+import java.nio.ByteBuffer;
+import java.nio.ByteOrder;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.stream.Stream;
+import org.apache.iceberg.Metrics;
+import org.apache.iceberg.MetricsConfig;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.mapping.NameMapping;
+import org.apache.iceberg.parquet.ParquetUtil;
+import org.apache.iceberg.types.Conversions;
+import org.apache.iceberg.types.Type.TypeID;
+import org.apache.iceberg.types.Types;
+import org.apache.parquet.column.ColumnDescriptor;
+import org.apache.parquet.column.statistics.Statistics;
+import org.apache.parquet.hadoop.metadata.BlockMetaData;
+import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData;
+import org.apache.parquet.hadoop.metadata.ColumnPath;
+import org.apache.parquet.hadoop.metadata.FileMetaData;
+import org.apache.parquet.hadoop.metadata.ParquetMetadata;
+import org.apache.parquet.schema.GroupType;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
+import 
org.apache.parquet.schema.LogicalTypeAnnotation.IntLogicalTypeAnnotation;
+import 
org.apache.parquet.schema.LogicalTypeAnnotation.TimeLogicalTypeAnnotation;
+import org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit;
+import 
org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation;
+import org.apache.parquet.schema.MessageType;
+import org.apache.parquet.schema.PrimitiveType;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.apache.parquet.schema.Type;
+import org.apache.parquet.schema.TypeConverter;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * Iceberg collects bounds in the file column's unit and width, but readers 
decode them with the
+ * table column's type: a millis or nanos timestamp under a micros column, or 
a millis or micros one
+ * under a nanos column, would be off by a factor of 1000 or more, and Iceberg 
types an unsigned
+ * 32-bit int as a long, so collecting its int statistics throws. Bounds are 
what partition
+ * inference and query pruning read, so they are rewritten in the table 
column's unit: the affected
+ * INT32 columns are presented to Iceberg without their annotation (so it 
computes plain int bounds
+ * instead of throwing) and every affected bound is converted afterwards, or 
dropped when it does
+ * not fit the table's unit. Value counts and null counts are unaffected. A 
uint32 column must stay
+ * below 2^31: see {@link #checkUnsignedRange}.
+ */
+enum BoundAdjustment {
+  /** Millis stored under a micros type: times 1000. */
+  MILLIS_TO_MICROS,
+  /** Millis stored under a nanos type: times 1,000,000. */
+  MILLIS_TO_NANOS,
+  /** Micros stored under a nanos type: times 1000. */
+  MICROS_TO_NANOS,
+  /**
+   * Nanos stored under a micros type: divided by 1000, rounded outward for 
pruning and down for
+   * partition inference.
+   */
+  NANOS_TO_MICROS,
+  /** Unsigned 32-bit int stored under a long. */
+  UINT32_TO_LONG,
+  /** Unsigned 32-bit int stored under an int: the bounds already are ints. */
+  UINT32_TO_INT;
+
+  static final String UNSIGNED_RANGE_ERROR =
+      "Iceberg readers return unsigned 32-bit values of 2^31 or more as 
negative numbers, and"
+          + " this column holds one: ";
+
+  static final String UNSIGNED_TYPE_ERROR =
+      "An unsigned 32-bit column can only be registered under an int or long 
column: ";
+
+  /**
+   * Iceberg's metrics for the file, with every bound in the table column's 
unit. Throws for a
+   * uint32 column that holds a value of 2^31 or more.
+   */
+  static Metrics footerMetrics(
+      ParquetFieldIds.Resolved resolved,
+      Schema tableSchema,
+      MetricsConfig config,
+      NameMapping mapping) {
+    return collect(resolved, tableSchema, config, mapping, true);
+  }
+
+  /**
+   * Bounds for partition inference: nanos are rounded down at both ends. A 
value's day, hour or
+   * identity partition is that of its floor in micros, while the upper bound 
pruning needs, rounded
+   * up, would put 23:59:59.999999999 in the next day.
+   */
+  static Metrics partitionMetrics(
+      ParquetFieldIds.Resolved resolved,
+      Schema tableSchema,
+      MetricsConfig config,
+      NameMapping mapping) {
+    return collect(resolved, tableSchema, config, mapping, false);
+  }
+
+  private static Metrics collect(
+      ParquetFieldIds.Resolved resolved,
+      Schema tableSchema,
+      MetricsConfig config,
+      NameMapping mapping,
+      boolean roundUpperUp) {
+    ParquetMetadata footer = resolved.footer();
+    Map<Integer, BoundAdjustment> adjustments =
+        forSchema(footer.getFileMetaData().getSchema(), tableSchema);
+    checkUnsignedRange(footer, adjustments);
+    if (adjustments.isEmpty()) {
+      return ParquetUtil.footerMetrics(footer, Stream.empty(), config, 
mapping);
+    }
+    Metrics raw =
+        ParquetUtil.footerMetrics(
+            withNeutralTypes(footer, adjustments), Stream.empty(), config, 
mapping);
+    return apply(raw, adjustments, roundUpperUp);
+  }
+
+  /** Whether any of the columns' stored bounds are rounded outward: nanos 
under a micros column. */
+  static boolean roundsBounds(
+      ParquetFieldIds.Resolved resolved, Schema tableSchema, List<Integer> 
fieldIds) {
+    Map<Integer, BoundAdjustment> adjustments =
+        forSchema(resolved.footer().getFileMetaData().getSchema(), 
tableSchema);
+    for (int fieldId : fieldIds) {
+      if (adjustments.get(fieldId) == NANOS_TO_MICROS) {
+        return true;
+      }
+    }
+    return false;
+  }
+
+  static Map<Integer, BoundAdjustment> forSchema(MessageType fileSchema, 
Schema tableSchema) {
+    Map<Integer, BoundAdjustment> adjustments = new HashMap<>();
+    for (ColumnDescriptor column : fileSchema.getColumns()) {
+      PrimitiveType primitive = column.getPrimitiveType();
+      Type.ID id = primitive.getId();
+      if (id == null) {
+        continue;
+      }
+      org.apache.iceberg.types.@Nullable Type tableType = 
tableSchema.findType(id.intValue());
+      if (tableType == null) {
+        continue;
+      }
+      @Nullable BoundAdjustment adjustment = forPrimitive(primitive, 
tableType.typeId());
+      if (adjustment == null && isUnsigned32(primitive)) {
+        // Iceberg would cast the int statistics to the long it maps uint32 
to, and throw.
+        throw new IllegalArgumentException(
+            UNSIGNED_TYPE_ERROR + String.join(".", column.getPath()) + " is " 
+ tableType);
+      }
+      if (adjustment != null) {
+        adjustments.put(id.intValue(), adjustment);
+      }
+    }
+    return adjustments;
+  }
+
+  /**
+   * The conversion that puts this file column's bounds into the table 
column's unit. There are
+   * three cases:
+   *
+   * <ul>
+   *   <li>a millis or nanos timestamp under a micros timestamp column, or a 
millis or micros one
+   *       under a nanos column;
+   *   <li>a millis or nanos time under a time column;
+   *   <li>an unsigned 32-bit int under a long or int column.
+   * </ul>
+   *
+   * <p>Null for anything else, including units that already match: those 
bounds stay as Iceberg
+   * computes them.
+   */
+  private static @Nullable BoundAdjustment forPrimitive(PrimitiveType 
primitive, TypeID tableType) {
+    LogicalTypeAnnotation annotation = primitive.getLogicalTypeAnnotation();
+    if (annotation instanceof TimestampLogicalTypeAnnotation) {
+      TimeUnit fileUnit = ((TimestampLogicalTypeAnnotation) 
annotation).getUnit();
+      if (tableType == TypeID.TIMESTAMP) {
+        return toMicros(fileUnit);
+      }
+      if (tableType == TypeID.TIMESTAMP_NANO) {
+        return toNanos(fileUnit);
+      }
+      return null;
+    }
+    if (annotation instanceof TimeLogicalTypeAnnotation) {
+      if (tableType != TypeID.TIME) {
+        return null;
+      }
+      return toMicros(((TimeLogicalTypeAnnotation) annotation).getUnit());
+    }
+    if (isUnsigned32(primitive)) {
+      if (tableType == TypeID.LONG) {
+        return UINT32_TO_LONG;
+      }
+      if (tableType == TypeID.INTEGER) {
+        return UINT32_TO_INT;
+      }
+    }
+    return null;
+  }
+
+  private static boolean isUnsigned32(PrimitiveType primitive) {
+    LogicalTypeAnnotation annotation = primitive.getLogicalTypeAnnotation();
+    if (!(annotation instanceof IntLogicalTypeAnnotation)) {
+      return false;
+    }
+    IntLogicalTypeAnnotation intType = (IntLogicalTypeAnnotation) annotation;
+    return intType.getBitWidth() == 32 && !intType.isSigned();
+  }
+
+  private static @Nullable BoundAdjustment toMicros(TimeUnit fileUnit) {
+    switch (fileUnit) {
+      case MILLIS:
+        return MILLIS_TO_MICROS;
+      case NANOS:
+        return NANOS_TO_MICROS;
+      default:
+        return null;
+    }
+  }
+
+  private static @Nullable BoundAdjustment toNanos(TimeUnit fileUnit) {
+    switch (fileUnit) {
+      case MILLIS:
+        return MILLIS_TO_NANOS;
+      case MICROS:
+        return MICROS_TO_NANOS;
+      default:
+        return null;
+    }
+  }
+
+  /**
+   * Refuses a file whose uint32 column holds a value of 2^31 or more. Iceberg 
readers return such a
+   * value as a negative number (3000000000 reads back as -1294967296), so the 
file would read back
+   * wrong whatever bounds are stored for it. Below 2^31 a value reads back 
unchanged, and the
+   * bounds Iceberg computes by comparing values as signed ints are correct.
+   */
+  static void checkUnsignedRange(
+      ParquetMetadata footer, Map<Integer, BoundAdjustment> adjustments) {
+    Set<ColumnPath> unsigned = new HashSet<>();
+    for (ColumnDescriptor column : 
footer.getFileMetaData().getSchema().getColumns()) {
+      Type.ID id = column.getPrimitiveType().getId();
+      if (id == null) {
+        continue;
+      }
+      @Nullable BoundAdjustment adjustment = adjustments.get(id.intValue());
+      if (adjustment == UINT32_TO_LONG || adjustment == UINT32_TO_INT) {
+        unsigned.add(ColumnPath.get(column.getPath()));
+      }
+    }
+    if (unsigned.isEmpty()) {
+      return;
+    }
+    for (BlockMetaData block : footer.getBlocks()) {
+      for (ColumnChunkMetaData chunk : block.getColumns()) {
+        if (unsigned.contains(chunk.getPath()) && 
maxIsNegative(chunk.getStatistics())) {
+          throw new IllegalArgumentException(UNSIGNED_RANGE_ERROR + 
chunk.getPath().toDotString());
+        }
+      }
+    }
+  }
+
+  /** A uint32 max that is negative as a signed int is 2^31 or more. */
+  private static boolean maxIsNegative(@Nullable Statistics<?> stats) {
+    if (stats == null || !stats.hasNonNullValue()) {
+      return false;
+    }
+    Object max = stats.genericGetMax();
+    return max instanceof Integer && (Integer) max < 0;
+  }
+
+  /** The footer with annotations removed from adjusted INT32 columns. */
+  static ParquetMetadata withNeutralTypes(
+      ParquetMetadata footer, Map<Integer, BoundAdjustment> adjustments) {
+    MessageType neutral =
+        (MessageType)
+            footer.getFileMetaData().getSchema().convertWith(new 
WithoutAnnotations(adjustments));
+    FileMetaData meta = footer.getFileMetaData();
+    return new ParquetMetadata(
+        new FileMetaData(neutral, meta.getKeyValueMetaData(), 
meta.getCreatedBy()),
+        footer.getBlocks());
+  }
+
+  /** Rebuilds a schema with the annotation removed from each adjusted INT32 
column. */
+  private static final class WithoutAnnotations implements TypeConverter<Type> 
{
+    private final Map<Integer, BoundAdjustment> adjustments;
+
+    WithoutAnnotations(Map<Integer, BoundAdjustment> adjustments) {
+      this.adjustments = adjustments;
+    }
+
+    @Override
+    public Type convertPrimitiveType(List<GroupType> path, PrimitiveType 
primitive) {
+      Type.ID id = primitive.getId();
+      if (id == null
+          || !adjustments.containsKey(id.intValue())
+          || primitive.getPrimitiveTypeName() != PrimitiveTypeName.INT32) {
+        return primitive;
+      }
+      return org.apache.parquet.schema.Types.primitive(
+              primitive.getPrimitiveTypeName(), primitive.getRepetition())
+          .id(id.intValue())
+          .named(primitive.getName());
+    }
+
+    @Override
+    public Type convertGroupType(List<GroupType> path, GroupType group, 
List<Type> children) {
+      return group.withNewFields(children);
+    }
+
+    @Override
+    public Type convertMessageType(MessageType message, List<Type> children) {
+      return new MessageType(message.getName(), children);
+    }
+  }
+
+  static Metrics apply(
+      Metrics metrics, Map<Integer, BoundAdjustment> adjustments, boolean 
roundUpperUp) {
+    Map<Integer, ByteBuffer> lower = metrics.lowerBounds();
+    Map<Integer, ByteBuffer> upper = metrics.upperBounds();
+    if (lower == null || upper == null) {
+      return metrics;
+    }
+    return new Metrics(
+        metrics.recordCount(),
+        metrics.columnSizes(),
+        metrics.valueCounts(),
+        metrics.nullValueCounts(),
+        metrics.nanValueCounts(),
+        adjust(lower, adjustments, false),
+        adjust(upper, adjustments, roundUpperUp));
+  }
+
+  private static Map<Integer, ByteBuffer> adjust(
+      Map<Integer, ByteBuffer> bounds, Map<Integer, BoundAdjustment> 
adjustments, boolean roundUp) {
+    Map<Integer, ByteBuffer> adjusted = new HashMap<>(bounds);
+    for (Map.Entry<Integer, BoundAdjustment> entry : adjustments.entrySet()) {
+      ByteBuffer bytes = bounds.get(entry.getKey());
+      if (bytes == null || entry.getValue() == UINT32_TO_INT) {
+        continue;
+      }
+      try {
+        long value = entry.getValue().convert(bytes, roundUp);
+        adjusted.put(entry.getKey(), 
Conversions.toByteBuffer(Types.LongType.get(), value));
+      } catch (ArithmeticException e) {
+        // Beyond the table unit's range (e.g. year 9999 in nanos): a missing 
bound is safe, a
+        // wrapped one is not.
+        adjusted.remove(entry.getKey());
+      }
+    }
+    return adjusted;
+  }
+
+  private long convert(ByteBuffer bytes, boolean roundUp) {
+    ByteBuffer little = bytes.duplicate().order(ByteOrder.LITTLE_ENDIAN);
+    switch (this) {
+      case UINT32_TO_LONG:
+        return Integer.toUnsignedLong(little.getInt(little.position()));
+      case MILLIS_TO_MICROS:
+        return Math.multiplyExact(readLong(little), 1000L);
+      case MILLIS_TO_NANOS:
+        return Math.multiplyExact(readLong(little), 1_000_000L);
+      case MICROS_TO_NANOS:
+        return Math.multiplyExact(readLong(little), 1000L);
+      case NANOS_TO_MICROS:
+        long nanos = little.getLong(little.position());
+        return roundUp ? -Math.floorDiv(-nanos, 1000L) : Math.floorDiv(nanos, 
1000L);
+      default:
+        throw new IllegalStateException(name());
+    }
+  }
+
+  /** A millis TIME bound has 4 bytes: its INT32 column is presented without 
the annotation. */
+  private static long readLong(ByteBuffer little) {
+    if (little.remaining() == 4) {
+      return little.getInt(little.position());
+    }
+    return little.getLong(little.position());
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/DryRunReport.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/DryRunReport.java
index ae069598e0a..7351c6f5c17 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/DryRunReport.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/DryRunReport.java
@@ -60,7 +60,7 @@ import org.slf4j.LoggerFactory;
  * allowed. Unreadable files always go to the error output in a real run; 
unchecked (ORC, Avro)
  * files count separately even when {@code unchecked_registered} says ACCEPT 
registers them. Pin
  * evidence is per file (the footer's null counts), so a pin violation or an 
unproven pin is not
- * predicted here.
+ * predicted here, nor is a file whose Parquet field ids disagree with the 
table.
  */
 class DryRunReport extends DoFn<List<CollectDistinctSchemas.SchemaGroup>, Row> 
{
   private static final Logger LOG = 
LoggerFactory.getLogger(DryRunReport.class);
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/NameMappingUtils.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/NameMappingUtils.java
index 1e8d31c03a3..c7af99e95d2 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/NameMappingUtils.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/NameMappingUtils.java
@@ -118,6 +118,18 @@ class NameMappingUtils {
     }
   }
 
+  /**
+   * The mapping readers will resolve a file without field ids with once it is 
committed: the stored
+   * one while it covers {@code schema}, else the one the commit regenerates 
from it. It keeps the
+   * old names of renamed columns, which a mapping created from the schema 
alone does not.
+   */
+  static NameMapping forReaders(Schema schema, @Nullable NameMapping stored) {
+    if (stored != null && covers(stored, schema.asStruct())) {
+      return stored;
+    }
+    return NameMappingParser.fromJson(regenerate(schema, stored));
+  }
+
   /**
    * Schema-derived mapping, with custom names (user aliases, pre-rename file 
names) carried over
    * from {@code existing} by field id. Schema names win: a carried name that 
collides with a name
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ParquetFieldIds.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ParquetFieldIds.java
new file mode 100644
index 00000000000..3c3711e3391
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/ParquetFieldIds.java
@@ -0,0 +1,185 @@
+/*
+ * 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.beam.sdk.io.iceberg;
+
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.TreeMap;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.mapping.MappedField;
+import org.apache.iceberg.mapping.MappingUtil;
+import org.apache.iceberg.mapping.NameMapping;
+import org.apache.iceberg.parquet.ParquetSchemaUtil;
+import org.apache.iceberg.types.TypeUtil;
+import org.apache.iceberg.types.Types;
+import org.apache.parquet.hadoop.metadata.FileMetaData;
+import org.apache.parquet.hadoop.metadata.ParquetMetadata;
+import org.apache.parquet.schema.MessageType;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * Decides, once per file, which table column each Parquet column is. Iceberg 
readers resolve a file
+ * by the field ids it carries and consult the table's name mapping only for a 
file that carries
+ * none, so a file's own ids must mean what the table means by them. 
Everything that reads a
+ * footer's columns by id takes the {@link Resolved} footer this produces.
+ */
+final class ParquetFieldIds {
+  private ParquetFieldIds() {}
+
+  /** A footer whose schema carries the table's field ids. */
+  static final class Resolved {
+    private final ParquetMetadata footer;
+
+    private Resolved(ParquetMetadata footer) {
+      this.footer = footer;
+    }
+
+    ParquetMetadata footer() {
+      return footer;
+    }
+
+    /**
+     * The file's dotted names for the given table field ids, skipping ids the 
file lacks. Iceberg
+     * looks metrics modes up by these names, which differ from the table's 
after a rename.
+     */
+    List<String> columnNames(List<Integer> fieldIds) {
+      Schema fileSchema = 
ParquetSchemaUtil.convertAndPrune(footer.getFileMetaData().getSchema());
+      List<String> names = new ArrayList<>();
+      for (int fieldId : fieldIds) {
+        @Nullable String name = fileSchema.findColumnName(fieldId);
+        if (name != null) {
+          names.add(name);
+        }
+      }
+      return names;
+    }
+  }
+
+  /** The file's own field ids disagree with the table; readers would misread 
its columns. */
+  static final class ConflictException extends IllegalArgumentException {
+    ConflictException(String message) {
+      super(message);
+    }
+  }
+
+  /**
+   * A file without ids gets the table's through {@code mapping}, which must 
be the one readers will
+   * use ({@link NameMappingUtils#forReaders}). A file with ids keeps them 
when they agree with the
+   * table and is refused otherwise: no metrics or mapping written at 
registration can change how
+   * readers resolve it.
+   */
+  static Resolved resolve(ParquetMetadata footer, Table table, NameMapping 
mapping) {
+    return resolve(footer, table.schema(), table.schemas().values(), mapping);
+  }
+
+  /** Resolves against one schema alone: no earlier versions, and a mapping 
made from it. */
+  @VisibleForTesting
+  static Resolved resolve(ParquetMetadata footer, Schema schema) {
+    return resolve(footer, schema, Collections.singletonList(schema), 
MappingUtil.create(schema));
+  }
+
+  private static Resolved resolve(
+      ParquetMetadata footer, Schema current, Collection<Schema> versions, 
NameMapping mapping) {
+    MessageType fileType = footer.getFileMetaData().getSchema();
+    if (!ParquetSchemaUtil.hasIds(fileType)) {
+      MessageType mapped = ParquetSchemaUtil.applyNameMapping(fileType, 
mapping);
+      FileMetaData meta = footer.getFileMetaData();
+      return new Resolved(
+          new ParquetMetadata(
+              new FileMetaData(mapped, meta.getKeyValueMetaData(), 
meta.getCreatedBy()),
+              footer.getBlocks()));
+    }
+    @Nullable String conflict = conflict(fileType, current, versions, mapping);
+    if (conflict != null) {
+      throw new ConflictException(conflict);
+    }
+    return new Resolved(footer);
+  }
+
+  /**
+   * An id agrees when some version of the table's schema, or the name 
mapping, gives it the file
+   * column's full path, so files written before a rename stay valid. 
Comparing only the last name
+   * would accept a shipping.city that carries billing.city's id, which 
readers find under neither.
+   * An id the table has never used agrees too, unless the column's path 
belongs to a table column
+   * with another id, which readers would then read as null.
+   */
+  private static @Nullable String conflict(
+      MessageType fileType, Schema current, Collection<Schema> versions, 
NameMapping mapping) {
+    Schema fileSchema;
+    try {
+      fileSchema = ParquetSchemaUtil.convert(fileType);
+    } catch (RuntimeException e) {
+      return "its field ids cannot be resolved: " + AddFiles.errorMessage(e);
+    }
+    Map<Integer, Types.NestedField> byId = new 
TreeMap<>(TypeUtil.indexById(fileSchema.asStruct()));
+    for (Types.NestedField field : byId.values()) {
+      int id = field.fieldId();
+      String path = fileSchema.findColumnName(id);
+      if (knownAs(current, versions, mapping, id, path)) {
+        continue;
+      }
+      Types.@Nullable NestedField sameId = current.findField(id);
+      if (sameId != null) {
+        return "column "
+            + path
+            + " carries field id "
+            + id
+            + ", which the table uses for column "
+            + current.findColumnName(id);
+      }
+      Types.@Nullable NestedField sameName = current.findField(path);
+      if (sameName != null) {
+        return "column "
+            + path
+            + " carries field id "
+            + id
+            + ", but the table's column "
+            + path
+            + " has field id "
+            + sameName.fieldId();
+      }
+    }
+    return null;
+  }
+
+  /** Paths are dotted full names, as both Iceberg schemas and name mappings 
index them. */
+  private static boolean knownAs(
+      Schema current, Collection<Schema> versions, NameMapping mapping, int 
id, String path) {
+    // The current schema first: versions come oldest first, and walking them 
for every column
+    // costs columns x versions per file.
+    if (path.equals(current.findColumnName(id))) {
+      return true;
+    }
+    for (Schema schema : versions) {
+      if (path.equals(schema.findColumnName(id))) {
+        return true;
+      }
+    }
+    @Nullable MappedField mapped = mapping.find(path);
+    if (mapped == null) {
+      return false;
+    }
+    @Nullable Integer mappedId = mapped.id();
+    return mappedId != null && mappedId == id;
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesMetricsTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesMetricsTest.java
new file mode 100644
index 00000000000..cb85908450e
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesMetricsTest.java
@@ -0,0 +1,680 @@
+/*
+ * 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.beam.sdk.io.iceberg;
+
+import static 
org.apache.beam.sdk.io.iceberg.AddFiles.ConvertToDataFile.FIELD_ID_ERROR;
+import static 
org.apache.beam.sdk.io.iceberg.AddFiles.ConvertToDataFile.UNKNOWN_PARTITION_ERROR;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.EPOCH_SECONDS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.FLAG;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.FULL_METRICS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.ID;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.NAME;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.TS_INT96;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.TS_MICROS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.TS_MILLIS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.TS_NANOS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.UNSIGNED;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.column;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.row;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import org.apache.beam.sdk.testing.PAssert;
+import org.apache.beam.sdk.testing.TestPipeline;
+import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.values.PCollectionRowTuple;
+import org.apache.beam.sdk.values.Row;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.iceberg.DataFile;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.data.parquet.GenericParquetReaders;
+import org.apache.iceberg.hadoop.HadoopCatalog;
+import org.apache.iceberg.io.CloseableIterable;
+import org.apache.iceberg.mapping.MappingUtil;
+import org.apache.iceberg.mapping.NameMappingParser;
+import org.apache.iceberg.parquet.Parquet;
+import org.apache.iceberg.types.Conversions;
+import org.apache.iceberg.types.Types;
+import org.apache.parquet.example.data.simple.NanoTime;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
+import org.apache.parquet.schema.PrimitiveType;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.Before;
+import org.junit.ClassRule;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+import org.junit.rules.TestName;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * Metrics-based partition inference, end to end, for files registered without 
a location prefix: a
+ * file lands in the partition readers find its rows in, or goes to the error 
output when that
+ * partition cannot be known. The bounds it reads are covered by {@link 
BoundAdjustmentTest}, and
+ * which table column each file column is by {@link ParquetFieldIdsTest}.
+ */
+@RunWith(JUnit4.class)
+public class AddFilesMetricsTest {
+  @ClassRule public static final TemporaryFolder TEMPORARY_FOLDER = new 
TemporaryFolder();
+  @Rule public TemporaryFolder temp = new TemporaryFolder();
+  @Rule public TestPipeline pipeline = TestPipeline.create();
+  @Rule public TestName testName = new TestName();
+
+  @Rule
+  public transient TestDataWarehouse warehouse = new 
TestDataWarehouse(TEMPORARY_FOLDER, "default");
+
+  /** Days from the epoch to 2024-01-01. */
+  private static final int EPOCH_DAY = 19723;
+
+  private static final Schema ID_FLAG =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "flag", Types.BooleanType.get()));
+  private static final Schema ID_NAME =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "name", Types.StringType.get()));
+  private static final Schema ID_TS =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "ts", Types.TimestampType.withZone()));
+  private static final Schema ID_TS_NS =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "ts", 
Types.TimestampNanoType.withZone()));
+  private static final Schema ID_NAME_AGE =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "name", Types.StringType.get()),
+          Types.NestedField.optional(3, "age", Types.IntegerType.get()));
+  private static final Schema UNSIGNED_ONLY =
+      new Schema(Types.NestedField.optional(1, "u", Types.LongType.get()));
+  private static final Schema ID_U_INT =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "u", Types.IntegerType.get()));
+  private static final Schema ID_FLAG_UNSIGNED =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "flag", Types.BooleanType.get()),
+          Types.NestedField.optional(3, "u", Types.LongType.get()));
+
+  private HadoopCatalog catalog;
+  private TableIdentifier tableId;
+  private IcebergCatalogConfig catalogConfig;
+  private ParquetTestFiles files;
+
+  @Before
+  public void setup() {
+    files = new ParquetTestFiles(temp.getRoot());
+    catalog = new HadoopCatalog(new Configuration(), warehouse.location);
+    tableId = TableIdentifier.of("default", testName.getMethodName());
+    catalogConfig =
+        IcebergCatalogConfig.builder()
+            .setCatalogProperties(
+                ImmutableMap.of("type", "hadoop", "warehouse", 
warehouse.location))
+            .build();
+  }
+
+  /** One Avro row {@code {id: 1}}; Avro files carry no column metrics. */
+  private String writeAvro(String name) throws IOException {
+    File file = new File(temp.getRoot(), name);
+    org.apache.avro.Schema avro =
+        
org.apache.avro.SchemaBuilder.record("r").fields().requiredInt("id").endRecord();
+    try 
(org.apache.avro.file.DataFileWriter<org.apache.avro.generic.GenericRecord> 
writer =
+        new org.apache.avro.file.DataFileWriter<>(
+            new org.apache.avro.generic.GenericDatumWriter<>(avro))) {
+      writer.create(avro, file);
+      org.apache.avro.generic.GenericData.Record record =
+          new org.apache.avro.generic.GenericData.Record(avro);
+      record.put("id", 1);
+      writer.append(record);
+    }
+    return file.getAbsolutePath();
+  }
+
+  // ---- metrics-based partition inference, end to end
+
+  private PCollectionRowTuple register(String... files) {
+    return pipeline
+        .apply("Create Input", Create.of(Arrays.asList(files)))
+        .apply(new AddFiles(catalogConfig, tableId.toString(), null, null, 
null, null, null, null));
+  }
+
+  private static void expectNoErrors(PCollectionRowTuple output) {
+    PAssert.that(output.get("errors")).empty();
+  }
+
+  /** The only error row is an unknown partition for {@code file} that names 
{@code column}. */
+  private static void expectUnknownPartition(
+      PCollectionRowTuple output, String file, String column) {
+    PAssert.that(output.get("errors"))
+        .satisfies(
+            rows -> {
+              Row error = Iterables.getOnlyElement(rows);
+              String message = String.valueOf(error.getString("error"));
+              assertEquals(file, error.getString("file"));
+              assertTrue(message, message.startsWith(UNKNOWN_PARTITION_ERROR));
+              assertTrue(message, message.contains(column));
+              return null;
+            });
+  }
+
+  private List<DataFile> registeredFiles() {
+    Table table = catalog.loadTable(tableId);
+    if (table.currentSnapshot() == null) {
+      return new ArrayList<>();
+    }
+    return 
Lists.newArrayList(table.currentSnapshot().addedDataFiles(table.io()));
+  }
+
+  private @Nullable Object onlyPartitionValue() {
+    DataFile file = Iterables.getOnlyElement(registeredFiles());
+    return file.partition().get(0, Object.class);
+  }
+
+  /** The only error row is for {@code file}, and its message contains {@code 
fragment}. */
+  private static void expectError(PCollectionRowTuple output, String file, 
String fragment) {
+    PAssert.that(output.get("errors"))
+        .satisfies(
+            rows -> {
+              Row error = Iterables.getOnlyElement(rows);
+              String message = String.valueOf(error.getString("error"));
+              assertEquals(file, error.getString("file"));
+              assertTrue(message, message.contains(fragment));
+              return null;
+            });
+  }
+
+  /**
+   * The registered file's rows as {@code column=value} strings, read the way 
an engine reads them:
+   * through the table's stored name mapping. IcebergGenerics passes no 
mapping at all.
+   */
+  private List<String> readBack(String... columns) throws IOException {
+    Table table = catalog.loadTable(tableId);
+    Schema schema = table.schema();
+    DataFile file = Iterables.getOnlyElement(registeredFiles());
+    List<String> rows = new ArrayList<>();
+    try (CloseableIterable<Record> records =
+        Parquet.read(table.io().newInputFile(file.location()))
+            .project(schema)
+            .withNameMapping(
+                NameMappingParser.fromJson(
+                    
table.properties().get(TableProperties.DEFAULT_NAME_MAPPING)))
+            .createReaderFunc(fileSchema -> 
GenericParquetReaders.buildReader(schema, fileSchema))
+            .build()) {
+      for (Record record : records) {
+        List<String> values = new ArrayList<>();
+        for (String column : columns) {
+          values.add(column + "=" + record.getField(column));
+        }
+        rows.add(String.join(" ", values));
+      }
+    }
+    return rows;
+  }
+
+  @Test
+  public void testAllNullBooleanColumnRegistersUnderTheNullPartition() throws 
IOException {
+    catalog.createTable(
+        tableId, ID_FLAG, 
PartitionSpec.builderFor(ID_FLAG).identity("flag").build(), FULL_METRICS);
+    String file =
+        files.write("nulls.parquet", true, Arrays.asList(ID, FLAG), row(1, 
null), row(2, null));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertNull(onlyPartitionValue());
+  }
+
+  @Test
+  public void testAllNullStringColumnRegistersUnderTheNullPartition() throws 
IOException {
+    catalog.createTable(
+        tableId, ID_NAME, 
PartitionSpec.builderFor(ID_NAME).identity("name").build(), FULL_METRICS);
+    String file =
+        files.write("nulls.parquet", true, Arrays.asList(ID, NAME), row(1, 
null), row(2, null));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertNull(onlyPartitionValue());
+  }
+
+  @Test
+  public void testAllNullTimestampColumnRegistersUnderTheNullDayPartition() 
throws IOException {
+    catalog.createTable(
+        tableId, ID_TS, PartitionSpec.builderFor(ID_TS).day("ts").build(), 
FULL_METRICS);
+    String file =
+        files.write(
+            "nulls.parquet", true, Arrays.asList(ID, TS_MICROS), row(1, null), 
row(2, null));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertNull(onlyPartitionValue());
+  }
+
+  /** A file written before the partition column existed reads as null for 
every row. */
+  @Test
+  public void 
testFileLackingThePartitionColumnRegistersUnderTheNullPartition() throws 
IOException {
+    catalog.createTable(
+        tableId, ID_FLAG, 
PartitionSpec.builderFor(ID_FLAG).identity("flag").build(), FULL_METRICS);
+    String file = files.write("older.parquet", true, Arrays.asList(ID), 
row(1), row(2));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertNull(onlyPartitionValue());
+  }
+
+  /** Null rows belong to the null partition, so such a file spans two 
partitions. */
+  @Test
+  public void testPartitionColumnWithNullsAndValuesIsAnUnknownPartition() 
throws IOException {
+    catalog.createTable(
+        tableId, ID_FLAG, 
PartitionSpec.builderFor(ID_FLAG).identity("flag").build(), FULL_METRICS);
+    String file =
+        files.write("mixed.parquet", true, Arrays.asList(ID, FLAG), row(1, 
true), row(2, null));
+
+    expectUnknownPartition(register(file), file, "flag");
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(Collections.emptyList(), registeredFiles());
+  }
+
+  @Test
+  public void testPartitionColumnWithoutStatisticsIsAnUnknownPartition() 
throws IOException {
+    catalog.createTable(
+        tableId, ID_FLAG, 
PartitionSpec.builderFor(ID_FLAG).identity("flag").build(), FULL_METRICS);
+    String file =
+        files.write("nostats.parquet", false, Arrays.asList(ID, FLAG), row(1, 
true), row(2, true));
+
+    expectUnknownPartition(register(file), file, "flag");
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(Collections.emptyList(), registeredFiles());
+  }
+
+  /** Iceberg collects no bounds for INT96; such a file must not land in the 
null partition. */
+  @Test
+  public void testInt96PartitionColumnIsAnUnknownPartition() throws 
IOException {
+    catalog.createTable(
+        tableId, ID_TS, PartitionSpec.builderFor(ID_TS).day("ts").build(), 
FULL_METRICS);
+    String file =
+        files.write(
+            "int96.parquet", true, Arrays.asList(ID, TS_INT96), row(1, new 
NanoTime(2460311, 0L)));
+
+    expectUnknownPartition(register(file), file, "ts");
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(Collections.emptyList(), registeredFiles());
+  }
+
+  @Test
+  public void testAvroFileOnAPartitionedTableIsAnUnknownPartition() throws 
IOException {
+    catalog.createTable(
+        tableId, ID_FLAG, 
PartitionSpec.builderFor(ID_FLAG).identity("flag").build(), FULL_METRICS);
+    String parquet = files.write("good.parquet", true, Arrays.asList(ID, 
FLAG), row(1, true));
+    String avro = writeAvro("rows.avro");
+
+    expectUnknownPartition(register(parquet, avro), avro, "flag");
+    pipeline.run().waitUntilFinish();
+
+    DataFile registered = Iterables.getOnlyElement(registeredFiles());
+    assertEquals(parquet, registered.location());
+    assertEquals(true, registered.partition().get(0, Object.class));
+  }
+
+  @Test
+  public void testMillisTimestampFileLandsInItsDayPartition() throws 
IOException {
+    catalog.createTable(
+        tableId, ID_TS, PartitionSpec.builderFor(ID_TS).day("ts").build(), 
FULL_METRICS);
+    String file =
+        files.write(
+            "millis.parquet",
+            true,
+            Arrays.asList(ID, TS_MILLIS),
+            row(1, EPOCH_SECONDS * 1000L),
+            row(2, EPOCH_SECONDS * 1000L + 1));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(EPOCH_DAY, onlyPartitionValue());
+  }
+
+  @Test
+  public void testNanosFileLandsInItsDayPartitionUnderANanosColumn() throws 
IOException {
+    catalog.createTable(
+        tableId,
+        ID_TS_NS,
+        PartitionSpec.builderFor(ID_TS_NS).day("ts").build(),
+        ImmutableMap.of("format-version", "3", 
"write.metadata.metrics.default", "full"));
+    String file =
+        files.write(
+            "nanos.parquet",
+            true,
+            Arrays.asList(ID, TS_NANOS),
+            row(1, EPOCH_SECONDS * 1_000_000_000L + 1),
+            row(2, EPOCH_SECONDS * 1_000_000_000L + 1_500));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(EPOCH_DAY, onlyPartitionValue());
+  }
+
+  @Test
+  public void testUnsigned32FileWithValuesFrom2To31GoesToTheErrorOutput() 
throws IOException {
+    catalog.createTable(tableId, UNSIGNED_ONLY, PartitionSpec.unpartitioned(), 
FULL_METRICS);
+    String file =
+        files.write(
+            "unsigned.parquet",
+            true,
+            2,
+            Arrays.asList(UNSIGNED),
+            row(5),
+            row(10),
+            row((int) 3_000_000_000L),
+            row((int) 3_000_000_005L));
+
+    expectError(register(file), file, BoundAdjustment.UNSIGNED_RANGE_ERROR);
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(Collections.emptyList(), registeredFiles());
+  }
+
+  /**
+   * Iceberg's default metrics mode, truncate(16), stores a 23-character value 
as two different
+   * bounds, which are a range for pruning but cannot name the partition.
+   */
+  @Test
+  public void testLongStringUnderDefaultMetricsLandsInItsIdentityPartition() 
throws IOException {
+    catalog.createTable(
+        tableId, ID_NAME, 
PartitionSpec.builderFor(ID_NAME).identity("name").build());
+    String value = "customer-00000000000001";
+    String file =
+        files.write("long.parquet", true, Arrays.asList(ID, NAME), row(1, 
value), row(2, value));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(value, String.valueOf(onlyPartitionValue()));
+  }
+
+  /** The partition comes from full bounds; the registered file keeps the 
table's own metrics. */
+  @Test
+  public void testTruncatedTableMetricsStillInferTheIdentityPartition() throws 
IOException {
+    catalog.createTable(
+        tableId,
+        ID_NAME,
+        PartitionSpec.builderFor(ID_NAME).identity("name").build(),
+        ImmutableMap.of("write.metadata.metrics.default", "truncate(4)"));
+    String file =
+        files.write(
+            "names.parquet", true, Arrays.asList(ID, NAME), row(1, 
"abcdefgh"), row(2, "abcdefgh"));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    DataFile registered = Iterables.getOnlyElement(registeredFiles());
+    assertEquals("abcdefgh", String.valueOf(registered.partition().get(0, 
Object.class)));
+    assertEquals(
+        "abcd",
+        Conversions.fromByteBuffer(Types.StringType.get(), 
registered.lowerBounds().get(2))
+            .toString());
+  }
+
+  // ---- files that carry field ids
+
+  /**
+   * Written for another table: its email column carries id 3, the table's 
age. Trusting that id put
+   * the file under age=779632737, the first four bytes of "[email protected]" read 
as an int. The refusal
+   * reaches the error output.
+   */
+  @Test
+  public void testForeignFieldIdsDoNotPlaceAFileInAPartition() throws 
IOException {
+    catalog.createTable(
+        tableId,
+        ID_NAME_AGE,
+        PartitionSpec.builderFor(ID_NAME_AGE).identity("age").build(),
+        FULL_METRICS);
+    PrimitiveType email =
+        column("email", PrimitiveTypeName.BINARY, 
LogicalTypeAnnotation.stringType()).withId(3);
+    String file =
+        files.write(
+            "foreign.parquet",
+            true,
+            Arrays.asList(NAME.withId(1), ID.withId(2), email),
+            row("alice", 1, "[email protected]"),
+            row("bob", 2, "[email protected]"));
+
+    expectError(register(file), file, FIELD_ID_ERROR);
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(Collections.emptyList(), registeredFiles());
+  }
+
+  /**
+   * Readers resolve a file without field ids through the table's stored name 
mapping, which keeps
+   * the old name of a renamed column, so the partition is inferred through 
that mapping too.
+   */
+  @Test
+  public void testFileWithoutIdsWrittenBeforeARenameLandsInItsPartition() 
throws IOException {
+    Table table =
+        catalog.createTable(
+            tableId,
+            ID_NAME,
+            PartitionSpec.builderFor(ID_NAME).identity("name").build(),
+            FULL_METRICS);
+    table
+        .updateProperties()
+        .set(
+            TableProperties.DEFAULT_NAME_MAPPING,
+            NameMappingParser.toJson(MappingUtil.create(table.schema())))
+        .commit();
+    table.updateSchema().renameColumn("name", "full_name").commit();
+    String file =
+        files.write(
+            "older.parquet", true, Arrays.asList(ID, NAME), row(1, "alice"), 
row(2, "alice"));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals("alice", String.valueOf(onlyPartitionValue()));
+    assertEquals(
+        Arrays.asList("id=1 full_name=alice", "id=2 full_name=alice"), 
readBack("id", "full_name"));
+  }
+
+  /**
+   * With metrics switched off for the table, the partition is inferred from 
metrics collected for
+   * the partition columns alone; an unrelated column must not be able to fail 
that.
+   */
+  @Test
+  public void testUnrelatedColumnDoesNotBreakInferenceUnderMetricsModeNone() 
throws IOException {
+    catalog.createTable(
+        tableId,
+        ID_FLAG_UNSIGNED,
+        PartitionSpec.builderFor(ID_FLAG_UNSIGNED).identity("flag").build(),
+        ImmutableMap.of("write.metadata.metrics.default", "none"));
+    String file =
+        files.write(
+            "unsigned.parquet",
+            true,
+            Arrays.asList(ID, FLAG, UNSIGNED),
+            row(1, true, 0),
+            row(2, true, 7));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(true, onlyPartitionValue());
+  }
+
+  // ---- a partition column's name, type or unit in the file differs from the 
table's
+
+  /**
+   * Iceberg looks metrics modes up by the file's column name, and a file 
written before a rename
+   * still uses the old one, so the recollection must ask for full metrics 
under that name.
+   */
+  @Test
+  public void 
testFileWithoutIdsWrittenBeforeARenameLandsInItsPartitionUnderDefaultMetrics()
+      throws IOException {
+    Table table =
+        catalog.createTable(
+            tableId, ID_NAME, 
PartitionSpec.builderFor(ID_NAME).identity("name").build());
+    table
+        .updateProperties()
+        .set(
+            TableProperties.DEFAULT_NAME_MAPPING,
+            NameMappingParser.toJson(MappingUtil.create(table.schema())))
+        .commit();
+    table.updateSchema().renameColumn("name", "full_name").commit();
+    String file =
+        files.write(
+            "older.parquet", true, Arrays.asList(ID, NAME), row(1, "alice"), 
row(2, "alice"));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals("alice", String.valueOf(onlyPartitionValue()));
+  }
+
+  @Test
+  public void 
testFileWithIdsWrittenBeforeARenameLandsInItsPartitionUnderDefaultMetrics()
+      throws IOException {
+    Table table =
+        catalog.createTable(
+            tableId, ID_NAME, 
PartitionSpec.builderFor(ID_NAME).identity("name").build());
+    table.updateSchema().renameColumn("name", "full_name").commit();
+    String file =
+        files.write(
+            "older.parquet",
+            true,
+            Arrays.asList(ID.withId(1), NAME.withId(2)),
+            row(1, "alice"),
+            row(2, "alice"));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals("alice", String.valueOf(onlyPartitionValue()));
+  }
+
+  /**
+   * Full mode set under the column's new name does not reach a file written 
before the rename:
+   * Iceberg stores that file's bounds truncated, so they must not be reused 
for the partition.
+   */
+  @Test
+  public void testTruncatedBoundsOfARenamedColumnAreNotReusedForItsPartition() 
throws IOException {
+    Table table =
+        catalog.createTable(
+            tableId, ID_NAME, 
PartitionSpec.builderFor(ID_NAME).identity("name").build());
+    table.updateSchema().renameColumn("name", "full_name").commit();
+    table.updateProperties().set("write.metadata.metrics.column.full_name", 
"full").commit();
+    String value = "customer-00000000000001";
+    String file =
+        files.write(
+            "older.parquet",
+            true,
+            Arrays.asList(ID.withId(1), NAME.withId(2)),
+            row(1, value),
+            row(2, value));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(value, String.valueOf(onlyPartitionValue()));
+  }
+
+  /**
+   * With metrics off for the table, the partition column is only read by the 
recollection; a uint32
+   * file column under an int table column used to throw there and fail the 
pipeline.
+   */
+  @Test
+  public void 
testUnsigned32UnderAnIntPartitionColumnRegistersUnderMetricsModeNone()
+      throws IOException {
+    catalog.createTable(
+        tableId,
+        ID_U_INT,
+        PartitionSpec.builderFor(ID_U_INT).identity("u").build(),
+        ImmutableMap.of("write.metadata.metrics.default", "none"));
+    String file = files.write("u.parquet", true, Arrays.asList(ID, UNSIGNED), 
row(1, 7), row(2, 7));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(7, onlyPartitionValue());
+  }
+
+  /**
+   * Rounding the upper bound up, as pruning needs, would put 
23:59:59.999999999 in the next day.
+   */
+  @Test
+  public void testNanosFileEndingInTheLastNanoOfADayLandsInThatDay() throws 
IOException {
+    catalog.createTable(
+        tableId, ID_TS, PartitionSpec.builderFor(ID_TS).day("ts").build(), 
FULL_METRICS);
+    long dayStartNanos = EPOCH_SECONDS * 1_000_000_000L;
+    String file =
+        files.write(
+            "nanos.parquet",
+            true,
+            Arrays.asList(ID, TS_NANOS),
+            row(1, dayStartNanos + 1_000_000_000L),
+            row(2, dayStartNanos + 86_400_000_000_000L - 1));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(EPOCH_DAY, onlyPartitionValue());
+  }
+
+  /** A micros column cannot hold the extra nanos: the partition is the 
value's micros, floored. */
+  @Test
+  public void testSubMicroNanosValueLandsInTheIdentityPartitionOfItsMicros() 
throws IOException {
+    catalog.createTable(
+        tableId, ID_TS, 
PartitionSpec.builderFor(ID_TS).identity("ts").build(), FULL_METRICS);
+    String file =
+        files.write(
+            "nanos.parquet",
+            true,
+            Arrays.asList(ID, TS_NANOS),
+            row(1, EPOCH_SECONDS * 1_000_000_000L + 1_500));
+
+    expectNoErrors(register(file));
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(EPOCH_SECONDS * 1_000_000L + 1, onlyPartitionValue());
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesTest.java
index b9fe94f27a5..0d5a41e0cc4 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/AddFilesTest.java
@@ -17,6 +17,7 @@
  */
 package org.apache.beam.sdk.io.iceberg;
 
+import static 
org.apache.beam.sdk.io.iceberg.AddFiles.ConvertToDataFile.FIELD_ID_ERROR;
 import static 
org.apache.beam.sdk.io.iceberg.AddFiles.ConvertToDataFile.PREFIX_ERROR;
 import static 
org.apache.beam.sdk.io.iceberg.AddFiles.ConvertToDataFile.UNKNOWN_PARTITION_ERROR;
 import static 
org.apache.beam.sdk.io.iceberg.AddFiles.ConvertToDataFile.getPartitionFromMetrics;
@@ -95,19 +96,21 @@ import org.apache.iceberg.TableProperties;
 import org.apache.iceberg.catalog.TableIdentifier;
 import org.apache.iceberg.data.GenericRecord;
 import org.apache.iceberg.data.Record;
+import org.apache.iceberg.data.parquet.GenericParquetReaders;
 import org.apache.iceberg.data.parquet.GenericParquetWriter;
 import org.apache.iceberg.exceptions.CommitFailedException;
 import org.apache.iceberg.hadoop.HadoopCatalog;
+import org.apache.iceberg.io.CloseableIterable;
 import org.apache.iceberg.io.DataWriter;
 import org.apache.iceberg.io.InputFile;
 import org.apache.iceberg.mapping.MappingUtil;
 import org.apache.iceberg.mapping.NameMapping;
 import org.apache.iceberg.mapping.NameMappingParser;
 import org.apache.iceberg.parquet.Parquet;
+import org.apache.iceberg.parquet.ParquetSchemaUtil;
 import org.apache.iceberg.types.Conversions;
 import org.apache.iceberg.types.Types;
 import org.apache.iceberg.util.SerializableFunction;
-import org.apache.parquet.hadoop.metadata.ParquetMetadata;
 import org.checkerframework.checker.nullness.qual.Nullable;
 import org.joda.time.Duration;
 import org.joda.time.Instant;
@@ -835,10 +838,12 @@ public class AddFilesTest {
       writer.close();
       InputFile file = table.io().newInputFile(fileName);
 
-      ParquetMetadata footer = ParquetFooters.read(fileName);
+      NameMapping mapping = MappingUtil.create(icebergSchema);
+      ParquetFieldIds.Resolved footer =
+          ParquetFieldIds.resolve(ParquetFooters.read(fileName), table, 
mapping);
       Metrics metrics =
           AddFiles.getFileMetrics(
-              file, FileFormat.PARQUET, metricsConfig, 
MappingUtil.create(icebergSchema), footer);
+              file, FileFormat.PARQUET, metricsConfig, mapping, icebergSchema, 
footer);
       for (int i = 0; i < partitionSpec.fields().size(); i++) {
         PartitionField partitionField = partitionSpec.fields().get(i);
         Types.NestedField field = 
icebergSchema.findField(partitionField.sourceId());
@@ -852,7 +857,8 @@ public class AddFilesTest {
         assertEquals(caze.expectedUpper.get(i), upper);
       }
 
-      String partitionPath = getPartitionFromMetrics(metrics, file, table, 
footer);
+      String partitionPath =
+          getPartitionFromMetrics(metrics, file, table, mapping, 
footer).toPath();
       assertEquals(caze.expectedPartition, partitionPath);
     }
   }
@@ -901,10 +907,12 @@ public class AddFilesTest {
       writer.close();
       InputFile file = table.io().newInputFile(fileName);
 
-      ParquetMetadata footer = ParquetFooters.read(fileName);
+      NameMapping mapping = MappingUtil.create(icebergSchema);
+      ParquetFieldIds.Resolved footer =
+          ParquetFieldIds.resolve(ParquetFooters.read(fileName), table, 
mapping);
       Metrics metrics =
           AddFiles.getFileMetrics(
-              file, FileFormat.PARQUET, metricsConfig, 
MappingUtil.create(icebergSchema), footer);
+              file, FileFormat.PARQUET, metricsConfig, mapping, icebergSchema, 
footer);
       // check that lower/upper stats are still fetched correctly
       for (int i = 0; i < partitionSpec.fields().size(); i++) {
         PartitionField partitionField = partitionSpec.fields().get(i);
@@ -921,7 +929,7 @@ public class AddFilesTest {
 
       assertThrows(
           AddFiles.UnknownPartitionException.class,
-          () -> getPartitionFromMetrics(metrics, file, table, footer));
+          () -> getPartitionFromMetrics(metrics, file, table, mapping, 
footer));
     }
   }
 
@@ -1425,8 +1433,13 @@ public class AddFilesTest {
 
   @Test
   public void testMissingTableIsCreatedFromTheFilesUnion() throws Exception {
-    String narrow = writeOneRecord("narrow.parquet");
-    String wide = writeWider("wide.parquet");
+    String narrow = writeWithoutFieldIds("narrow.parquet", icebergSchema, 
record(1, "a", 1));
+    Record wider = GenericRecord.create(WIDER);
+    wider.setField("id", 1);
+    wider.setField("name", "a");
+    wider.setField("age", 1);
+    wider.setField("email", "e");
+    String wide = writeWithoutFieldIds("wide.parquet", WIDER, wider);
 
     PCollectionRowTuple output =
         pipeline.apply("Create Input", Create.of(narrow, 
wide)).apply(addFiles(ADDITIONS));
@@ -1438,6 +1451,93 @@ public class AddFilesTest {
     assertTrue(
         "created columns are optional",
         catalog.loadTable(tableId).schema().findField("id").isOptional());
+    assertEquals(
+        Arrays.asList("id=1 name=a age=1 email=e", "id=1 name=a age=1 
email=null"),
+        readRows("id", "name", "age", "email"));
+  }
+
+  /**
+   * The created table numbers its columns by sorted name, not by the ids an 
Iceberg-written file
+   * carries, and readers resolve such a file by its own ids: registered, its 
id values read back as
+   * age and its names as id. It goes to the error output instead.
+   */
+  @Test
+  public void testFileCarryingOtherFieldIdsIsRefusedByTheTableCreatedForIt() 
throws Exception {
+    String narrow = writeOneRecord("narrow.parquet");
+
+    PCollectionRowTuple output =
+        pipeline.apply("Create Input", 
Create.of(narrow)).apply(addFiles(ADDITIONS));
+    PAssert.that(output.get("errors"))
+        .satisfies(
+            rows -> {
+              Row error = Iterables.getOnlyElement(rows);
+              assertEquals(narrow, error.getString("file"));
+              assertThat(error.getString("error"), 
containsString(FIELD_ID_ERROR));
+              return null;
+            });
+
+    pipeline.run().waitUntilFinish();
+
+    assertEquals(0, 
Iterables.size(catalog.loadTable(tableId).newScan().planFiles()));
+  }
+
+  /** Plain Parquet, as pyarrow or Beam's WriteToParquet write it: no field 
ids. */
+  private String writeWithoutFieldIds(String name, Schema schema, Record... 
records)
+      throws IOException {
+    org.apache.avro.Schema avro = 
org.apache.iceberg.avro.AvroSchemaUtil.convert(schema, "row");
+    String file = root + name;
+    try 
(org.apache.parquet.hadoop.ParquetWriter<org.apache.avro.generic.GenericData.Record>
+        writer =
+            org.apache.parquet.avro.AvroParquetWriter
+                .<org.apache.avro.generic.GenericData.Record>builder(
+                    new org.apache.hadoop.fs.Path(file))
+                .withSchema(avro)
+                .build()) {
+      for (Record record : records) {
+        org.apache.avro.generic.GenericData.Record row =
+            new org.apache.avro.generic.GenericData.Record(avro);
+        for (Types.NestedField field : schema.columns()) {
+          row.put(field.name(), record.getField(field.name()));
+        }
+        writer.write(row);
+      }
+    }
+    assertFalse(
+        "fixture must carry no field ids",
+        
ParquetSchemaUtil.hasIds(ParquetFooters.read(file).getFileMetaData().getSchema()));
+    return file;
+  }
+
+  /**
+   * The table's rows as sorted {@code column=value} strings, read the way an 
engine that honors the
+   * table's name mapping reads them; IcebergGenerics passes no mapping at all.
+   */
+  private List<String> readRows(String... columns) throws IOException {
+    Table table = catalog.loadTable(tableId);
+    Schema schema = table.schema();
+    NameMapping mapping = MappingUtil.create(schema);
+    List<String> rows = new ArrayList<>();
+    try (CloseableIterable<FileScanTask> tasks = table.newScan().planFiles()) {
+      for (FileScanTask task : tasks) {
+        try (CloseableIterable<Record> records =
+            Parquet.read(table.io().newInputFile(task.file().location()))
+                .project(schema)
+                .withNameMapping(mapping)
+                .createReaderFunc(
+                    fileSchema -> GenericParquetReaders.buildReader(schema, 
fileSchema))
+                .build()) {
+          for (Record record : records) {
+            List<String> values = new ArrayList<>();
+            for (String column : columns) {
+              values.add(column + "=" + record.getField(column));
+            }
+            rows.add(String.join(" ", values));
+          }
+        }
+      }
+    }
+    Collections.sort(rows);
+    return rows;
   }
 
   /** Streaming schema evolution comes in a follow-up: until then the front 
door rejects it. */
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundAdjustmentTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundAdjustmentTest.java
new file mode 100644
index 00000000000..2a2e6c79711
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundAdjustmentTest.java
@@ -0,0 +1,391 @@
+/*
+ * 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.beam.sdk.io.iceberg;
+
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.EPOCH_SECONDS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.FULL_METRICS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.TS_MICROS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.TS_MILLIS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.TS_NANOS;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.UNSIGNED;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.column;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.row;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.Arrays;
+import java.util.Map;
+import org.apache.iceberg.Metrics;
+import org.apache.iceberg.MetricsConfig;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.mapping.MappingUtil;
+import org.apache.iceberg.parquet.ParquetSchemaUtil;
+import org.apache.iceberg.types.Conversions;
+import org.apache.iceberg.types.Types;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
+import org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit;
+import org.apache.parquet.schema.PrimitiveType;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.Before;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * Parquet stores a column's bounds in the file's unit; readers and partition 
inference decode them
+ * with the table column's type. Each test writes a file and checks the bounds 
Iceberg would store.
+ */
+@RunWith(JUnit4.class)
+public class BoundAdjustmentTest {
+  @Rule public TemporaryFolder temp = new TemporaryFolder();
+
+  private static final Schema ID_TS =
+      new Schema(Types.NestedField.optional(1, "ts", 
Types.TimestampType.withZone()));
+
+  private static final Schema TS_NS =
+      new Schema(Types.NestedField.optional(1, "ts", 
Types.TimestampNanoType.withZone()));
+
+  private ParquetTestFiles files;
+
+  @Before
+  public void setup() {
+    files = new ParquetTestFiles(temp.getRoot());
+  }
+
+  private static Metrics metricsOf(String file, Schema tableSchema) throws 
IOException {
+    return BoundAdjustment.footerMetrics(
+        ParquetFieldIds.resolve(ParquetFooters.read(file), tableSchema),
+        tableSchema,
+        MetricsConfig.fromProperties(FULL_METRICS),
+        MappingUtil.create(tableSchema));
+  }
+
+  private static Schema convertedSchema(String file) throws IOException {
+    return 
ParquetSchemaUtil.convert(ParquetFooters.read(file).getFileMetaData().getSchema());
+  }
+
+  private static @Nullable Object lower(Metrics metrics, Schema schema) {
+    return decode(metrics.lowerBounds(), schema);
+  }
+
+  private static @Nullable Object upper(Metrics metrics, Schema schema) {
+    return decode(metrics.upperBounds(), schema);
+  }
+
+  private static @Nullable Object decode(@Nullable Map<Integer, ByteBuffer> 
bounds, Schema schema) {
+    Types.NestedField field = schema.columns().get(0);
+    ByteBuffer bytes = bounds == null ? null : bounds.get(field.fieldId());
+    return bytes == null ? null : Conversions.fromByteBuffer(field.type(), 
bytes);
+  }
+
+  @Test
+  public void testMillisTimestampBoundsAreMicros() throws IOException {
+    String file =
+        files.write(
+            "millis.parquet",
+            true,
+            Arrays.asList(TS_MILLIS),
+            row(EPOCH_SECONDS * 1000L),
+            row(EPOCH_SECONDS * 1000L + 1));
+    Schema schema = convertedSchema(file);
+    assertEquals(Types.TimestampType.withZone(), 
schema.columns().get(0).type());
+
+    Metrics metrics = metricsOf(file, schema);
+
+    assertEquals(EPOCH_SECONDS * 1_000_000L, lower(metrics, schema));
+    assertEquals(EPOCH_SECONDS * 1_000_000L + 1_000L, upper(metrics, schema));
+  }
+
+  @Test
+  public void testNanosTimestampBoundsAreMicrosAndConservative() throws 
IOException {
+    PrimitiveType nanos =
+        column(
+            "ts",
+            PrimitiveTypeName.INT64,
+            LogicalTypeAnnotation.timestampType(false, TimeUnit.NANOS));
+    String file =
+        files.write(
+            "nanos.parquet",
+            true,
+            Arrays.asList(nanos),
+            row(EPOCH_SECONDS * 1_000_000_000L + 1),
+            row(EPOCH_SECONDS * 1_000_000_000L + 1_500));
+    Schema schema = convertedSchema(file);
+    assertEquals(Types.TimestampType.withoutZone(), 
schema.columns().get(0).type());
+
+    Metrics metrics = metricsOf(file, schema);
+
+    // lower rounds down, upper rounds up, so the bounds still contain every 
value
+    assertEquals(EPOCH_SECONDS * 1_000_000L, lower(metrics, schema));
+    assertEquals(EPOCH_SECONDS * 1_000_000L + 2, upper(metrics, schema));
+  }
+
+  @Test
+  public void testMillisTimeBoundsAreMicros() throws IOException {
+    PrimitiveType millisTime =
+        column("t", PrimitiveTypeName.INT32, 
LogicalTypeAnnotation.timeType(true, TimeUnit.MILLIS));
+    String file =
+        files.write(
+            "time.parquet", true, Arrays.asList(millisTime), row(3_600_000), 
row(3_600_001));
+    Schema schema = convertedSchema(file);
+    assertEquals(Types.TimeType.get(), schema.columns().get(0).type());
+
+    Metrics metrics = metricsOf(file, schema);
+
+    assertEquals(3_600_000_000L, lower(metrics, schema));
+    assertEquals(3_600_001_000L, upper(metrics, schema));
+  }
+
+  @Test
+  public void testNanosTimeBoundsAreMicros() throws IOException {
+    PrimitiveType nanosTime =
+        column("t", PrimitiveTypeName.INT64, 
LogicalTypeAnnotation.timeType(true, TimeUnit.NANOS));
+    String file =
+        files.write(
+            "nanotime.parquet",
+            true,
+            Arrays.asList(nanosTime),
+            row(3_600_000_000_001L),
+            row(3_600_000_000_999L));
+    Schema schema = convertedSchema(file);
+
+    Metrics metrics = metricsOf(file, schema);
+
+    assertEquals(3_600_000_000L, lower(metrics, schema));
+    assertEquals(3_600_000_001L, upper(metrics, schema));
+  }
+
+  /** With evolution off the table may declare the column int: bounds stay 
ints. */
+  @Test
+  public void testUnsigned32BoundsUnderAnIntColumnStayInts() throws 
IOException {
+    String file = files.write("unsigned.parquet", true, 
Arrays.asList(UNSIGNED), row(7), row(9));
+    Schema schema = new Schema(Types.NestedField.optional(1, "u", 
Types.IntegerType.get()));
+
+    Metrics metrics = metricsOf(file, schema);
+
+    assertEquals(7, lower(metrics, schema));
+    assertEquals(9, upper(metrics, schema));
+  }
+
+  @Test
+  public void testUnsigned32ValueFrom2To31UnderAnIntColumnIsRefused() throws 
IOException {
+    String file =
+        files.write(
+            "unsigned.parquet", true, Arrays.asList(UNSIGNED), row(7), 
row((int) 3_000_000_000L));
+    Schema schema = new Schema(Types.NestedField.optional(1, "u", 
Types.IntegerType.get()));
+
+    IllegalArgumentException e =
+        assertThrows(IllegalArgumentException.class, () -> metricsOf(file, 
schema));
+    assertTrue(e.getMessage(), 
e.getMessage().startsWith(BoundAdjustment.UNSIGNED_RANGE_ERROR));
+  }
+
+  /** Iceberg would cast the int statistics to the long it maps uint32 to, and 
throw. */
+  @Test
+  public void testUnsigned32UnderAStringColumnIsRefused() throws IOException {
+    String file = files.write("unsigned.parquet", true, 
Arrays.asList(UNSIGNED), row(7), row(9));
+    Schema schema = new Schema(Types.NestedField.optional(1, "u", 
Types.StringType.get()));
+
+    IllegalArgumentException e =
+        assertThrows(IllegalArgumentException.class, () -> metricsOf(file, 
schema));
+    assertEquals(BoundAdjustment.UNSIGNED_TYPE_ERROR + "u is string", 
e.getMessage());
+  }
+
+  @Test
+  public void testUnsigned32BoundsAreLongs() throws IOException {
+    String file =
+        files.write(
+            "unsigned.parquet", true, Arrays.asList(UNSIGNED), row(0), 
row(Integer.MAX_VALUE));
+    Schema schema = convertedSchema(file);
+    assertEquals(Types.LongType.get(), schema.columns().get(0).type());
+
+    Metrics metrics = metricsOf(file, schema);
+
+    assertEquals(0L, lower(metrics, schema));
+    assertEquals((long) Integer.MAX_VALUE, upper(metrics, schema));
+  }
+
+  /** Below 2^31 signed and unsigned order agree, so folding the row groups is 
safe. */
+  @Test
+  public void testUnsigned32BoundsSpanRowGroupsBelow2To31() throws IOException 
{
+    String file =
+        files.write(
+            "unsigned.parquet",
+            true,
+            2,
+            Arrays.asList(UNSIGNED),
+            row(5),
+            row(10),
+            row(2_000_000_000),
+            row(2_000_000_005));
+    assertEquals(2, ParquetFooters.read(file).getBlocks().size());
+    Schema schema = convertedSchema(file);
+
+    Metrics metrics = metricsOf(file, schema);
+
+    assertEquals(5L, lower(metrics, schema));
+    assertEquals(2_000_000_005L, upper(metrics, schema));
+  }
+
+  /**
+   * The raw int -1 is 4294967295 as a uint32, and Iceberg readers would 
return it as -1, so the
+   * file is refused.
+   */
+  @Test
+  public void testUnsigned32ValueFrom2To31IsRefused() throws IOException {
+    String file = files.write("unsigned.parquet", true, 
Arrays.asList(UNSIGNED), row(0), row(-1));
+    Schema schema = convertedSchema(file);
+
+    IllegalArgumentException e =
+        assertThrows(IllegalArgumentException.class, () -> metricsOf(file, 
schema));
+    assertTrue(e.getMessage(), 
e.getMessage().startsWith(BoundAdjustment.UNSIGNED_RANGE_ERROR));
+  }
+
+  /** Every row group is checked: only the second one holds values of 2^31 or 
more. */
+  @Test
+  public void testUnsigned32ValuesFrom2To31AcrossRowGroupsAreRefused() throws 
IOException {
+    String file =
+        files.write(
+            "unsigned.parquet",
+            true,
+            2,
+            Arrays.asList(UNSIGNED),
+            row(5),
+            row(10),
+            row((int) 3_000_000_000L),
+            row((int) 3_000_000_005L));
+    assertEquals(2, ParquetFooters.read(file).getBlocks().size());
+    Schema schema = convertedSchema(file);
+
+    IllegalArgumentException e =
+        assertThrows(IllegalArgumentException.class, () -> metricsOf(file, 
schema));
+    assertTrue(e.getMessage(), 
e.getMessage().startsWith(BoundAdjustment.UNSIGNED_RANGE_ERROR));
+  }
+
+  @Test
+  public void testMicrosTimestampBoundsAreUnchanged() throws IOException {
+    String file =
+        files.write(
+            "micros.parquet",
+            true,
+            Arrays.asList(TS_MICROS),
+            row(EPOCH_SECONDS * 1_000_000L),
+            row(EPOCH_SECONDS * 1_000_000L + 7));
+    Schema schema = convertedSchema(file);
+
+    Metrics metrics = metricsOf(file, schema);
+
+    assertEquals(EPOCH_SECONDS * 1_000_000L, lower(metrics, schema));
+    assertEquals(EPOCH_SECONDS * 1_000_000L + 7, upper(metrics, schema));
+  }
+
+  @Test
+  public void testNanosTimestampBoundsAreUnchangedUnderANanosColumn() throws 
IOException {
+    String file =
+        files.write(
+            "nanos.parquet",
+            true,
+            Arrays.asList(TS_NANOS),
+            row(EPOCH_SECONDS * 1_000_000_000L + 1),
+            row(EPOCH_SECONDS * 1_000_000_000L + 1_500));
+
+    Metrics metrics = metricsOf(file, TS_NS);
+
+    assertEquals(EPOCH_SECONDS * 1_000_000_000L + 1, lower(metrics, TS_NS));
+    assertEquals(EPOCH_SECONDS * 1_000_000_000L + 1_500, upper(metrics, 
TS_NS));
+  }
+
+  @Test
+  public void testMillisTimestampBoundsAreNanosUnderANanosColumn() throws 
IOException {
+    String file =
+        files.write(
+            "millis.parquet",
+            true,
+            Arrays.asList(TS_MILLIS),
+            row(EPOCH_SECONDS * 1000L),
+            row(EPOCH_SECONDS * 1000L + 1));
+
+    Metrics metrics = metricsOf(file, TS_NS);
+
+    assertEquals(EPOCH_SECONDS * 1_000_000_000L, lower(metrics, TS_NS));
+    assertEquals(EPOCH_SECONDS * 1_000_000_000L + 1_000_000L, upper(metrics, 
TS_NS));
+  }
+
+  @Test
+  public void testMicrosTimestampBoundsAreNanosUnderANanosColumn() throws 
IOException {
+    String file =
+        files.write(
+            "micros.parquet",
+            true,
+            Arrays.asList(TS_MICROS),
+            row(EPOCH_SECONDS * 1_000_000L),
+            row(EPOCH_SECONDS * 1_000_000L + 7));
+
+    Metrics metrics = metricsOf(file, TS_NS);
+
+    assertEquals(EPOCH_SECONDS * 1_000_000_000L, lower(metrics, TS_NS));
+    assertEquals(EPOCH_SECONDS * 1_000_000_000L + 7_000L, upper(metrics, 
TS_NS));
+  }
+
+  /** Stored bounds round outward for pruning; bounds for partition inference 
round down. */
+  @Test
+  public void testPartitionBoundsRoundNanosDown() throws IOException {
+    String file =
+        files.write(
+            "nanos.parquet",
+            true,
+            Arrays.asList(TS_NANOS),
+            row(EPOCH_SECONDS * 1_000_000_000L + 1),
+            row(EPOCH_SECONDS * 1_000_000_000L + 1_500));
+    ParquetFieldIds.Resolved footer = 
ParquetFieldIds.resolve(ParquetFooters.read(file), ID_TS);
+    MetricsConfig config = MetricsConfig.fromProperties(FULL_METRICS);
+
+    Metrics stored =
+        BoundAdjustment.footerMetrics(footer, ID_TS, config, 
MappingUtil.create(ID_TS));
+    Metrics partition =
+        BoundAdjustment.partitionMetrics(footer, ID_TS, config, 
MappingUtil.create(ID_TS));
+
+    assertEquals(EPOCH_SECONDS * 1_000_000L + 2, upper(stored, ID_TS));
+    assertEquals(EPOCH_SECONDS * 1_000_000L, lower(partition, ID_TS));
+    assertEquals(EPOCH_SECONDS * 1_000_000L + 1, upper(partition, ID_TS));
+  }
+
+  @Test
+  public void testBoundBeyondTheNanosRangeIsDropped() throws IOException {
+    // 9999-12-31T23:59:59Z; nanos since 1970 overflow a long after 2262
+    long lastSecondMillis = 253402300799000L;
+    String file =
+        files.write(
+            "far.parquet",
+            true,
+            Arrays.asList(TS_MILLIS),
+            row(EPOCH_SECONDS * 1000L),
+            row(lastSecondMillis));
+
+    Metrics metrics = metricsOf(file, TS_NS);
+
+    assertEquals(EPOCH_SECONDS * 1_000_000_000L, lower(metrics, TS_NS));
+    assertNull(upper(metrics, TS_NS));
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/NameMappingUtilsTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/NameMappingUtilsTest.java
index 80f5069fab5..fa216f12e29 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/NameMappingUtilsTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/NameMappingUtilsTest.java
@@ -21,6 +21,7 @@ import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertSame;
 import static org.junit.Assert.assertTrue;
 
 import org.apache.iceberg.Schema;
@@ -422,4 +423,34 @@ public class NameMappingUtilsTest {
     assertEquals(1, merged.find("new_name").id().intValue());
     assertEquals(1, merged.find("old_name").id().intValue());
   }
+
+  @Test
+  public void testForReadersKeepsAStoredMappingThatCoversTheSchema() {
+    Schema schema = new Schema(Types.NestedField.required(1, "new_name", 
Types.IntegerType.get()));
+    NameMapping stored = mapping("[ {'field-id': 1, 'names': ['new_name', 
'old_name']} ]");
+
+    assertSame(stored, NameMappingUtils.forReaders(schema, stored));
+  }
+
+  /** The commit regenerates a mapping that misses a column; the old names it 
had are kept. */
+  @Test
+  public void testForReadersRegeneratesAStoredMappingThatMissesAColumn() {
+    Schema schema =
+        new Schema(
+            Types.NestedField.required(1, "new_name", Types.IntegerType.get()),
+            Types.NestedField.optional(2, "added", Types.StringType.get()));
+    NameMapping stored = mapping("[ {'field-id': 1, 'names': ['new_name', 
'old_name']} ]");
+
+    NameMapping forReaders = NameMappingUtils.forReaders(schema, stored);
+
+    assertEquals(1, forReaders.find("old_name").id().intValue());
+    assertEquals(2, forReaders.find("added").id().intValue());
+  }
+
+  @Test
+  public void testForReadersWithoutAStoredMappingUsesTheSchemas() {
+    assertEquals(
+        NameMappingParser.toJson(MappingUtil.create(FULL_SCHEMA)),
+        NameMappingParser.toJson(NameMappingUtils.forReaders(FULL_SCHEMA, 
null)));
+  }
 }
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetFieldIdsTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetFieldIdsTest.java
new file mode 100644
index 00000000000..4eb0fd5cc8b
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetFieldIdsTest.java
@@ -0,0 +1,412 @@
+/*
+ * 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.beam.sdk.io.iceberg;
+
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.ID;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.NAME;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.column;
+import static org.apache.beam.sdk.io.iceberg.ParquetTestFiles.row;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.data.parquet.GenericParquetReaders;
+import org.apache.iceberg.hadoop.HadoopCatalog;
+import org.apache.iceberg.io.CloseableIterable;
+import org.apache.iceberg.mapping.MappedField;
+import org.apache.iceberg.mapping.MappingUtil;
+import org.apache.iceberg.mapping.NameMapping;
+import org.apache.iceberg.mapping.NameMappingParser;
+import org.apache.iceberg.parquet.Parquet;
+import org.apache.iceberg.types.Types;
+import org.apache.parquet.hadoop.metadata.ParquetMetadata;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
+import org.apache.parquet.schema.PrimitiveType;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.Before;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+import org.junit.rules.TestName;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * Which table column each Parquet column is. Readers resolve a file by the 
field ids it carries,
+ * and by the table's name mapping only when it carries none, so each test 
also reads the file back
+ * the way readers do to show why it is accepted or refused.
+ */
+@RunWith(JUnit4.class)
+public class ParquetFieldIdsTest {
+  @Rule public TemporaryFolder temp = new TemporaryFolder();
+  @Rule public TestName testName = new TestName();
+
+  private static final Schema ID_NAME =
+      new Schema(
+          Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "name", Types.StringType.get()));
+
+  private static final PrimitiveType EMAIL =
+      column("email", PrimitiveTypeName.BINARY, 
LogicalTypeAnnotation.stringType());
+
+  private HadoopCatalog catalog;
+  private ParquetTestFiles files;
+
+  @Before
+  public void setup() throws IOException {
+    catalog = new HadoopCatalog(new Configuration(), 
temp.newFolder("warehouse").getAbsolutePath());
+    files = new ParquetTestFiles(temp.newFolder("files"));
+  }
+
+  /** Iceberg assigns fresh ids when it creates a table; tests read them back 
from the result. */
+  private Table table(Schema schema) {
+    return catalog.createTable(
+        TableIdentifier.of("default", testName.getMethodName()),
+        schema,
+        PartitionSpec.unpartitioned());
+  }
+
+  private static ParquetFieldIds.Resolved resolve(String file, Table table) 
throws IOException {
+    return ParquetFieldIds.resolve(
+        ParquetFooters.read(file), table, MappingUtil.create(table.schema()));
+  }
+
+  private static String refusal(String file, Table table) {
+    ParquetFieldIds.ConflictException e =
+        assertThrows(ParquetFieldIds.ConflictException.class, () -> 
resolve(file, table));
+    return String.valueOf(e.getMessage());
+  }
+
+  /** The field id the resolved footer gives a top-level column, or null for 
none. */
+  private static @Nullable Integer idOf(ParquetFieldIds.Resolved resolved, 
String column) {
+    org.apache.parquet.schema.Type.ID id =
+        
resolved.footer().getFileMetaData().getSchema().getType(column).getId();
+    return id == null ? null : id.intValue();
+  }
+
+  /** The file's rows as readers return them: by the file's ids, or by {@code 
mapping} without. */
+  private static List<String> read(String file, Schema schema, NameMapping 
mapping)
+      throws IOException {
+    List<String> rows = new ArrayList<>();
+    try (CloseableIterable<Record> records =
+        Parquet.read(org.apache.iceberg.Files.localInput(file))
+            .project(schema)
+            .withNameMapping(mapping)
+            .createReaderFunc(fileSchema -> 
GenericParquetReaders.buildReader(schema, fileSchema))
+            .build()) {
+      for (Record record : records) {
+        rows.add(record.toString());
+      }
+    }
+    return rows;
+  }
+
+  private static List<String> read(String file, Schema schema) throws 
IOException {
+    return read(file, schema, MappingUtil.create(schema));
+  }
+
+  // ---- files without field ids
+
+  @Test
+  public void testFileWithoutIdsTakesTheIdsTheMappingGivesItsNames() throws 
IOException {
+    String file =
+        files.write("plain.parquet", true, Arrays.asList(ID, NAME, EMAIL), 
row(1, "a", "[email protected]"));
+
+    ParquetFieldIds.Resolved resolved = 
ParquetFieldIds.resolve(ParquetFooters.read(file), ID_NAME);
+
+    assertEquals(Integer.valueOf(1), idOf(resolved, "id"));
+    assertEquals(Integer.valueOf(2), idOf(resolved, "name"));
+    assertNull(idOf(resolved, "email"));
+  }
+
+  /**
+   * The table's stored mapping keeps a renamed column's old name, and readers 
use it; a mapping
+   * created from the current schema alone would leave the old column without 
an id.
+   */
+  @Test
+  public void testFileWithoutIdsWrittenBeforeARenameTakesTheRenamedColumnsId() 
throws IOException {
+    Table table = table(ID_NAME);
+    table
+        .updateProperties()
+        .set(
+            TableProperties.DEFAULT_NAME_MAPPING,
+            NameMappingParser.toJson(MappingUtil.create(table.schema())))
+        .commit();
+    table.updateSchema().renameColumn("name", "full_name").commit();
+    String file = files.write("older.parquet", true, Arrays.asList(ID, NAME), 
row(1, "a"));
+    ParquetMetadata footer = ParquetFooters.read(file);
+
+    NameMapping forReaders =
+        NameMappingUtils.forReaders(
+            table.schema(),
+            NameMappingUtils.parseOrNull(
+                table.properties().get(TableProperties.DEFAULT_NAME_MAPPING)));
+    assertEquals(
+        Integer.valueOf(2), idOf(ParquetFieldIds.resolve(footer, table, 
forReaders), "name"));
+    assertNull(
+        idOf(ParquetFieldIds.resolve(footer, table, 
MappingUtil.create(table.schema())), "name"));
+    assertEquals(Arrays.asList("Record(1, a)"), read(file, table.schema(), 
forReaders));
+  }
+
+  // ---- files with field ids
+
+  /**
+   * Ids that agree with the table are what Iceberg's own writers produce for 
it; an id the table
+   * has never used (email's 50) is one readers ignore.
+   */
+  @Test
+  public void testFileWithTheTablesFieldIdsKeepsThem() throws IOException {
+    Table table = table(ID_NAME);
+    String file =
+        files.write(
+            "own.parquet",
+            true,
+            Arrays.asList(ID.withId(1), NAME.withId(2), EMAIL.withId(50)),
+            row(1, "a", "[email protected]"));
+    ParquetMetadata footer = ParquetFooters.read(file);
+
+    assertSame(
+        footer,
+        ParquetFieldIds.resolve(footer, table, 
MappingUtil.create(table.schema())).footer());
+    assertEquals(Arrays.asList("Record(1, a)"), read(file, table.schema()));
+  }
+
+  /** Readers use a file's own ids, so no name mapping can correct swapped 
ones. */
+  @Test
+  public void testSwappedFieldIdsAreRefused() throws IOException {
+    Table table = table(ID_NAME);
+    String file =
+        files.write(
+            "swapped.parquet", true, Arrays.asList(NAME.withId(1), 
ID.withId(2)), row("a", 1));
+
+    assertEquals(
+        "column name carries field id 1, which the table uses for column id", 
refusal(file, table));
+  }
+
+  /** Readers find no column with the table's id for name in the file, so they 
read it as null. */
+  @Test
+  public void testTableColumnUnderAnotherFieldIdIsRefused() throws IOException 
{
+    Table table = table(ID_NAME);
+    String file =
+        files.write(
+            "renumbered.parquet", true, Arrays.asList(ID.withId(1), 
NAME.withId(9)), row(1, "a"));
+
+    assertEquals(
+        "column name carries field id 9, but the table's column name has field 
id 2",
+        refusal(file, table));
+    assertEquals(Arrays.asList("Record(1, null)"), read(file, table.schema()));
+  }
+
+  /** A rename keeps the column's id, and an earlier schema version still has 
the old name. */
+  @Test
+  public void testFileWrittenBeforeARenameKeepsItsIds() throws IOException {
+    Table table = table(ID_NAME);
+    table.updateSchema().renameColumn("name", "full_name").commit();
+    String file =
+        files.write(
+            "older.parquet", true, Arrays.asList(ID.withId(1), 
NAME.withId(2)), row(1, "a"));
+
+    resolve(file, table);
+
+    assertEquals(Arrays.asList("Record(1, a)"), read(file, table.schema()));
+  }
+
+  /** Iceberg looks metrics modes up by the file's name, which a rename does 
not change. */
+  @Test
+  public void testColumnNameIsTheFilesNameAfterARename() throws IOException {
+    Table table = table(ID_NAME);
+    table.updateSchema().renameColumn("name", "full_name").commit();
+    String file =
+        files.write(
+            "older.parquet", true, Arrays.asList(ID.withId(1), 
NAME.withId(2)), row(1, "a"));
+
+    ParquetFieldIds.Resolved resolved = resolve(file, table);
+
+    assertEquals(Arrays.asList("name"), resolved.columnNames(Arrays.asList(2, 
3)));
+  }
+
+  /** A name the mapping gives an id counts as that column's, like a name from 
a schema version. */
+  @Test
+  public void testNameTheMappingGivesTheIdKeepsIt() throws IOException {
+    Table table = table(ID_NAME);
+    NameMapping aliased =
+        NameMapping.of(
+            MappedField.of(1, "id"), MappedField.of(2, Arrays.asList("name", 
"customer")));
+    PrimitiveType customer =
+        column("customer", PrimitiveTypeName.BINARY, 
LogicalTypeAnnotation.stringType());
+    String file =
+        files.write(
+            "alias.parquet", true, Arrays.asList(ID.withId(1), 
customer.withId(2)), row(1, "a"));
+    ParquetMetadata footer = ParquetFooters.read(file);
+
+    ParquetFieldIds.resolve(footer, table, aliased);
+
+    assertEquals(
+        "column customer carries field id 2, which the table uses for column 
name",
+        refusal(file, table));
+  }
+
+  /**
+   * The file's shipping.city carries the id of the table's billing.city. 
Readers find neither
+   * column in it, yet trusting the id would store its values as 
billing.city's statistics.
+   */
+  @Test
+  public void testNestedFieldCarryingTheIdOfAnotherStructsFieldIsRefused() 
throws IOException {
+    Table table =
+        table(
+            new Schema(
+                Types.NestedField.optional(
+                    1,
+                    "billing",
+                    Types.StructType.of(
+                        Types.NestedField.optional(2, "city", 
Types.StringType.get()))),
+                Types.NestedField.optional(
+                    3,
+                    "shipping",
+                    Types.StructType.of(
+                        Types.NestedField.optional(4, "city", 
Types.StringType.get())))));
+    Schema schema = table.schema();
+    int billingCity = schema.findField("billing.city").fieldId();
+    Types.StructType shipping =
+        Types.StructType.of(
+            Types.NestedField.optional(billingCity, "city", 
Types.StringType.get()));
+    Schema fileSchema =
+        new Schema(
+            Types.NestedField.optional(
+                schema.findField("shipping").fieldId(), "shipping", shipping));
+    Record city = GenericRecord.create(shipping);
+    city.setField("city", "Paris");
+    Record record = GenericRecord.create(fileSchema);
+    record.setField("shipping", city);
+    String file = files.writeWithIds("moved.parquet", fileSchema, record);
+
+    assertEquals(
+        "column shipping.city carries field id "
+            + billingCity
+            + ", which the table uses for column billing.city",
+        refusal(file, table));
+    assertEquals(Arrays.asList("Record(null, null)"), read(file, schema));
+  }
+
+  /**
+   * Struct fields, list elements and map entries are compared by full path, 
as Iceberg names them.
+   */
+  @Test
+  public void testNestedFileWithTheTablesFieldIdsKeepsThem() throws 
IOException {
+    Table table =
+        table(
+            new Schema(
+                Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+                Types.NestedField.optional(
+                    2,
+                    "address",
+                    Types.StructType.of(
+                        Types.NestedField.optional(3, "city", 
Types.StringType.get()))),
+                Types.NestedField.optional(
+                    4, "tags", Types.ListType.ofOptional(5, 
Types.StringType.get())),
+                Types.NestedField.optional(
+                    6,
+                    "attrs",
+                    Types.MapType.ofOptional(
+                        7, 8, Types.StringType.get(), 
Types.StringType.get()))));
+    String file =
+        files.writeWithIds("nested.parquet", table.schema(), 
nestedRecord(table.schema()));
+
+    resolve(file, table);
+
+    assertEquals(Arrays.asList("Record(1, Record(Paris), [a], {k=v})"), 
read(file, table.schema()));
+  }
+
+  /**
+   * A file written before its struct was renamed matches the schema version 
it was written with.
+   */
+  @Test
+  public void testFileWrittenBeforeAStructRenameKeepsItsIds() throws 
IOException {
+    Table table =
+        table(
+            new Schema(
+                Types.NestedField.optional(1, "id", Types.IntegerType.get()),
+                Types.NestedField.optional(
+                    2,
+                    "address",
+                    Types.StructType.of(
+                        Types.NestedField.optional(3, "city", 
Types.StringType.get())))));
+    Schema before = table.schema();
+    table.updateSchema().renameColumn("address", "location").commit();
+    String file = files.writeWithIds("older.parquet", before, 
nestedRecord(before));
+
+    resolve(file, table);
+
+    assertEquals(Arrays.asList("Record(1, Record(Paris))"), read(file, 
table.schema()));
+  }
+
+  @Test
+  public void testDuplicateFieldIdsAreRefused() throws IOException {
+    Table table = table(ID_NAME);
+    String file =
+        files.write(
+            "duplicate.parquet", true, Arrays.asList(ID.withId(1), 
NAME.withId(1)), row(1, "a"));
+
+    String message = refusal(file, table);
+
+    assertTrue(message, message.startsWith("its field ids cannot be 
resolved"));
+  }
+
+  /** Readers resolve by ids once any column has one, so a column without an 
id reads as null. */
+  @Test
+  public void testFileWithIdsOnSomeColumnsIsRefused() throws IOException {
+    Table table = table(ID_NAME);
+    String file =
+        files.write("partial.parquet", true, Arrays.asList(ID.withId(1), 
NAME), row(1, "a"));
+
+    String message = refusal(file, table);
+
+    assertTrue(message, message.startsWith("column name "));
+    assertEquals(Arrays.asList("Record(1, null)"), read(file, table.schema()));
+  }
+
+  /** Fills whichever of id, address.city, tags and attrs {@code schema} has. 
*/
+  private static Record nestedRecord(Schema schema) {
+    Record record = GenericRecord.create(schema);
+    record.setField("id", 1);
+    Record address = 
GenericRecord.create(schema.findType("address").asStructType());
+    address.setField("city", "Paris");
+    record.setField("address", address);
+    if (schema.findField("tags") != null) {
+      record.setField("tags", Arrays.asList("a"));
+      record.setField("attrs", ImmutableMap.of("k", "v"));
+    }
+    return record;
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetTestFiles.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetTestFiles.java
new file mode 100644
index 00000000000..33b0e542c12
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/ParquetTestFiles.java
@@ -0,0 +1,170 @@
+/*
+ * 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.beam.sdk.io.iceberg;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import org.apache.hadoop.fs.Path;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.data.parquet.GenericParquetWriter;
+import org.apache.iceberg.io.DataWriter;
+import org.apache.iceberg.parquet.Parquet;
+import org.apache.parquet.example.data.Group;
+import org.apache.parquet.example.data.simple.NanoTime;
+import org.apache.parquet.example.data.simple.SimpleGroupFactory;
+import org.apache.parquet.hadoop.ParquetWriter;
+import org.apache.parquet.hadoop.example.ExampleParquetWriter;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
+import org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit;
+import org.apache.parquet.schema.MessageType;
+import org.apache.parquet.schema.PrimitiveType;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.apache.parquet.schema.Type.Repetition;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * Writes Parquet files for the AddFiles metrics tests, with exactly the 
columns, field ids, row
+ * groups and statistics a test asks for.
+ */
+final class ParquetTestFiles {
+  /** 2024-01-01T00:00:00Z. */
+  static final long EPOCH_SECONDS = 1704067200L;
+
+  static final Map<String, String> FULL_METRICS =
+      ImmutableMap.of("write.metadata.metrics.default", "full");
+
+  static final PrimitiveType ID = column("id", PrimitiveTypeName.INT32, null);
+  static final PrimitiveType FLAG = column("flag", PrimitiveTypeName.BOOLEAN, 
null);
+  static final PrimitiveType NAME =
+      column("name", PrimitiveTypeName.BINARY, 
LogicalTypeAnnotation.stringType());
+  static final PrimitiveType TS_MICROS =
+      column(
+          "ts",
+          PrimitiveTypeName.INT64,
+          LogicalTypeAnnotation.timestampType(true, TimeUnit.MICROS));
+  static final PrimitiveType TS_MILLIS =
+      column(
+          "ts",
+          PrimitiveTypeName.INT64,
+          LogicalTypeAnnotation.timestampType(true, TimeUnit.MILLIS));
+  static final PrimitiveType TS_NANOS =
+      column(
+          "ts", PrimitiveTypeName.INT64, 
LogicalTypeAnnotation.timestampType(true, TimeUnit.NANOS));
+  static final PrimitiveType TS_INT96 = column("ts", PrimitiveTypeName.INT96, 
null);
+  static final PrimitiveType UNSIGNED =
+      column("u", PrimitiveTypeName.INT32, LogicalTypeAnnotation.intType(32, 
false));
+
+  private final File dir;
+
+  ParquetTestFiles(File dir) {
+    this.dir = dir;
+  }
+
+  static PrimitiveType column(
+      String name, PrimitiveTypeName physical, @Nullable LogicalTypeAnnotation 
annotation) {
+    org.apache.parquet.schema.Types.PrimitiveBuilder<PrimitiveType> builder =
+        org.apache.parquet.schema.Types.primitive(physical, 
Repetition.OPTIONAL);
+    if (annotation != null) {
+      builder = builder.as(annotation);
+    }
+    return builder.named(name);
+  }
+
+  static Object[] row(@Nullable Object... values) {
+    return values;
+  }
+
+  /** One row per array, one value per column; a null value leaves the column 
unset. */
+  String write(String name, boolean statistics, List<PrimitiveType> columns, 
Object[]... rows)
+      throws IOException {
+    return write(name, statistics, 0, columns, rows);
+  }
+
+  /** Cuts a row group every {@code rowsPerGroup} rows; 0 writes a single row 
group. */
+  String write(
+      String name,
+      boolean statistics,
+      int rowsPerGroup,
+      List<PrimitiveType> columns,
+      Object[]... rows)
+      throws IOException {
+    MessageType type = new MessageType("root", new ArrayList<>(columns));
+    File file = new File(dir, name);
+    SimpleGroupFactory factory = new SimpleGroupFactory(type);
+    ExampleParquetWriter.Builder builder =
+        ExampleParquetWriter.builder(new Path(file.getAbsolutePath()))
+            .withType(type)
+            .withStatisticsEnabled(statistics);
+    if (rowsPerGroup > 0) {
+      builder =
+          builder
+              .withRowGroupSize(1L)
+              .withMinRowCountForPageSizeCheck(rowsPerGroup)
+              .withMaxRowCountForPageSizeCheck(rowsPerGroup);
+    }
+    try (ParquetWriter<Group> writer = builder.build()) {
+      for (Object[] values : rows) {
+        Group group = factory.newGroup();
+        for (int i = 0; i < columns.size(); i++) {
+          add(group, columns.get(i).getName(), values[i]);
+        }
+        writer.write(group);
+      }
+    }
+    return file.getAbsolutePath();
+  }
+
+  private static void add(Group group, String column, @Nullable Object value) {
+    if (value instanceof Long) {
+      group.add(column, (Long) value);
+    } else if (value instanceof Integer) {
+      group.add(column, (Integer) value);
+    } else if (value instanceof Boolean) {
+      group.add(column, (Boolean) value);
+    } else if (value instanceof String) {
+      group.add(column, (String) value);
+    } else if (value instanceof NanoTime) {
+      group.add(column, (NanoTime) value);
+    }
+  }
+
+  /** Writes with Iceberg's own writer, so every column carries {@code 
fileSchema}'s field id. */
+  String writeWithIds(String name, Schema fileSchema, Record... records) 
throws IOException {
+    String file = new File(dir, name).getAbsolutePath();
+    DataWriter<Record> writer =
+        Parquet.writeData(org.apache.iceberg.Files.localOutput(file))
+            .schema(fileSchema)
+            .withSpec(PartitionSpec.unpartitioned())
+            .createWriterFunc(GenericParquetWriter::create)
+            .build();
+    try {
+      for (Record record : records) {
+        writer.write(record);
+      }
+    } finally {
+      writer.close();
+    }
+    return file;
+  }
+}

Reply via email to