ahmedabu98 commented on code in PR #40252:
URL: https://github.com/apache/beam/pull/40252#discussion_r4109254199


##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java:
##########
@@ -942,16 +950,16 @@ private String getPartitionFromFilePath(String filePath) {
      * <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(
+    static PartitionKey getPartitionFromMetrics(
         Metrics metrics, InputFile inputFile, Table table, @Nullable 
ParquetMetadata preReadFooter)
         throws UnknownPartitionException {
       List<PartitionField> fields = table.spec().fields();
       List<Integer> sourceIds =
           
fields.stream().map(PartitionField::sourceId).collect(Collectors.toList());
       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)) {
         partitionMetrics = metrics;
       } else {

Review Comment:
   We need to add another condition to this fast path:
   All partition source metric columns (specifically string and binary types) 
should be "full", and not truncated.
   
   Iceberg does a weird thing to make its metric bounds make sense to readers. 
When it [truncates the upper 
bound](https://github.com/apache/iceberg/blob/729e650b9e0a9f706ea8bea795b3a6dfb58e306b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetMetrics.java#L728),
 it [increments the very last 
character](https://github.com/apache/iceberg/blob/729e650b9e0a9f706ea8bea795b3a6dfb58e306b/api/src/main/java/org/apache/iceberg/util/UnicodeUtil.java#L88-L106)
 to make sure it's greater than all the values in that column. So a min/max of 
("aaabb", "aaabb") with metric mode truncate(3) gets converted to bounds 
("aaa", "aab"). So this way the upper bound covers "aaabb".
   
   So if we have a truncated string metric column, and compare it with an 
identity partition, we'll find that the lower and upper bounds don't match and 
incorrectly pass it to the error DLQ.
   Instead, we should let it pass to the `else` block where it reconstructs the 
Metric object with full mode.
   
   Maybe add a test that uses truncated table metrics, and create a file with 
long string min/max values.



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java:
##########
@@ -1218,7 +1308,18 @@ public static Metrics getFileMetrics(
         if (!ParquetSchemaUtil.hasIds(originalMessageType)) {
           footer = getFooterWithTypeIds(originalMessageType, footer, mapping);
         }

Review Comment:
   Existing issue that is identical with the last comment. We can't trust that 
`originalMessageType` has IDs that correspond to the table's, so we should 
always apply `getFooterWithTypeIds(..)`



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java:
##########
@@ -1261,6 +1362,234 @@ static ParquetMetadata getFooterWithTypeIds(
     return new ParquetMetadata(newFileMeta, footer.getBlocks());
   }
 
+  /**
+   * 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 
an unsigned 32-bit int
+   * under a long column throws when the int is cast to a long. 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.
+   */
+  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, lower rounded down, 
upper rounded up. */
+    NANOS_TO_MICROS,
+    /** Unsigned 32-bit int stored under a long. */
+    UINT32_TO_LONG;
+
+    static Map<Integer, BoundAdjustment> forSchema(
+        MessageType fileSchema, org.apache.iceberg.Schema tableSchema) {
+      Map<Integer, BoundAdjustment> adjustments = new HashMap<>();
+      collect(fileSchema, tableSchema, adjustments);
+      return adjustments;
+    }
+
+    private static void collect(
+        org.apache.parquet.schema.GroupType group,
+        org.apache.iceberg.Schema tableSchema,
+        Map<Integer, BoundAdjustment> out) {
+      for (org.apache.parquet.schema.Type field : group.getFields()) {
+        if (!field.isPrimitive()) {
+          collect(field.asGroupType(), tableSchema, out);
+          continue;
+        }
+        org.apache.parquet.schema.Type.ID id = field.getId();
+        if (id == null) {
+          continue;
+        }
+        @Nullable Type tableType = tableSchema.findType(id.intValue());
+        if (tableType == null) {
+          continue;
+        }
+        @Nullable BoundAdjustment adjustment = 
forPrimitive(field.asPrimitiveType(), tableType);
+        if (adjustment != null) {
+          out.put(id.intValue(), adjustment);
+        }
+      }
+    }
+
+    /**
+     * Null when the file and table units agree, or when the table type is not 
the matching
+     * timestamp, time or long type: such a column is left as Iceberg computes 
it.
+     */
+    private static @Nullable BoundAdjustment forPrimitive(
+        org.apache.parquet.schema.PrimitiveType primitive, Type tableType) {
+      org.apache.parquet.schema.LogicalTypeAnnotation annotation =
+          primitive.getLogicalTypeAnnotation();
+      if (annotation
+          instanceof
+          
org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation) 
{
+        org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit fileUnit =
+            
((org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation)
+                    annotation)
+                .getUnit();
+        if (tableType.typeId() == Type.TypeID.TIMESTAMP) {
+          return toMicros(fileUnit);
+        }
+        if (tableType.typeId() == Type.TypeID.TIMESTAMP_NANO) {
+          return toNanos(fileUnit);
+        }
+        return null;
+      }
+      if (annotation
+          instanceof 
org.apache.parquet.schema.LogicalTypeAnnotation.TimeLogicalTypeAnnotation) {
+        if (tableType.typeId() != Type.TypeID.TIME) {
+          return null;
+        }
+        return toMicros(
+            
((org.apache.parquet.schema.LogicalTypeAnnotation.TimeLogicalTypeAnnotation) 
annotation)
+                .getUnit());
+      }
+      if (annotation
+          instanceof 
org.apache.parquet.schema.LogicalTypeAnnotation.IntLogicalTypeAnnotation) {
+        
org.apache.parquet.schema.LogicalTypeAnnotation.IntLogicalTypeAnnotation 
intType =
+            
(org.apache.parquet.schema.LogicalTypeAnnotation.IntLogicalTypeAnnotation) 
annotation;
+        if (intType.getBitWidth() == 32
+            && !intType.isSigned()
+            && tableType.typeId() == Type.TypeID.LONG) {
+          return UINT32_TO_LONG;
+        }
+      }
+      return null;
+    }
+
+    private static @Nullable BoundAdjustment toMicros(
+        org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit fileUnit) {
+      switch (fileUnit) {
+        case MILLIS:
+          return MILLIS_TO_MICROS;
+        case NANOS:
+          return NANOS_TO_MICROS;
+        default:
+          return null;
+      }
+    }
+
+    private static @Nullable BoundAdjustment toNanos(
+        org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit fileUnit) {
+      switch (fileUnit) {
+        case MILLIS:
+          return MILLIS_TO_NANOS;
+        case MICROS:
+          return MICROS_TO_NANOS;
+        default:
+          return null;
+      }
+    }
+
+    /** The footer with annotations removed from adjusted INT32 columns. */
+    static ParquetMetadata withNeutralTypes(
+        ParquetMetadata footer, Map<Integer, BoundAdjustment> adjustments) {
+      MessageType schema = footer.getFileMetaData().getSchema();
+      List<org.apache.parquet.schema.Type> fields = neutralFields(schema, 
adjustments);
+      MessageType neutral = new MessageType(schema.getName(), fields);
+      FileMetaData meta = footer.getFileMetaData();
+      return new ParquetMetadata(
+          new FileMetaData(neutral, meta.getKeyValueMetaData(), 
meta.getCreatedBy()),
+          footer.getBlocks());
+    }
+
+    private static List<org.apache.parquet.schema.Type> neutralFields(
+        org.apache.parquet.schema.GroupType group, Map<Integer, 
BoundAdjustment> adjustments) {
+      List<org.apache.parquet.schema.Type> fields = new ArrayList<>();
+      for (org.apache.parquet.schema.Type field : group.getFields()) {
+        if (field.isPrimitive()) {
+          fields.add(neutralPrimitive(field.asPrimitiveType(), adjustments));
+        } else {
+          fields.add(
+              
field.asGroupType().withNewFields(neutralFields(field.asGroupType(), 
adjustments)));
+        }
+      }
+      return fields;
+    }
+
+    private static org.apache.parquet.schema.Type neutralPrimitive(
+        org.apache.parquet.schema.PrimitiveType primitive,
+        Map<Integer, BoundAdjustment> adjustments) {
+      boolean adjusted =
+          primitive.getId() != null && 
adjustments.containsKey(primitive.getId().intValue());
+      if (!adjusted
+          || primitive.getPrimitiveTypeName()
+              != 
org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT32) {
+        return primitive;
+      }
+      return org.apache.parquet.schema.Types.primitive(
+              primitive.getPrimitiveTypeName(), primitive.getRepetition())
+          .id(primitive.getId().intValue())
+          .named(primitive.getName());

Review Comment:
   AI review complained that this creates an issue when (1) the file has 
multiple row groups, and (2) one row group's bound is a small unsigned int and 
the other row group's bound is a very large unsigned int (2^31 and above).
   
   This block recreates the type without preserving the unsigned info, so 
Iceberg treats it as a signed int. When Iceberg tries merging row groups, it 
compares their bounds to create one pair of bounds for the file. In this case, 
the comparison can interpret large unsigned ints (e.g. `2147483649`) as 
negative signed ints (`-2147483647`) and output incorrect bounds.
   
   We can trust the min/max bounds that Parquet shows for each row group, but 
we can't trust Iceberg's min/max merge output across row groups in this case. 
We should still strip the annotation but we'd need to compute the bounds of 
unsigned int columns ourselves (using `Integer.compareUnsigned`).
   
   Maybe also include a test with two row groups to confirm the failure and 
also make sure the fix works as intended



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFiles.java:
##########
@@ -1004,11 +1027,77 @@ static String getPartitionFromMetrics(
           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;
+    }
+
+    /** 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(@Nullable ParquetMetadata footer, Table 
table, int fieldId) {
+      if (footer == null) {
+        return false;
+      }
+      MessageType fileType = footer.getFileMetaData().getSchema();
+      if (!ParquetSchemaUtil.hasIds(fileType)) {
+        fileType = ParquetSchemaUtil.applyNameMapping(fileType, 
MappingUtil.create(table.schema()));

Review Comment:
   If a file has IDs, this would assume that those IDs correlate exactly with 
the Table's IDs? I don't think we can make that assumption, since we don't know 
how this file was created. We should always apply the name mapping.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to