From dac69a5aeff9c2a1b349ca908b9398493a197019 Mon Sep 17 00:00:00 2001 From: Vova Kolmakov Date: Mon, 3 Aug 2026 13:07:32 +0700 Subject: [PATCH 1/2] fix(trino): skip predicate pushdown on type-evolved parquet columns --- .../plugin/hudi/HudiPageSourceProvider.java | 6 +- .../hudi/util/ParquetStatisticsDomains.java | 172 ++++++++ .../hudi/TestHudiEvolvedColumnPredicates.java | 388 ++++++++++++++++++ .../TestHudiSchemaEvolutionPredicates.java | 98 +++++ ...diSchemaEvolutionPredicatesPositional.java | 38 ++ .../SchemaEvolutionHudiTablesInitializer.java | 164 ++++++++ .../util/TestParquetStatisticsDomains.java | 297 ++++++++++++++ 7 files changed, 1162 insertions(+), 1 deletion(-) create mode 100644 hudi-trino/src/main/java/io/trino/plugin/hudi/util/ParquetStatisticsDomains.java create mode 100644 hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiEvolvedColumnPredicates.java create mode 100644 hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicates.java create mode 100644 hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicatesPositional.java create mode 100644 hudi-trino/src/test/java/io/trino/plugin/hudi/testing/SchemaEvolutionHudiTablesInitializer.java create mode 100644 hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java diff --git a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPageSourceProvider.java b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPageSourceProvider.java index 2803f8ae231a..674c91b61a8f 100644 --- a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPageSourceProvider.java +++ b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPageSourceProvider.java @@ -113,6 +113,7 @@ import static io.trino.plugin.hudi.HudiUtil.prependHudiMetaAndMergeRequiredColumns; import static io.trino.plugin.hudi.HudiUtil.resolveMergeModeAndStrategyId; import static io.trino.plugin.hudi.HudiUtil.usesNonProjectionCompatibleMerger; +import static io.trino.plugin.hudi.util.ParquetStatisticsDomains.dropIncomparableDomains; import static io.trino.spi.StandardErrorCode.NOT_SUPPORTED; import static java.lang.String.format; import static java.util.Objects.requireNonNull; @@ -406,9 +407,12 @@ static ConnectorPageSource createPageSource( Map, ColumnDescriptor> descriptorsByPath = getDescriptors(fileSchema, requestedSchema); + // A domain typed by the metastore cannot be matched against the statistics of a column the file stores + // under the type it had before a schema evolution, so those are dropped before the parquet predicate + // ever sees them. See ParquetStatisticsDomains. TupleDomain parquetTupleDomain = options.isIgnoreStatistics() || !enablePredicatePushDown ? TupleDomain.all() - : getParquetTupleDomain(descriptorsByPath, getPushdownPredicate(hudiSplit, dynamicFilter, physicalIndexMap), fileSchema, useColumnNames); + : dropIncomparableDomains(getParquetTupleDomain(descriptorsByPath, getPushdownPredicate(hudiSplit, dynamicFilter, physicalIndexMap), fileSchema, useColumnNames)); TupleDomainParquetPredicate parquetPredicate = buildPredicate(requestedSchema, parquetTupleDomain, descriptorsByPath, timeZone); diff --git a/hudi-trino/src/main/java/io/trino/plugin/hudi/util/ParquetStatisticsDomains.java b/hudi-trino/src/main/java/io/trino/plugin/hudi/util/ParquetStatisticsDomains.java new file mode 100644 index 000000000000..0a8cdf694c9e --- /dev/null +++ b/hudi-trino/src/main/java/io/trino/plugin/hudi/util/ParquetStatisticsDomains.java @@ -0,0 +1,172 @@ +/* + * Licensed 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 io.trino.plugin.hudi.util; + +import io.airlift.log.Logger; +import io.trino.spi.predicate.Domain; +import io.trino.spi.predicate.TupleDomain; +import io.trino.spi.type.DecimalType; +import io.trino.spi.type.TimestampType; +import io.trino.spi.type.Type; +import io.trino.spi.type.VarcharType; +import org.apache.parquet.column.ColumnDescriptor; +import org.apache.parquet.schema.LogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName; + +import java.util.LinkedHashMap; +import java.util.Map; + +import static io.trino.spi.type.BigintType.BIGINT; +import static io.trino.spi.type.BooleanType.BOOLEAN; +import static io.trino.spi.type.DateType.DATE; +import static io.trino.spi.type.DoubleType.DOUBLE; +import static io.trino.spi.type.IntegerType.INTEGER; +import static io.trino.spi.type.RealType.REAL; +import static io.trino.spi.type.SmallintType.SMALLINT; +import static io.trino.spi.type.TinyintType.TINYINT; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.BINARY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.FLOAT; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT32; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT64; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT96; + +/** + * Keeps a pushed-down predicate from being matched against statistics it cannot be compared with. + *

+ * {@code TupleDomainParquetPredicate.getDomain} selects its branch on the type of the pushed-down DOMAIN and then + * reads the parquet statistics as that type - {@code Double min = (Double) minimums.get(i)} and so on. The domain's + * type comes from the metastore while the statistics come from the file, and Hudi's type evolution is exactly what + * makes those two disagree: after a column evolves and the metastore is synced, every base file written before the + * evolution still stores the old physical type. Handing such a domain to the parquet predicate either fails the + * whole split with {@code Malformed Parquet file. Corrupted statistics for column ...} wrapping a + * {@link ClassCastException}, or - where the two types happen to share a representation, as a decimal and a varchar + * both do through {@code Slice} - silently prunes row groups on a comparison that means nothing. + *

+ * The read path has no such problem: {@code ColumnReaderFactory} decodes parquet {@code FLOAT} into Trino + * {@code DOUBLE} and {@code INT32} into {@code BIGINT} natively, and {@code ParquetTypeTranslator.createCoercer} + * covers the rest of the promotions the page source is asked for. Only the statistics side is blind, so only the + * statistics side needs the guard. + */ +public final class ParquetStatisticsDomains +{ + private static final Logger log = Logger.get(ParquetStatisticsDomains.class); + + private ParquetStatisticsDomains() {} + + /** + * Drops every domain whose type cannot be compared against its column's statistics, keeping the rest untouched. + *

+ * Dropping loses row group pruning for that column but never a row: {@code HudiMetadata.applyFilter} hands the + * whole regular predicate back to the engine as the remaining filter and does not precalculate statistics for + * the pushdown, so a connector-side domain is an optimization and nothing else. The dynamic half of the + * predicate is redundant with the join above the scan by construction. It is the same trade + * {@code HudiPageSourceProvider.remapPredicateColumnIndicesToPhysical} already makes for a predicate column the + * file does not carry, and the same one {@code HudiColumnStatsIndexSupport.getDomainFromColumnStats} makes when + * the metadata table's column statistics do not match the column's type. + *

+ * The filter runs on the descriptor-keyed domain rather than on the column handles it was built from. That is + * what the parquet predicate itself will be evaluated against, so the check and the evaluation cannot disagree + * about which column or which physical type is meant; it covers both values of + * {@code hudi.parquet.use-column-names} in one pass, since a handle is resolved to a descriptor before either; + * and a dereference handle contributes the leaf field's type without any extra work. + */ + public static TupleDomain dropIncomparableDomains(TupleDomain parquetTupleDomain) + { + if (parquetTupleDomain.isAll() || parquetTupleDomain.isNone()) { + return parquetTupleDomain; + } + + Map domains = parquetTupleDomain.getDomains().orElseThrow(); + Map comparableDomains = new LinkedHashMap<>(); + for (Map.Entry entry : domains.entrySet()) { + if (hasComparableStatistics(entry.getValue().getType(), entry.getKey().getPrimitiveType())) { + comparableDomains.put(entry.getKey(), entry.getValue()); + } + else { + log.debug("Not pushing down a %s predicate on %s: the file stores it as %s, so the column statistics cannot answer it", + entry.getValue().getType(), entry.getKey(), entry.getKey().getPrimitiveType()); + } + } + if (comparableDomains.size() == domains.size()) { + return parquetTupleDomain; + } + return TupleDomain.withColumnDomains(comparableDomains); + } + + /** + * Whether {@code TupleDomainParquetPredicate.getDomain} builds a meaningful domain out of a {@code fileType} + * column's statistics when asked for {@code domainType}, which is the case only when the two describe the same + * physical values. + *

+ * The accepted pairs mirror that method's dispatch branch for branch. Everything else is rejected, which for a + * type it has no branch for - {@code CHAR}, {@code VARBINARY}, {@code UUID}, {@code TIME}, a timestamp with time + * zone - costs nothing at all: its fallthrough returns a domain covering every value, which prunes exactly as + * much as pushing nothing down. Only the accepted pairs can be wrong, and + * {@code TestParquetStatisticsDomains} pins each of them against the real {@code getDomain}. + */ + public static boolean hasComparableStatistics(Type domainType, PrimitiveType fileType) + { + PrimitiveTypeName primitiveType = fileType.getPrimitiveTypeName(); + LogicalTypeAnnotation annotation = fileType.getLogicalTypeAnnotation(); + + if (BOOLEAN.equals(domainType)) { + return primitiveType == PrimitiveTypeName.BOOLEAN; + } + if (TINYINT.equals(domainType) || SMALLINT.equals(domainType) || INTEGER.equals(domainType) + || BIGINT.equals(domainType) || DATE.equals(domainType)) { + // asLong takes any integral box the statistics can hold, so INT32 and INT64 are interchangeable here. + // A decimal column reports its UNSCALED value though, which is the integer it stands for only at scale 0. + return (primitiveType == INT32 || primitiveType == INT64) && !isScaledDecimal(annotation); + } + if (domainType instanceof DecimalType) { + // getShortDecimal and getLongDecimal rescale against the column's own annotation. Without one they read + // the raw value as an unscaled decimal at the DOMAIN's scale, so an int or a string that evolved into a + // decimal would be compared a factor of ten-to-the-scale off, or as raw UTF-8 bytes. + return isDecimalPrimitive(primitiveType) && annotation instanceof DecimalLogicalTypeAnnotation; + } + if (REAL.equals(domainType)) { + return primitiveType == FLOAT; + } + if (DOUBLE.equals(domainType)) { + return primitiveType == PrimitiveTypeName.DOUBLE; + } + if (domainType instanceof VarcharType) { + // Both sides compare raw bytes, which is varchar's own ordering - unless the bytes are a decimal's + // big-endian two's complement, which orders nothing like the digits it prints as. + return (primitiveType == BINARY || primitiveType == FIXED_LEN_BYTE_ARRAY) + && !(annotation instanceof DecimalLogicalTypeAnnotation); + } + if (domainType instanceof TimestampType) { + // INT96 is read from the binary statistics and INT64 from the long ones, but an INT64 column has to say + // which unit it counts in before its bounds mean anything. + return primitiveType == INT96 + || (primitiveType == INT64 && annotation instanceof TimestampLogicalTypeAnnotation timestampAnnotation && timestampAnnotation.getUnit() != null); + } + return false; + } + + private static boolean isDecimalPrimitive(PrimitiveTypeName primitiveType) + { + return primitiveType == INT32 || primitiveType == INT64 || primitiveType == BINARY || primitiveType == FIXED_LEN_BYTE_ARRAY; + } + + private static boolean isScaledDecimal(LogicalTypeAnnotation annotation) + { + return annotation instanceof DecimalLogicalTypeAnnotation decimalAnnotation && decimalAnnotation.getScale() != 0; + } +} diff --git a/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiEvolvedColumnPredicates.java b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiEvolvedColumnPredicates.java new file mode 100644 index 000000000000..ddbc6fa95490 --- /dev/null +++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiEvolvedColumnPredicates.java @@ -0,0 +1,388 @@ +/* + * Licensed 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 io.trino.plugin.hudi; + +import io.airlift.slice.Slices; +import io.trino.filesystem.local.LocalInputFile; +import io.trino.metastore.HiveType; +import io.trino.parquet.ParquetReaderOptions; +import io.trino.plugin.base.metrics.FileFormatDataSourceStats; +import io.trino.plugin.hive.HiveColumnHandle; +import io.trino.plugin.hive.parquet.ParquetReaderConfig; +import io.trino.plugin.hudi.file.HudiBaseFile; +import io.trino.spi.SplitWeight; +import io.trino.spi.connector.ColumnHandle; +import io.trino.spi.connector.ConnectorPageSource; +import io.trino.spi.connector.ConnectorSession; +import io.trino.spi.connector.DynamicFilter; +import io.trino.spi.predicate.Domain; +import io.trino.spi.predicate.Range; +import io.trino.spi.predicate.TupleDomain; +import io.trino.spi.predicate.ValueSet; +import io.trino.spi.type.Type; +import io.trino.testing.MaterializedResult; +import io.trino.testing.TestingConnectorSession; +import org.apache.parquet.conf.PlainParquetConfiguration; +import org.apache.parquet.example.data.Group; +import org.apache.parquet.example.data.simple.SimpleGroupFactory; +import org.apache.parquet.hadoop.ParquetFileReader; +import org.apache.parquet.hadoop.ParquetWriter; +import org.apache.parquet.hadoop.example.ExampleParquetWriter; +import org.apache.parquet.io.LocalOutputFile; +import org.apache.parquet.schema.LogicalTypeAnnotation; +import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.Types; +import org.joda.time.DateTimeZone; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.OptionalLong; +import java.util.Set; +import java.util.concurrent.CompletableFuture; + +import static io.trino.plugin.hive.HiveColumnHandle.createBaseColumn; +import static io.trino.plugin.hudi.HudiPageSourceProvider.createPageSource; +import static io.trino.spi.type.BigintType.BIGINT; +import static io.trino.spi.type.DoubleType.DOUBLE; +import static io.trino.spi.type.IntegerType.INTEGER; +import static io.trino.spi.type.VarcharType.VARCHAR; +import static io.trino.testing.MaterializedResult.materializeSourceDataStream; +import static org.apache.hudi.common.model.HoodieRecord.HOODIE_META_COLUMNS; +import static org.apache.parquet.schema.Type.Repetition.OPTIONAL; +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Covers what a pushed-down predicate does to a base file written before the column it constrains was evolved. + *

+ * Hudi lets a column's type widen and hive sync then reports the NEW type, while every base file written before the + * evolution keeps storing the old one. The parquet reader copes with that on its own -- {@code ColumnReaderFactory} + * decodes {@code FLOAT} into {@code DOUBLE} and {@code INT32} into {@code BIGINT}, and + * {@code ParquetTypeTranslator.createCoercer} handles the rest -- but the statistics do not: {@code + * TupleDomainParquetPredicate.getDomain} reads them as whatever the DOMAIN's type says, so a {@code DOUBLE} domain + * over a {@code FLOAT} column casts a {@code Float} to a {@code Double} and fails the whole split with {@code + * HUDI_BAD_DATA}. See apache/hudi#19457. + *

+ * The fixture writes every data column as the type it had BEFORE the evolution and every handle carries the type the + * metastore reports AFTER it, which is exactly the state an unrewritten base file is in. Values grow with the row + * index so pruning stays observable: a predicate that survives pushdown reads fewer rows than the file holds, and one + * that was dropped reads all of them. Do not "simplify" that to asserting the matching rows alone -- both a working + * pushdown and no pushdown at all produce the same matching rows, since the connector's pushdown is an optimization + * and the engine re-applies the predicate above the scan. + */ +class TestHudiEvolvedColumnPredicates +{ + private static final String STABLE_COLUMN = "stable_int"; + /** Written as parquet FLOAT, reported by the metastore as double. */ + private static final String FLOAT_TO_DOUBLE_COLUMN = "evolved_double"; + /** Written as parquet INT32, reported by the metastore as bigint. */ + private static final String INT_TO_BIGINT_COLUMN = "evolved_bigint"; + /** Written as parquet INT32, reported by the metastore as string. */ + private static final String INT_TO_VARCHAR_COLUMN = "evolved_varchar"; + + private static final int ROW_COUNT = 1000; + private static final long THRESHOLD = 900; + private static final int MATCHING_ROW_COUNT = (int) (ROW_COUNT - THRESHOLD - 1); + + @TempDir + static Path tempDir; + + private static Path baseFile; + + @BeforeAll + static void writeBaseFile() + throws IOException + { + MessageType schema = preEvolutionFileSchema(); + baseFile = tempDir.resolve("evolved_base_file.parquet"); + SimpleGroupFactory groupFactory = new SimpleGroupFactory(schema); + try (ParquetWriter writer = ExampleParquetWriter.builder(new LocalOutputFile(baseFile)) + .withType(schema) + .withConf(new PlainParquetConfiguration()) + .withRowGroupSize(1024L) + .withPageSize(512) + .build()) { + for (int row = 0; row < ROW_COUNT; row++) { + Group group = groupFactory.newGroup(); + for (String metaColumn : HOODIE_META_COLUMNS) { + group.append(metaColumn, metaColumn + "_" + row); + } + group.append(STABLE_COLUMN, row); + group.append(FLOAT_TO_DOUBLE_COLUMN, (float) row); + group.append(INT_TO_BIGINT_COLUMN, row); + group.append(INT_TO_VARCHAR_COLUMN, row); + writer.write(group); + } + } + // With a single row group there would be nothing to prune and every "still prunes" assertion below would + // hold without proving anything, so assert the outcome rather than the writer knobs that produce it. + assertThat(rowGroupCount(baseFile)).as("row groups written").isGreaterThan(1); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void testPredicateOnFloatColumnEvolvedToDouble(boolean useParquetColumnNames) + throws Exception + { + HiveColumnHandle evolved = column(FLOAT_TO_DOUBLE_COLUMN, HiveType.HIVE_DOUBLE, DOUBLE); + List projection = List.of(column(STABLE_COLUMN, HiveType.HIVE_INT, INTEGER), evolved); + + MaterializedResult result = read(projection, greaterThanThreshold(evolved, DOUBLE, (double) THRESHOLD), + useParquetColumnNames, DynamicFilter.EMPTY); + + // The domain cannot be matched against FLOAT statistics, so it is dropped and nothing is pruned. Before the + // fix this threw HUDI_BAD_DATA ("Corrupted statistics for column") instead of reading anything at all. + assertThat(result.getRowCount()).as("rows read").isEqualTo(ROW_COUNT); + // The read itself promotes, so the rows the engine will filter carry the widened values + assertThat(result.getMaterializedRows().get(7).getField(1)).as("promoted value of row 7").isEqualTo(7.0d); + assertThat(valuesOver(result, 1)).as("rows matching %s > %s", FLOAT_TO_DOUBLE_COLUMN, THRESHOLD).isEqualTo(MATCHING_ROW_COUNT); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void testPredicateOnIntColumnEvolvedToVarchar(boolean useParquetColumnNames) + throws Exception + { + HiveColumnHandle evolved = column(INT_TO_VARCHAR_COLUMN, HiveType.HIVE_STRING, VARCHAR); + List projection = List.of(column(STABLE_COLUMN, HiveType.HIVE_INT, INTEGER), evolved); + + MaterializedResult result = read(projection, + greaterThanThreshold(evolved, VARCHAR, Slices.utf8Slice("900")), + useParquetColumnNames, DynamicFilter.EMPTY); + + // A varchar domain over an INT32 column would cast an Integer to a Slice + assertThat(result.getRowCount()).as("rows read").isEqualTo(ROW_COUNT); + assertThat(result.getMaterializedRows().get(7).getField(1)).as("promoted value of row 7").isEqualTo("7"); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void testPredicateOnIntColumnEvolvedToBigintStillPrunes(boolean useParquetColumnNames) + throws Exception + { + HiveColumnHandle evolved = column(INT_TO_BIGINT_COLUMN, HiveType.HIVE_LONG, BIGINT); + List projection = List.of(column(STABLE_COLUMN, HiveType.HIVE_INT, INTEGER), evolved); + + MaterializedResult result = read(projection, greaterThanThreshold(evolved, BIGINT, THRESHOLD), + useParquetColumnNames, DynamicFilter.EMPTY); + + // asLong takes an Integer as happily as a Long, so this promotion is one the statistics CAN answer and the + // guard must leave it alone. This is what catches a check that drops more than it should. + assertThat(result.getRowCount()).as("rows read out of %s", ROW_COUNT).isLessThan(ROW_COUNT); + assertThat(valuesOver(result, 1)).as("rows matching %s > %s after pruning", INT_TO_BIGINT_COLUMN, THRESHOLD).isEqualTo(MATCHING_ROW_COUNT); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void testPredicateOnUnevolvedColumnStillPrunes(boolean useParquetColumnNames) + throws Exception + { + HiveColumnHandle stable = column(STABLE_COLUMN, HiveType.HIVE_INT, INTEGER); + List projection = List.of(stable); + + MaterializedResult result = read(projection, greaterThanThreshold(stable, INTEGER, THRESHOLD), + useParquetColumnNames, DynamicFilter.EMPTY); + + assertThat(result.getRowCount()).as("rows read out of %s", ROW_COUNT).isLessThan(ROW_COUNT); + assertThat(valuesOver(result, 0)).as("rows matching %s > %s after pruning", STABLE_COLUMN, THRESHOLD).isEqualTo(MATCHING_ROW_COUNT); + } + + @Test + public void testOnlyTheEvolvedColumnsDomainIsDropped() + throws Exception + { + HiveColumnHandle stable = column(STABLE_COLUMN, HiveType.HIVE_INT, INTEGER); + HiveColumnHandle evolved = column(FLOAT_TO_DOUBLE_COLUMN, HiveType.HIVE_DOUBLE, DOUBLE); + List projection = List.of(stable, evolved); + + MaterializedResult result = read(projection, + greaterThanThreshold(stable, INTEGER, THRESHOLD) + .intersect(greaterThanThreshold(evolved, DOUBLE, (double) THRESHOLD)), + false, DynamicFilter.EMPTY); + + // One unusable domain must not cost the whole predicate its pushdown: the stable column's domain still + // prunes, which is only visible because reading everything and reading nothing are both wrong here. + assertThat(result.getRowCount()).as("rows read out of %s", ROW_COUNT).isLessThan(ROW_COUNT); + assertThat(valuesOver(result, 0)).as("rows matching %s > %s after pruning", STABLE_COLUMN, THRESHOLD).isEqualTo(MATCHING_ROW_COUNT); + } + + @Test + public void testEvolvedColumnArrivingThroughADynamicFilter() + throws Exception + { + HiveColumnHandle evolved = column(FLOAT_TO_DOUBLE_COLUMN, HiveType.HIVE_DOUBLE, DOUBLE); + List projection = List.of(column(STABLE_COLUMN, HiveType.HIVE_INT, INTEGER), evolved); + + // A dynamic filter reaches getCombinedPredicate by its own route and its handles carry the metastore type + // just the same, so it has to be guarded on the same path + MaterializedResult result = read(projection, TupleDomain.all(), false, + dynamicFilterOn(greaterThanThreshold(evolved, DOUBLE, (double) THRESHOLD))); + + assertThat(result.getRowCount()).as("rows read").isEqualTo(ROW_COUNT); + } + + /** + * Reads the whole base file through the page source the connector builds for a split with no log files, which is + * the only path on which it enables predicate pushdown. + */ + private static MaterializedResult read( + List projection, + TupleDomain predicate, + boolean useParquetColumnNames, + DynamicFilter dynamicFilter) + throws Exception + { + long fileSize = Files.size(baseFile); + HudiSplit split = new HudiSplit( + new HudiBaseFile(baseFile.toString(), baseFile.getFileName().toString(), fileSize, 0, 0, fileSize), + List.of(), + "000", + predicate, + List.of(), + SplitWeight.standard()); + HudiSessionProperties sessionProperties = new HudiSessionProperties( + new HudiConfig().setUseParquetColumnNames(useParquetColumnNames), + new ParquetReaderConfig()); + ConnectorSession session = TestingConnectorSession.builder() + .setPropertyMetadata(sessionProperties.getSessionProperties()) + .build(); + + List types = projection.stream().map(HiveColumnHandle::getType).toList(); + try (ConnectorPageSource pageSource = createPageSource( + session, + projection, + split, + new LocalInputFile(baseFile.toFile()), + baseFile.toString(), + 0L, + fileSize, + OptionalLong.of(fileSize), + new FileFormatDataSourceStats(), + ParquetReaderOptions.builder().build(), + DateTimeZone.UTC, + dynamicFilter, + true)) { + return materializeSourceDataStream(session, pageSource, types).toTestTypes(); + } + } + + /** + * The base file as it was written BEFORE the evolution: the five {@code _hoodie_*} meta columns followed by the + * data columns in their original types. The metastore column list the handles below model reports the widened + * types instead, which is the whole point of the fixture. + */ + private static MessageType preEvolutionFileSchema() + { + List fields = new ArrayList<>(); + for (String metaColumn : HOODIE_META_COLUMNS) { + fields.add(Types.primitive(PrimitiveType.PrimitiveTypeName.BINARY, OPTIONAL).as(LogicalTypeAnnotation.stringType()).named(metaColumn)); + } + fields.add(Types.primitive(PrimitiveType.PrimitiveTypeName.INT32, OPTIONAL).named(STABLE_COLUMN)); + fields.add(Types.primitive(PrimitiveType.PrimitiveTypeName.FLOAT, OPTIONAL).named(FLOAT_TO_DOUBLE_COLUMN)); + fields.add(Types.primitive(PrimitiveType.PrimitiveTypeName.INT32, OPTIONAL).named(INT_TO_BIGINT_COLUMN)); + fields.add(Types.primitive(PrimitiveType.PrimitiveTypeName.INT32, OPTIONAL).named(INT_TO_VARCHAR_COLUMN)); + return new MessageType("hudi_base_file", fields); + } + + /** + * A handle as the metastore reports the column AFTER the evolution, on its physical ordinal. Stale ordinals are + * {@code TestHudiPageSourceProviderTest}'s subject, not this one, so the two resolution modes see the same + * column here and any difference between them is about the type alone. + */ + private static HiveColumnHandle column(String columnName, HiveType hiveType, Type trinoType) + { + return createBaseColumn(columnName, physicalIndexOf(columnName), hiveType, trinoType, + HiveColumnHandle.ColumnType.REGULAR, Optional.empty()); + } + + private static int physicalIndexOf(String columnName) + { + List fields = preEvolutionFileSchema().getFields(); + for (int i = 0; i < fields.size(); i++) { + if (fields.get(i).getName().equals(columnName)) { + return i; + } + } + throw new IllegalArgumentException("No such column in the fixture: " + columnName); + } + + private static TupleDomain greaterThanThreshold(HiveColumnHandle handle, Type type, Object threshold) + { + return TupleDomain.withColumnDomains(Map.of(handle, + Domain.create(ValueSet.ofRanges(Range.greaterThan(type, threshold)), false))); + } + + /** Counts the rows whose {@code fieldIndex}-th field is over {@link #THRESHOLD}, whatever numeric type it read as. */ + private static long valuesOver(MaterializedResult result, int fieldIndex) + { + return result.getMaterializedRows().stream() + .map(row -> row.getField(fieldIndex)) + .filter(value -> value != null && ((Number) value).doubleValue() > THRESHOLD) + .count(); + } + + private static int rowGroupCount(Path path) + throws IOException + { + try (ParquetFileReader reader = ParquetFileReader.open(new org.apache.parquet.io.LocalInputFile(path))) { + return reader.getRowGroups().size(); + } + } + + private static DynamicFilter dynamicFilterOn(TupleDomain predicate) + { + return new DynamicFilter() + { + @Override + public Set getColumnsCovered() + { + return Set.copyOf(predicate.getDomains().orElseThrow().keySet()); + } + + @Override + public CompletableFuture isBlocked() + { + return CompletableFuture.completedFuture(null); + } + + @Override + public boolean isComplete() + { + return true; + } + + @Override + public boolean isAwaitable() + { + return false; + } + + @Override + public TupleDomain getCurrentPredicate() + { + return predicate.transformKeys(ColumnHandle.class::cast); + } + }; + } +} diff --git a/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicates.java b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicates.java new file mode 100644 index 000000000000..0359bf41be39 --- /dev/null +++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicates.java @@ -0,0 +1,98 @@ +/* + * Licensed 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 io.trino.plugin.hudi; + +import io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer; +import io.trino.testing.AbstractTestQueryFramework; +import io.trino.testing.QueryRunner; +import org.junit.jupiter.api.Test; + +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.BIGINT_THRESHOLD; +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.DOUBLE_THRESHOLD; +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.FLOAT_TO_DOUBLE_COLUMN; +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.INT_TO_BIGINT_COLUMN; +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.INT_TO_VARCHAR_COLUMN; +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.TABLE_NAME; +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.VARCHAR_THRESHOLD; +import static io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer.expectedRowsFrom; + +/** + * apache/hudi#19457: a predicate on a column whose type was widened after a base file was written used to fail the + * whole query with {@code Malformed Parquet file. Corrupted statistics for column ...}, because the domain carries + * the metastore's widened type while the file's statistics are still of the original one. The connector now leaves + * such a domain out of the parquet predicate, and the engine applies it above the scan as it always did. + *

+ * Selecting the evolved column without constraining it was never affected, so every test here has to put a predicate + * ON the evolved column -- and every projection has to include it, otherwise the domain would resolve to no + * descriptor and be discarded for an unrelated reason, which is exactly the shape that passes against the unfixed + * code. {@link #testReadingAnEvolvedColumnWithoutAPredicate} is the anchor that shows the read path itself was + * always fine. + * + * @see TestHudiSchemaEvolutionPredicatesPositional for the same suite with columns resolved by ordinal + */ +public class TestHudiSchemaEvolutionPredicates + extends AbstractTestQueryFramework +{ + @Override + protected QueryRunner createQueryRunner() + throws Exception + { + return HudiQueryRunner.builder() + .addConnectorProperty("hudi.parquet.use-column-names", "true") + .setDataLoader(new SchemaEvolutionHudiTablesInitializer()) + .build(); + } + + @Test + public void testPredicateOnColumnEvolvedFromFloatToDouble() + { + assertQuery(selectWhere(FLOAT_TO_DOUBLE_COLUMN + " > " + DOUBLE_THRESHOLD), expectedRowsFrom(3)); + } + + @Test + public void testPredicateOnColumnEvolvedFromIntToBigint() + { + assertQuery(selectWhere(INT_TO_BIGINT_COLUMN + " > " + BIGINT_THRESHOLD), expectedRowsFrom(4)); + } + + @Test + public void testPredicateOnColumnEvolvedFromIntToVarchar() + { + assertQuery(selectWhere("%s > '%s'".formatted(INT_TO_VARCHAR_COLUMN, VARCHAR_THRESHOLD)), expectedRowsFrom(4)); + } + + @Test + public void testPredicatesOnEvolvedAndUnevolvedColumnsTogether() + { + // The key column never evolved, so its domain is pushed down while the evolved column's is dropped. One + // unusable domain must not cost the whole predicate its pushdown, nor the query its rows. + assertQuery( + selectWhere("%s > %s AND key >= 'k4'".formatted(FLOAT_TO_DOUBLE_COLUMN, DOUBLE_THRESHOLD)), + expectedRowsFrom(4)); + } + + @Test + public void testReadingAnEvolvedColumnWithoutAPredicate() + { + // The anchor: widening is a read-path feature that already worked, so a regression here would mean the + // fixture stopped modelling an evolved table rather than that the guard misbehaved + assertQuery(selectWhere("true"), expectedRowsFrom(1)); + } + + private static String selectWhere(String predicate) + { + return "SELECT key, %s, %s, %s FROM %s WHERE %s ORDER BY key".formatted( + FLOAT_TO_DOUBLE_COLUMN, INT_TO_BIGINT_COLUMN, INT_TO_VARCHAR_COLUMN, TABLE_NAME, predicate); + } +} diff --git a/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicatesPositional.java b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicatesPositional.java new file mode 100644 index 000000000000..846f8843d004 --- /dev/null +++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSchemaEvolutionPredicatesPositional.java @@ -0,0 +1,38 @@ +/* + * Licensed 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 io.trino.plugin.hudi; + +import io.trino.plugin.hudi.testing.SchemaEvolutionHudiTablesInitializer; +import io.trino.testing.QueryRunner; + +/** + * {@link TestHudiSchemaEvolutionPredicates} with {@code hudi.parquet.use-column-names=false}, so columns are resolved + * by ordinal instead of by name. apache/hudi#19457 was reported against both modes and neither has anything to do + * with the other: resolution decides WHICH file column a handle denotes, and the failure is about what its + * statistics are made of once it has been found. Running the suite twice is what keeps a fix aimed at one resolution + * path from quietly leaving the other broken. + */ +public class TestHudiSchemaEvolutionPredicatesPositional + extends TestHudiSchemaEvolutionPredicates +{ + @Override + protected QueryRunner createQueryRunner() + throws Exception + { + return HudiQueryRunner.builder() + .addConnectorProperty("hudi.parquet.use-column-names", "false") + .setDataLoader(new SchemaEvolutionHudiTablesInitializer()) + .build(); + } +} diff --git a/hudi-trino/src/test/java/io/trino/plugin/hudi/testing/SchemaEvolutionHudiTablesInitializer.java b/hudi-trino/src/test/java/io/trino/plugin/hudi/testing/SchemaEvolutionHudiTablesInitializer.java new file mode 100644 index 000000000000..bc1028fba89a --- /dev/null +++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/testing/SchemaEvolutionHudiTablesInitializer.java @@ -0,0 +1,164 @@ +/* + * Licensed 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 io.trino.plugin.hudi.testing; + +import com.google.common.collect.ImmutableList; +import io.trino.metastore.Column; +import org.apache.avro.Schema; +import org.apache.avro.generic.GenericData; +import org.apache.avro.generic.GenericRecord; +import org.apache.hudi.client.HoodieJavaWriteClient; +import org.apache.hudi.client.WriteStatus; +import org.apache.hudi.common.config.RecordMergeMode; +import org.apache.hudi.common.model.HoodieAvroPayload; +import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.config.HoodieWriteConfig; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.Optional; + +import static io.trino.metastore.HiveType.HIVE_DOUBLE; +import static io.trino.metastore.HiveType.HIVE_LONG; +import static io.trino.metastore.HiveType.HIVE_STRING; + +/** + * Creates a table whose base file was written BEFORE its columns were widened, while the metastore reports the types + * they were widened TO -- the state every unrewritten base file is in after a Hudi type evolution followed by a hive + * sync. The two halves are declared independently on purpose: {@link #avroSchema()} is what the write client puts in + * the file, {@link #dataColumns()} is what the metastore hands the connector, and only their disagreement is being + * modelled. No schema-evolution write path is exercised, because none is needed to reproduce apache/hudi#19457. + *

+ * The three widenings are the ones Hudi allows and the parquet reader can serve: {@code float -> double}, + * {@code int -> long} and {@code int -> string}. Reading them has always worked. Putting a predicate on one is what + * used to fail the query, because the statistics in the file are still of the original type while the pushed-down + * domain carries the widened one. + *

+ * Unlike {@link OmittedMetaColumnsHudiTablesInitializer} this fixture registers the Hudi meta columns, so every + * metastore ordinal equals its physical one and nothing here depends on how columns are resolved. What the two + * {@code hudi.parquet.use-column-names} modes must agree on is the TYPE handling alone. + *

+ * A single bulk-insert commit, so the file slice has no log files: predicate pushdown is only enabled for + * base-file-only splits. + */ +public class SchemaEvolutionHudiTablesInitializer + extends AbstractMergerHudiTablesInitializer +{ + public static final String TABLE_NAME = "schema_evolved_mor"; + + /** Written as Avro {@code float}, reported by the metastore as {@code double}. */ + public static final String FLOAT_TO_DOUBLE_COLUMN = "float_value"; + /** Written as Avro {@code int}, reported by the metastore as {@code bigint}. */ + public static final String INT_TO_BIGINT_COLUMN = "int_value"; + /** Written as Avro {@code int}, reported by the metastore as {@code string}. */ + public static final String INT_TO_VARCHAR_COLUMN = "string_value"; + + /** Sits between the third and the fourth row, so a predicate on it keeps some rows and drops others. */ + public static final String DOUBLE_THRESHOLD = "1003.0"; + public static final long BIGINT_THRESHOLD = 1003; + public static final String VARCHAR_THRESHOLD = "1003"; + + private static final int ROW_COUNT = 5; + private static final int BASE_VALUE = 1000; + + public SchemaEvolutionHudiTablesInitializer() + { + super(TABLE_NAME); + } + + @Override + protected List dataColumns() + { + return ImmutableList.of( + new Column(FLOAT_TO_DOUBLE_COLUMN, HIVE_DOUBLE, Optional.empty(), Map.of()), + new Column(INT_TO_BIGINT_COLUMN, HIVE_LONG, Optional.empty(), Map.of()), + new Column(INT_TO_VARCHAR_COLUMN, HIVE_STRING, Optional.empty(), Map.of()), + new Column(RECORD_KEY_FIELD, HIVE_STRING, Optional.empty(), Map.of()), + new Column(ORDERING_FIELD, HIVE_LONG, Optional.empty(), Map.of())); + } + + @Override + protected Schema avroSchema() + { + List fields = new ArrayList<>(); + fields.add(new Schema.Field(FLOAT_TO_DOUBLE_COLUMN, Schema.create(Schema.Type.FLOAT))); + fields.add(new Schema.Field(INT_TO_BIGINT_COLUMN, Schema.create(Schema.Type.INT))); + fields.add(new Schema.Field(INT_TO_VARCHAR_COLUMN, Schema.create(Schema.Type.INT))); + fields.add(new Schema.Field(RECORD_KEY_FIELD, Schema.create(Schema.Type.STRING))); + fields.add(new Schema.Field(ORDERING_FIELD, Schema.create(Schema.Type.LONG))); + return Schema.createRecord(TABLE_NAME, null, null, false, fields); + } + + @Override + protected void configureTableConfig(HoodieTableMetaClient.TableBuilder tableBuilder) + { + tableBuilder.setRecordMergeMode(RecordMergeMode.COMMIT_TIME_ORDERING); + } + + @Override + protected void configureWriteConfig(HoodieWriteConfig.Builder writeConfigBuilder) + { + writeConfigBuilder.withRecordMergeMode(RecordMergeMode.COMMIT_TIME_ORDERING); + } + + @Override + protected void writeInitialCommits(HoodieJavaWriteClient client) + { + Schema schema = avroSchema(); + List> records = new ArrayList<>(); + for (int row = 1; row <= ROW_COUNT; row++) { + records.add(record(schema, row)); + } + // One commit only: a file slice with log files would take the merge path, which disables pushdown. + String commit = client.startCommit(); + List statuses = client.bulkInsert(records, commit); + client.commit(commit, statuses); + } + + /** + * The expected rows of {@code SELECT key, float_value, int_value, string_value ... WHERE > } + * for the rows whose index is over {@code firstMatchingRow}. + *

+ * Every float value is an exact binary fraction, so widening it to double is lossless and the expected literals + * can be written out in full rather than compared with a tolerance. + */ + public static String expectedRowsFrom(int firstMatchingRow) + { + List rows = new ArrayList<>(); + for (int row = firstMatchingRow; row <= ROW_COUNT; row++) { + rows.add("('k%s', CAST(%s AS DOUBLE), CAST(%s AS BIGINT), '%s')".formatted( + row, floatValue(row), BASE_VALUE + row, BASE_VALUE + row)); + } + return "VALUES " + String.join(", ", rows); + } + + private static float floatValue(int row) + { + return BASE_VALUE + row + 0.5f; + } + + private static HoodieRecord record(Schema schema, int row) + { + String key = "k" + row; + GenericRecord record = new GenericData.Record(schema); + record.put(FLOAT_TO_DOUBLE_COLUMN, floatValue(row)); + record.put(INT_TO_BIGINT_COLUMN, BASE_VALUE + row); + record.put(INT_TO_VARCHAR_COLUMN, BASE_VALUE + row); + record.put(RECORD_KEY_FIELD, key); + record.put(ORDERING_FIELD, 100L); + return avroRecord(record, key); + } +} diff --git a/hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java b/hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java new file mode 100644 index 000000000000..277ae7ea50cf --- /dev/null +++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java @@ -0,0 +1,297 @@ +/* + * Licensed 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 io.trino.plugin.hudi.util; + +import io.trino.parquet.ParquetDataSourceId; +import io.trino.parquet.predicate.TupleDomainParquetPredicate; +import io.trino.spi.predicate.Domain; +import io.trino.spi.predicate.TupleDomain; +import io.trino.spi.type.DecimalType; +import io.trino.spi.type.Type; +import org.apache.parquet.column.ColumnDescriptor; +import org.apache.parquet.column.statistics.Statistics; +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.apache.parquet.schema.Types; +import org.joda.time.DateTimeZone; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.MethodSource; + +import java.nio.ByteBuffer; +import java.util.List; +import java.util.Map; + +import static io.trino.plugin.hudi.util.ParquetStatisticsDomains.dropIncomparableDomains; +import static io.trino.plugin.hudi.util.ParquetStatisticsDomains.hasComparableStatistics; +import static io.trino.plugin.hudi.util.TestParquetStatisticsDomains.LibraryOutcome.ALL; +import static io.trino.plugin.hudi.util.TestParquetStatisticsDomains.LibraryOutcome.NARROW; +import static io.trino.plugin.hudi.util.TestParquetStatisticsDomains.LibraryOutcome.THROWS; +import static io.trino.spi.type.BigintType.BIGINT; +import static io.trino.spi.type.BooleanType.BOOLEAN; +import static io.trino.spi.type.DateType.DATE; +import static io.trino.spi.type.DoubleType.DOUBLE; +import static io.trino.spi.type.IntegerType.INTEGER; +import static io.trino.spi.type.RealType.REAL; +import static io.trino.spi.type.TimestampType.TIMESTAMP_MILLIS; +import static io.trino.spi.type.TinyintType.TINYINT; +import static io.trino.spi.type.VarbinaryType.VARBINARY; +import static io.trino.spi.type.VarcharType.VARCHAR; +import static java.nio.ByteOrder.LITTLE_ENDIAN; +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.BINARY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.FLOAT; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT32; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT64; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT96; +import static org.apache.parquet.schema.Type.Repetition.OPTIONAL; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** + * Pins {@link ParquetStatisticsDomains#hasComparableStatistics} against the method it exists to protect. Every case + * below states BOTH what the guard decides and what {@code TupleDomainParquetPredicate.getDomain} actually does with + * the same pair, and the check is run against the real {@code getDomain}, not a description of it. A Trino upgrade + * that moves a branch therefore fails here, where the mismatch is a line of test output, instead of in a query. + *

+ * The three library outcomes are worth telling apart, because the guard exists for two different reasons: + *

    + *
  • {@link LibraryOutcome#THROWS} - the cast fails and the whole split dies with {@code HUDI_BAD_DATA}. This + * is apache/hudi#19457 as reported.
  • + *
  • {@link LibraryOutcome#NARROW} on a pair the guard drops - far worse: a domain IS produced, from bytes that + * mean something else entirely, and row groups get pruned on a comparison that is simply false. Nothing fails, + * rows just go missing.
  • + *
  • {@link LibraryOutcome#ALL} - the library declines the pair itself, so dropping it changes nothing.
  • + *
+ * The invariant that ties them together is asserted for every case: whatever the guard keeps must be a pair the + * library reads a real range out of. + */ +class TestParquetStatisticsDomains +{ + private static final ParquetDataSourceId DATA_SOURCE_ID = new ParquetDataSourceId("test"); + private static final long VALUE_COUNT = 10; + + enum LibraryOutcome + { + /** getDomain read the statistics and returned a range narrower than "any value". */ + NARROW, + /** getDomain declined to use the statistics and returned a domain covering every value. */ + ALL, + /** getDomain failed, which the connector reports as a corrupt-statistics error over the whole split. */ + THROWS, + } + + private record TypePair(String description, Type domainType, PrimitiveType fileType, boolean comparable, LibraryOutcome outcome) + { + @Override + public String toString() + { + return description; + } + } + + private static List typePairs() + { + return List.of( + // A column that never evolved: the domain's type is the one the file was written with + new TypePair("boolean over BOOLEAN", BOOLEAN, plain(PrimitiveTypeName.BOOLEAN), true, NARROW), + new TypePair("integer over INT32", INTEGER, plain(INT32), true, NARROW), + new TypePair("bigint over INT64", BIGINT, plain(INT64), true, NARROW), + new TypePair("tinyint over INT32", TINYINT, plain(INT32), true, NARROW), + new TypePair("date over INT32 date", DATE, annotated(INT32, LogicalTypeAnnotation.dateType()), true, NARROW), + new TypePair("real over FLOAT", REAL, plain(FLOAT), true, NARROW), + new TypePair("double over DOUBLE", DOUBLE, plain(PrimitiveTypeName.DOUBLE), true, NARROW), + new TypePair("varchar over BINARY string", VARCHAR, annotated(BINARY, LogicalTypeAnnotation.stringType()), true, NARROW), + new TypePair("decimal(9,2) over INT32 decimal(9,2)", DecimalType.createDecimalType(9, 2), decimal(INT32, 9, 2), true, NARROW), + new TypePair("timestamp over INT64 timestamp", TIMESTAMP_MILLIS, annotated(INT64, LogicalTypeAnnotation.timestampType(false, TimeUnit.MILLIS)), true, NARROW), + new TypePair("timestamp over INT96", TIMESTAMP_MILLIS, plain(INT96), true, NARROW), + + // Promotions the statistics can answer, so pushdown must survive them + new TypePair("int -> long", BIGINT, plain(INT32), true, NARROW), + new TypePair("decimal(9,2) -> decimal(9,4)", DecimalType.createDecimalType(9, 4), decimal(INT32, 9, 2), true, NARROW), + new TypePair("decimal(20,2) -> decimal(38,4)", DecimalType.createDecimalType(38, 4), decimal(FIXED_LEN_BYTE_ARRAY, 20, 2), true, NARROW), + new TypePair("integer over a zero-scale INT32 decimal", INTEGER, decimal(INT32, 9, 0), true, NARROW), + + // Promotions that fail the split today: apache/hudi#19457 and its neighbours + new TypePair("float -> double", DOUBLE, plain(FLOAT), false, THROWS), + new TypePair("int -> double", DOUBLE, plain(INT32), false, THROWS), + new TypePair("long -> double", DOUBLE, plain(INT64), false, THROWS), + new TypePair("int -> float", REAL, plain(INT32), false, THROWS), + new TypePair("int -> string", VARCHAR, plain(INT32), false, THROWS), + new TypePair("long -> string", VARCHAR, plain(INT64), false, THROWS), + new TypePair("float -> string", VARCHAR, plain(FLOAT), false, THROWS), + new TypePair("double -> string", VARCHAR, plain(PrimitiveTypeName.DOUBLE), false, THROWS), + new TypePair("string -> date", DATE, annotated(BINARY, LogicalTypeAnnotation.stringType()), false, THROWS), + + // Promotions that silently prune on a comparison that means nothing, which is why the guard cannot + // be a try/catch around the cast + new TypePair("decimal -> string", VARCHAR, decimal(FIXED_LEN_BYTE_ARRAY, 20, 2), false, NARROW), + new TypePair("string -> decimal", DecimalType.createDecimalType(9, 2), annotated(BINARY, LogicalTypeAnnotation.stringType()), false, NARROW), + new TypePair("int -> decimal", DecimalType.createDecimalType(9, 2), plain(INT32), false, NARROW), + new TypePair("integer over a scaled INT32 decimal", INTEGER, decimal(INT32, 9, 2), false, NARROW), + + // Pairs the library declines on its own, where dropping costs nothing + new TypePair("timestamp over an unannotated INT64", TIMESTAMP_MILLIS, plain(INT64), false, ALL), + new TypePair("varbinary over BINARY", VARBINARY, plain(BINARY), false, ALL)); + } + + @ParameterizedTest + @MethodSource("typePairs") + public void testGuardMatchesTheParquetPredicate(TypePair pair) + throws Exception + { + assertThat(hasComparableStatistics(pair.domainType(), pair.fileType())) + .as("guard verdict for %s", pair) + .isEqualTo(pair.comparable()); + + ColumnDescriptor descriptor = descriptorOf(pair.fileType()); + Statistics statistics = statisticsOf(pair.fileType()); + if (pair.outcome() == THROWS) { + assertThatThrownBy(() -> TupleDomainParquetPredicate.getDomain(descriptor, pair.domainType(), VALUE_COUNT, statistics, DATA_SOURCE_ID, DateTimeZone.UTC)) + .as("getDomain for %s", pair) + .hasMessageContaining("Corrupted statistics"); + return; + } + + Domain domain = TupleDomainParquetPredicate.getDomain(descriptor, pair.domainType(), VALUE_COUNT, statistics, DATA_SOURCE_ID, DateTimeZone.UTC); + assertThat(domain.getValues().isAll()) + .as("getDomain for %s returned %s", pair, domain) + .isEqualTo(pair.outcome() == ALL); + } + + @ParameterizedTest + @MethodSource("typePairs") + public void testEveryKeptPairYieldsUsableStatistics(TypePair pair) + { + // The invariant behind the whole table: a false negative only costs pruning, but a false positive is either + // a failed query or a wrong one, so nothing may be kept that the library does not read a real range out of. + if (hasComparableStatistics(pair.domainType(), pair.fileType())) { + assertThat(pair.outcome()).as("library outcome for the kept pair %s", pair).isEqualTo(NARROW); + } + } + + @Test + public void testAllAndNonePassThroughUntouched() + { + assertThat(dropIncomparableDomains(TupleDomain.all())).isEqualTo(TupleDomain.all()); + assertThat(dropIncomparableDomains(TupleDomain.none())).isEqualTo(TupleDomain.none()); + } + + @Test + public void testOnlyTheIncomparableDomainIsDropped() + { + // Distinct names on purpose: a ColumnDescriptor is keyed by its path, so two columns sharing one name would + // collapse into a single map entry and the test would pass without ever exercising the filtering + ColumnDescriptor evolved = descriptorNamed("evolved", FLOAT); + ColumnDescriptor stable = descriptorNamed("stable", INT32); + Domain doubleDomain = Domain.singleValue(DOUBLE, 1.0d); + Domain intDomain = Domain.singleValue(INTEGER, 1L); + + TupleDomain filtered = dropIncomparableDomains( + TupleDomain.withColumnDomains(Map.of(evolved, doubleDomain, stable, intDomain))); + + assertThat(filtered.getDomains().orElseThrow()).containsExactly(Map.entry(stable, intDomain)); + } + + @Test + public void testAPredicateWithNothingToDropIsReturnedAsIs() + { + TupleDomain predicate = TupleDomain.withColumnDomains( + Map.of(descriptorOf(plain(INT32)), Domain.singleValue(INTEGER, 1L))); + + assertThat(dropIncomparableDomains(predicate)).isSameAs(predicate); + } + + @Test + public void testDroppingEveryDomainLeavesAnUnconstrainedPredicate() + { + TupleDomain predicate = TupleDomain.withColumnDomains( + Map.of(descriptorOf(plain(FLOAT)), Domain.singleValue(DOUBLE, 1.0d))); + + // Not TupleDomain.none(): dropping means "do not prune on this", never "this matches nothing" + assertThat(dropIncomparableDomains(predicate).isAll()).as("everything dropped").isTrue(); + } + + private static PrimitiveType plain(PrimitiveTypeName primitiveTypeName) + { + if (primitiveTypeName == INT96) { + return Types.primitive(primitiveTypeName, OPTIONAL).named("c"); + } + if (primitiveTypeName == FIXED_LEN_BYTE_ARRAY) { + return Types.primitive(primitiveTypeName, OPTIONAL).length(16).named("c"); + } + return Types.primitive(primitiveTypeName, OPTIONAL).named("c"); + } + + private static PrimitiveType annotated(PrimitiveTypeName primitiveTypeName, LogicalTypeAnnotation annotation) + { + return Types.primitive(primitiveTypeName, OPTIONAL).as(annotation).named("c"); + } + + private static PrimitiveType decimal(PrimitiveTypeName primitiveTypeName, int precision, int scale) + { + Types.PrimitiveBuilder builder = Types.primitive(primitiveTypeName, OPTIONAL); + if (primitiveTypeName == FIXED_LEN_BYTE_ARRAY) { + builder = builder.length(16); + } + return builder.as(LogicalTypeAnnotation.decimalType(scale, precision)).named("c"); + } + + private static ColumnDescriptor descriptorOf(PrimitiveType fileType) + { + return new ColumnDescriptor(new String[] {fileType.getName()}, fileType, 0, 1); + } + + private static ColumnDescriptor descriptorNamed(String name, PrimitiveTypeName primitiveTypeName) + { + return descriptorOf(Types.primitive(primitiveTypeName, OPTIONAL).named(name)); + } + + /** + * Statistics holding a small, non-degenerate range, so that a pair the library CAN read produces a domain + * narrower than "any value" and the {@link LibraryOutcome#NARROW} cases stay distinguishable from + * {@link LibraryOutcome#ALL}. Two nearby values also keep every integral type clear of + * {@code isStatisticsOverflow}, which would otherwise widen tinyint back to everything. + */ + private static Statistics statisticsOf(PrimitiveType fileType) + { + return Statistics.getBuilderForReading(fileType) + .withMin(statisticsBytes(fileType, 1)) + .withMax(statisticsBytes(fileType, 2)) + .withNumNulls(0) + .build(); + } + + private static byte[] statisticsBytes(PrimitiveType fileType, int value) + { + return switch (fileType.getPrimitiveTypeName()) { + // Both bounds false, so the boolean branch reports "only false" rather than "true and false", which it + // would report as every value + case BOOLEAN -> new byte[] {0}; + case INT32 -> ByteBuffer.allocate(4).order(LITTLE_ENDIAN).putInt(value).array(); + case INT64 -> ByteBuffer.allocate(8).order(LITTLE_ENDIAN).putLong(value).array(); + case FLOAT -> ByteBuffer.allocate(4).order(LITTLE_ENDIAN).putFloat(value).array(); + case DOUBLE -> ByteBuffer.allocate(8).order(LITTLE_ENDIAN).putDouble(value).array(); + // INT96 statistics are only usable when the bounds are equal (PARQUET-1065), so ignore the value: + // 8 bytes of nanos-within-the-day followed by the julian day of the epoch + case INT96 -> ByteBuffer.allocate(12).order(LITTLE_ENDIAN).putLong(0).putInt(2440588).array(); + case BINARY -> Integer.toString(value).getBytes(UTF_8); + // A decimal's unscaled value is big-endian two's complement, padded to the column's length + case FIXED_LEN_BYTE_ARRAY -> ByteBuffer.allocate(fileType.getTypeLength()).put(fileType.getTypeLength() - 1, (byte) value).array(); + }; + } +} From cd48bb14dbe4a917c82e768ad8364b64e0c3ff4f Mon Sep 17 00:00:00 2001 From: Vova Kolmakov Date: Thu, 6 Aug 2026 10:09:12 +0700 Subject: [PATCH 2/2] test(trino): remove the dead INT96 branch in the parquet statistics test helper --- .../trino/plugin/hudi/util/TestParquetStatisticsDomains.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java b/hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java index 277ae7ea50cf..eeafae9a75dc 100644 --- a/hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java +++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/util/TestParquetStatisticsDomains.java @@ -228,9 +228,6 @@ public void testDroppingEveryDomainLeavesAnUnconstrainedPredicate() private static PrimitiveType plain(PrimitiveTypeName primitiveTypeName) { - if (primitiveTypeName == INT96) { - return Types.primitive(primitiveTypeName, OPTIONAL).named("c"); - } if (primitiveTypeName == FIXED_LEN_BYTE_ARRAY) { return Types.primitive(primitiveTypeName, OPTIONAL).length(16).named("c"); }