claudevdm commented on code in PR #40252:
URL: https://github.com/apache/beam/pull/40252#discussion_r4126870402
##########
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:
Agreed. Always applying the name mapping isn't enough: it only replaces ids
for names it knows, and readers use a file's own ids whenever it has any,
ignoring the 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]