claudevdm commented on code in PR #40252:
URL: https://github.com/apache/beam/pull/40252#discussion_r4126847874
##########
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:
We now reject a file whose uint32 column holds a value of 2^31 or more.
Iceberg's readers read uint32 values as signed, so 3000000000 comes back as
-1294967296. Comparing the row-group min/max values as unsigned would give
bounds that don't match what readers return.
Below 2^31 signed and unsigned order are the same, so the bounds Iceberg
computes are correct and we keep them.
--
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]