From 95803566ec8985e5c3f2102b0b7eb71c265c71da Mon Sep 17 00:00:00 2001 From: Denys Kuzmenko Date: Mon, 31 Aug 2026 11:58:25 +0300 Subject: [PATCH 1/2] HIVE-29857: Iceberg: Fix ROW__POSITION on vectorized VARIANT reads ROW__POSITION came back as Long.MIN_VALUE for every VARIANT table read vectorized, and row lineage and positional deletes read the same value. HiveBatchIterator takes the position of a batch from the record reader, but only when the reader is a RowPositionAwareVectorizedRecordReader: if (batch.size != 0 && recordReader instanceof RowPositionAwareVectorizedRecordReader) { rowOffset = ((RowPositionAwareVectorizedRecordReader) recordReader).getRowNumber(); } A VARIANT table's reader is wrapped in ParquetVariantRecordReader, which declared only RecordReader, so the test failed and rowOffset kept the marker that stands for an unknown position. Let the wrapper carry the interface and answer from the reader it wraps, and let it fail rather than hand back the marker. That uncovered a second fault underneath. Variant row group pruning handed the reader a footer holding only the row groups that survived, and the reader counted row positions by walking that footer and summing row counts. A position is absolute within the file, so counting it over a footer that is missing row groups placed every row group after a pruned one too low - the first row of the second row group reported position 0 rather than 200. The same footer also fed the scan statistics, which under reported the rows and bytes a variant scan covers. The reader now keeps the file's own footer and is told separately which row groups to read, as ORC's readers are. Parquet records where each row group's first row sits in the file, so the position is read from the footer rather than counted up, and stays right whatever subset of row groups a split reads. VariantParquetFilters.pickRowGroups already returned that decision as a boolean per row group for the non-vectorized reader; the vectorized path now takes the same answer instead of a pruned footer, so pruneVariantRowGroups has no caller left and is removed. variant_type_row_position.q prunes an early row group and checks that the surviving rows keep the positions they were written at. Without the first fix the positions are the unknown marker, without the second they are short by the pruned row group's row count. --- .../mr/hive/vector/HiveVectorizedReader.java | 11 ++- .../vector/ParquetVariantRecordReader.java | 18 +++- .../org/apache/iceberg/parquet/ReadConf.java | 6 +- .../parquet/VariantParquetFilters.java | 51 +---------- .../mr/hive/TestHiveIcebergVariantType.java | 26 ++++-- .../positive/variant_type_row_position.q | 39 ++++++++ .../positive/variant_type_row_position.q.out | 88 +++++++++++++++++++ .../io/parquet/ParquetRecordReaderBase.java | 15 +++- .../parquet/VectorizedParquetInputFormat.java | 12 ++- .../vector/VectorizedParquetRecordReader.java | 21 ++++- 10 files changed, 215 insertions(+), 72 deletions(-) create mode 100644 iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q create mode 100644 iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out diff --git a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java index 0bbc5aa8f08e..7fe7f8f64ecb 100644 --- a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java +++ b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java @@ -266,9 +266,12 @@ private static RecordReader parquetRecordReade ParquetMetadata parquetMetadata = HiveParquetUtil.readFooter(task.file(), io, job, footerData); MessageType fileSchema = parquetMetadata.getFileMetaData().getSchema(); - ParquetMetadata prunedMetadata = - VariantParquetFilters.pruneVariantRowGroups(parquetMetadata, fileSchema, residual); - inputFormat.setMetadata(prunedMetadata); + // The reader keeps the file's own footer, so row positions and scan statistics still describe the + // file. Row groups variant pruning ruled out are named separately; null means read them all. + inputFormat.setMetadata(parquetMetadata); + + inputFormat.setIncludedRowGroups( + VariantParquetFilters.pickRowGroups(fileSchema, residual, parquetMetadata.getBlocks())); MessageType typeWithIds = null; Schema expectedSchema = task.spec().schema(); @@ -287,7 +290,7 @@ private static RecordReader parquetRecordReade inputFormat.seInitialColumnDefaults(initialColumnDefaults); RecordReader reader = inputFormat.getRecordReader(split, job, reporter); return ParquetVariantRecordReader - .tryWrap(reader, job, task, path, start, length, prunedMetadata) + .tryWrap(reader, job, task, path, start, length, parquetMetadata) .orElse(reader); } diff --git a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java index 7ce1f9645f68..93d49ce3361e 100644 --- a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java +++ b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java @@ -31,6 +31,7 @@ import org.apache.hadoop.hive.ql.exec.vector.ColumnVector; import org.apache.hadoop.hive.ql.exec.vector.StructColumnVector; import org.apache.hadoop.hive.ql.exec.vector.VectorizedRowBatch; +import org.apache.hadoop.hive.ql.io.RowPositionAwareVectorizedRecordReader; import org.apache.hadoop.hive.ql.io.parquet.ParquetRecordReaderBase; import org.apache.hadoop.io.NullWritable; import org.apache.hadoop.mapred.JobConf; @@ -52,7 +53,8 @@ import org.apache.parquet.schema.MessageType; import org.apache.parquet.schema.Type; -final class ParquetVariantRecordReader implements RecordReader { +final class ParquetVariantRecordReader + implements RecordReader, RowPositionAwareVectorizedRecordReader { private static final String INVALID_VARIANT_STRUCT = "Invalid Variant struct for column "; @@ -154,6 +156,20 @@ private static List blocksForSplit( return splitBlocks; } + /** + * Row positions come from the reader this wraps. Without this the wrapper hides the delegate's + * {@link RowPositionAwareVectorizedRecordReader}, and every row of a VARIANT table is handed the unknown + * position marker instead - which ROW__POSITION, row lineage and positional deletes all rely on. + */ + @Override + public long getRowNumber() throws IOException { + if (delegate instanceof RowPositionAwareVectorizedRecordReader positionAware) { + return positionAware.getRowNumber(); + } + throw new UnsupportedOperationException( + "The reader under " + delegate.getClass().getName() + " cannot report row positions"); + } + @Override public boolean next(NullWritable key, VectorizedRowBatch value) throws IOException { boolean hasNext = delegate.next(key, value); diff --git a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/ReadConf.java b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/ReadConf.java index c350a3d7d562..a62430ba8008 100644 --- a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/ReadConf.java +++ b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/ReadConf.java @@ -101,14 +101,14 @@ class ReadConf { bloomFilter = new ParquetBloomRowGroupFilter(expectedSchema, filter, caseSensitive); } - boolean[] variantRowGroupMayMatch = - VariantParquetFilters.variantRowGroupMayMatch(fileSchema, filter, rowGroups); + boolean[] mayMatch = + VariantParquetFilters.pickRowGroups(fileSchema, filter, rowGroups); long computedTotalValues = 0L; for (int i = 0; i < shouldSkip.length; i += 1) { BlockMetaData rowGroup = rowGroups.get(i); boolean shouldRead = - (variantRowGroupMayMatch == null || variantRowGroupMayMatch[i]) && + (mayMatch == null || mayMatch[i]) && (filter == null || statsFilter.shouldRead(typeWithIds, rowGroup) && dictFilter.shouldRead(typeWithIds, rowGroup, reader.getDictionaryReader(rowGroup)) && diff --git a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java index 32149e28314f..aac163ea1ff4 100644 --- a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java +++ b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java @@ -41,7 +41,6 @@ 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.ParquetMetadata; import org.apache.parquet.io.api.Binary; import org.apache.parquet.schema.GroupType; import org.apache.parquet.schema.LogicalTypeAnnotation.StringLogicalTypeAnnotation; @@ -95,7 +94,7 @@ private static ResolvedVariantFilter resolveVariantFilter(MessageType schema, Ex return new ResolvedVariantFilter(predicate, visitor.fallbackValueColumns()); } - public static boolean[] variantRowGroupMayMatch( + public static boolean[] pickRowGroups( MessageType fileSchema, Expression filter, List rowGroups) { if (fileSchema == null || filter == null || rowGroups == null || rowGroups.isEmpty()) { return null; @@ -135,54 +134,6 @@ private static boolean mayMatchViaFallback(BlockMetaData rowGroup, Set pruneVariantRowGroups( - MessageType fileSchema, Expression filter, List rowGroups) { - boolean[] mayMatch = variantRowGroupMayMatch(fileSchema, filter, rowGroups); - if (mayMatch == null) { - return rowGroups; - } - - List kept = Lists.newArrayListWithCapacity(rowGroups.size()); - for (int i = 0; i < rowGroups.size(); i++) { - if (mayMatch[i]) { - kept.add(rowGroups.get(i)); - } - } - - return kept.size() == rowGroups.size() ? rowGroups : kept; - } - - /** Returns Parquet metadata with row groups pruned using best-effort VARIANT pruning. */ - public static ParquetMetadata pruneVariantRowGroups( - ParquetMetadata parquetMetadata, MessageType fileSchema, Expression filter) { - if (parquetMetadata == null || filter == null) { - return parquetMetadata; - } - - List rowGroups = parquetMetadata.getBlocks(); - if (rowGroups == null || rowGroups.isEmpty()) { - return parquetMetadata; - } - - MessageType schema = fileSchema; - if (schema == null) { - if (parquetMetadata.getFileMetaData() == null) { - return parquetMetadata; - } - schema = parquetMetadata.getFileMetaData().getSchema(); - } - if (schema == null) { - return parquetMetadata; - } - - List kept = pruneVariantRowGroups(schema, filter, rowGroups); - if (kept == rowGroups || parquetMetadata.getFileMetaData() == null) { - return parquetMetadata; - } - - return new ParquetMetadata(parquetMetadata.getFileMetaData(), kept); - } - private static boolean isColumnAllNull(ColumnChunkMetaData meta) { if (meta == null || meta.getStatistics() == null || !meta.getStatistics().isNumNullsSet()) { return false; diff --git a/iceberg/iceberg-handler/src/test/java/org/apache/iceberg/mr/hive/TestHiveIcebergVariantType.java b/iceberg/iceberg-handler/src/test/java/org/apache/iceberg/mr/hive/TestHiveIcebergVariantType.java index 107b40a817d0..ce9bff630c69 100644 --- a/iceberg/iceberg-handler/src/test/java/org/apache/iceberg/mr/hive/TestHiveIcebergVariantType.java +++ b/iceberg/iceberg-handler/src/test/java/org/apache/iceberg/mr/hive/TestHiveIcebergVariantType.java @@ -339,21 +339,31 @@ private static boolean mapVectorized(List explain) { private static void assertVectorizedParquetRowGroupsPruned(Path parquetPath, Expression filter) { assertParquetRowGroupsPruned( parquetPath, filter, - (parquetMetadata, fileSchema, expr) -> - // Simulate what HiveVectorizedReader.parquetRecordReader() does - VariantParquetFilters - .pruneVariantRowGroups(parquetMetadata, fileSchema, expr) - .getBlocks() - .size()); + (parquetMetadata, fileSchema, expr) -> { + // Simulate what HiveVectorizedReader.parquetRecordReader() does: the footer stays whole and the + // row groups to read are named separately + boolean[] mayMatch = VariantParquetFilters + .pickRowGroups(fileSchema, expr, parquetMetadata.getBlocks()); + if (mayMatch == null) { + return parquetMetadata.getBlocks().size(); + } + int matching = 0; + for (boolean match : mayMatch) { + if (match) { + matching++; + } + } + return matching; + }); } private static void assertNonVectorizedParquetRowGroupsPruned(Path parquetPath, Expression filter) { assertParquetRowGroupsPruned( parquetPath, filter, (parquetMetadata, fileSchema, expr) -> { - // Simulate what ReadConf does - uses variantRowGroupMayMatch to compute shouldSkip array + // Simulate what ReadConf does - uses pickRowGroups to compute shouldSkip array boolean[] mayMatch = VariantParquetFilters - .variantRowGroupMayMatch(fileSchema, expr, parquetMetadata.getBlocks()); + .pickRowGroups(fileSchema, expr, parquetMetadata.getBlocks()); int matching = 0; for (boolean match : mayMatch) { if (match) { diff --git a/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q b/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q new file mode 100644 index 000000000000..c0d6a77b9980 --- /dev/null +++ b/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q @@ -0,0 +1,39 @@ +-- Row positions are absolute within a data file. Variant row-group pruning must not shift them: a row +-- group the predicate drops still occupies its rows, so every later row group keeps the position it had. +set hive.explain.user=false; +set hive.fetch.task.conversion=none; +set hive.vectorized.execution.enabled=true; + +drop table if exists variant_row_position; + +CREATE EXTERNAL TABLE variant_row_position ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024' +); + +-- The first rows carry tier=bronze, the later ones tier=gold, and the small row group size puts them in +-- different row groups. A predicate on tier drops the bronze ones, which is what shifts the positions of +-- the gold ones when pruning is applied to the footer the reader counts over. +INSERT INTO variant_row_position +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val; + +-- ROW__POSITION of the surviving rows must match their id, which was written in file order. +SELECT id, variant_row_position.ROW__POSITION +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +ORDER BY id; + +-- the lowest surviving position is the first gold row, not zero +SELECT min(variant_row_position.ROW__POSITION) AS first_gold_position, + max(variant_row_position.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold'; + +drop table variant_row_position; diff --git a/iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out b/iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out new file mode 100644 index 000000000000..2f1fac1f4b41 --- /dev/null +++ b/iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out @@ -0,0 +1,88 @@ +PREHOOK: query: drop table if exists variant_row_position +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_row_position +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_row_position ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position +POSTHOOK: query: CREATE EXTERNAL TABLE variant_row_position ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position +PREHOOK: query: INSERT INTO variant_row_position +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_row_position +POSTHOOK: query: INSERT INTO variant_row_position +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_row_position +PREHOOK: query: SELECT id, variant_row_position.ROW__POSITION +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +ORDER BY id +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT id, variant_row_position.ROW__POSITION +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +ORDER BY id +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position +POSTHOOK: Output: hdfs://### HDFS PATH ### +200 200 +201 201 +202 202 +203 203 +204 204 +PREHOOK: query: SELECT min(variant_row_position.ROW__POSITION) AS first_gold_position, + max(variant_row_position.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT min(variant_row_position.ROW__POSITION) AS first_gold_position, + max(variant_row_position.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position +POSTHOOK: Output: hdfs://### HDFS PATH ### +200 399 200 +PREHOOK: query: drop table variant_row_position +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_row_position +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position +POSTHOOK: query: drop table variant_row_position +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_row_position +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position diff --git a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/ParquetRecordReaderBase.java b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/ParquetRecordReaderBase.java index 50c30e2941a3..b4ea3f399c2c 100644 --- a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/ParquetRecordReaderBase.java +++ b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/ParquetRecordReaderBase.java @@ -95,6 +95,15 @@ protected void setupMetadataAndParquetSplit(JobConf conf, ParquetMetadata metada // having null as parquetInputSplit seems to be a valid case based on this file's history } + /** + * Whether the row group at this index in the file's footer is one to read. A caller that has ruled some + * out with a filter Parquet cannot express drops them here, which leaves the footer the file's own - and + * the footer is what row positions are read from. + */ + protected boolean includeRowGroup(int rowGroupIndex) { + return true; + } + /** * gets a ParquetInputSplit corresponding to a split given by Hive * @@ -130,9 +139,11 @@ protected ParquetInputSplit getSplit( final List splitGroup = new ArrayList(); final long splitStart = fileSplit.getStart(); final long splitLength = fileSplit.getLength(); - for (final BlockMetaData block : blocks) { + for (int rowGroupIndex = 0; rowGroupIndex < blocks.size(); ++rowGroupIndex) { + final BlockMetaData block = blocks.get(rowGroupIndex); final long firstDataPage = block.getColumns().get(0).getFirstDataPageOffset(); - if (firstDataPage >= splitStart && firstDataPage < splitStart + splitLength) { + if (firstDataPage >= splitStart && firstDataPage < splitStart + splitLength + && includeRowGroup(rowGroupIndex)) { splitGroup.add(block); } } diff --git a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java index a9b1402c9ddc..a24ed83e8489 100644 --- a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java +++ b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java @@ -46,6 +46,7 @@ public class VectorizedParquetInputFormat private DataCache dataCache = null; private Configuration cacheConf = null; private ParquetMetadata metadata; + private boolean[] includedRowGroups; private Map initialDefaults; public VectorizedParquetInputFormat() { @@ -57,13 +58,22 @@ public RecordReader getRecordReader( JobConf jobConf, Reporter reporter) throws IOException { return new VectorizedParquetRecordReader( - inputSplit, jobConf, metadataCache, dataCache, cacheConf, metadata, initialDefaults); + inputSplit, jobConf, metadataCache, dataCache, cacheConf, metadata, initialDefaults, includedRowGroups); } public void setMetadata(ParquetMetadata metadata) throws IOException { this.metadata = metadata; } + /** + * Restricts the reader to these of the file's row groups, by footer index, for a caller that has already + * ruled some out with a filter Parquet cannot express. The footer given to {@link #setMetadata} stays the + * file's own, so row positions remain absolute. + */ + public void setIncludedRowGroups(boolean[] includedRowGroups) { + this.includedRowGroups = includedRowGroups; + } + public void seInitialColumnDefaults(Map initialDefaults) { this.initialDefaults = initialDefaults; } diff --git a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java index 236f6f3095f0..c4d0237cf7be 100644 --- a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java +++ b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java @@ -107,6 +107,9 @@ public class VectorizedParquetRecordReader extends ParquetRecordReaderBase protected MessageType requestedSchema; private List columnNamesList; private List columnTypesList; + /** Which of the file's row groups to read, by footer index, or null to read every one. */ + private boolean[] includedRowGroups; + private VectorizedRowBatchCtx rbCtx; private Object[] partitionValues; private boolean addPartitionCols = true; @@ -160,8 +163,15 @@ public VectorizedParquetRecordReader(InputSplit oldInputSplit, JobConf conf) thr public VectorizedParquetRecordReader(InputSplit oldInputSplit, JobConf conf, FileMetadataCache metadataCache, DataCache dataCache, Configuration cacheConf, ParquetMetadata parquetMetadata, Map initialDefaults) throws IOException { + this(oldInputSplit, conf, metadataCache, dataCache, cacheConf, parquetMetadata, initialDefaults, null); + } + + public VectorizedParquetRecordReader(InputSplit oldInputSplit, JobConf conf, FileMetadataCache metadataCache, + DataCache dataCache, Configuration cacheConf, ParquetMetadata parquetMetadata, + Map initialDefaults, boolean[] includedRowGroups) throws IOException { super(conf, oldInputSplit); try { + this.includedRowGroups = includedRowGroups; this.metadataCache = metadataCache; this.cache = dataCache; this.cacheConf = cacheConf; @@ -205,6 +215,11 @@ public VectorizedParquetRecordReader(InputSplit oldInputSplit, JobConf conf, Fil this(oldInputSplit, conf, metadataCache, dataCache, cacheConf, null, null); } + @Override + protected boolean includeRowGroup(int rowGroupIndex) { + return includedRowGroups == null || includedRowGroups[rowGroupIndex]; + } + private void initPartitionValues(FileSplit fileSplit, JobConf conf) throws IOException { int partitionColumnCount = rbCtx.getPartitionColumnCount(); if (partitionColumnCount > 0) { @@ -241,14 +256,14 @@ public void initialize( offsets.add(offset); } blocks = new ArrayList<>(); - long allRowsInFile = 0; int blockIndex = 0; for (BlockMetaData block : parquetMetadata.getBlocks()) { if (offsets.contains(block.getStartingPos())) { - rowGroupNumToRowPos.put(blockIndex++, allRowsInFile); + // Parquet records where each row group's first row sits in the file, so the position is read + // rather than counted up, and stays right whatever subset of row groups this split reads. + rowGroupNumToRowPos.put(blockIndex++, block.getRowIndexOffset()); blocks.add(block); } - allRowsInFile += block.getRowCount(); } // verify we found them all if (blocks.size() != rowGroupOffsets.length) { From ee45098cc2731058ee13fc5164275cd1c961b6a3 Mon Sep 17 00:00:00 2001 From: Denys Kuzmenko Date: Mon, 31 Aug 2026 19:16:47 +0300 Subject: [PATCH 2/2] HIVE-29857: Review fixes and coverage for VARIANT row positions --- .../mr/hive/vector/HiveVectorizedReader.java | 6 +- .../vector/ParquetVariantRecordReader.java | 9 +- .../parquet/VariantParquetFilters.java | 4 + .../positive/variant_type_row_position.q | 56 ++++ .../llap/variant_type_row_position.q.out | 247 ++++++++++++++++++ .../positive/variant_type_row_position.q.out | 159 +++++++++++ .../resources/testconfiguration.properties | 1 + .../parquet/VectorizedParquetInputFormat.java | 18 +- .../vector/VectorizedParquetRecordReader.java | 10 +- .../parquet/TestVectorizedColumnReader.java | 52 ++++ 10 files changed, 544 insertions(+), 18 deletions(-) create mode 100644 iceberg/iceberg-handler/src/test/results/positive/llap/variant_type_row_position.q.out diff --git a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java index 7fe7f8f64ecb..17fe3d7cc6e6 100644 --- a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java +++ b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java @@ -267,10 +267,8 @@ private static RecordReader parquetRecordReade ParquetMetadata parquetMetadata = HiveParquetUtil.readFooter(task.file(), io, job, footerData); MessageType fileSchema = parquetMetadata.getFileMetaData().getSchema(); // The reader keeps the file's own footer, so row positions and scan statistics still describe the - // file. Row groups variant pruning ruled out are named separately; null means read them all. - inputFormat.setMetadata(parquetMetadata); - - inputFormat.setIncludedRowGroups( + // file. Row groups variant pruning ruled out are named alongside it; null means read them all. + inputFormat.setMetadata(parquetMetadata, VariantParquetFilters.pickRowGroups(fileSchema, residual, parquetMetadata.getBlocks())); MessageType typeWithIds = null; diff --git a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java index 93d49ce3361e..860dbdf739ea 100644 --- a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java +++ b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/ParquetVariantRecordReader.java @@ -23,6 +23,7 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.Optional; import java.util.function.ToIntFunction; @@ -139,11 +140,11 @@ private static List blocksForSplit( // If the underlying Hive Parquet reader already computed row-group filtering (e.g. from SARG), // we must use the exact same blocks to keep this reader aligned with the delegate. if (delegate instanceof ParquetRecordReaderBase parquetDelegate) { + // The delegate's answer is the whole answer, including when it holds none: an empty list means it + // filtered every row group out, and null means it found none to read at all. Falling back to the + // split's own row groups would read one the delegate has already ruled out. List filteredBlocks = parquetDelegate.getFilteredBlocks(); - // Treat an empty list as authoritative (delegate filtered out all row groups). - if (filteredBlocks != null) { - return filteredBlocks; - } + return filteredBlocks != null ? filteredBlocks : Collections.emptyList(); } // Fallback: compute blocks from split boundaries List splitBlocks = Lists.newArrayList(); diff --git a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java index aac163ea1ff4..170889f27938 100644 --- a/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java +++ b/iceberg/iceberg-handler/src/main/java/org/apache/iceberg/parquet/VariantParquetFilters.java @@ -94,6 +94,10 @@ private static ResolvedVariantFilter resolveVariantFilter(MessageType schema, Ex return new ResolvedVariantFilter(predicate, visitor.fallbackValueColumns()); } + /** + * Which of these row groups a VARIANT predicate could match, one flag each, or null when the predicate + * says nothing about them and every row group is to be read. + */ public static boolean[] pickRowGroups( MessageType fileSchema, Expression filter, List rowGroups) { if (fileSchema == null || filter == null || rowGroups == null || rowGroups.isEmpty()) { diff --git a/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q b/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q index c0d6a77b9980..6b640e7e5234 100644 --- a/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q +++ b/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q @@ -37,3 +37,59 @@ FROM variant_row_position WHERE variant_get(data, '$.tier', 'string') = 'gold'; drop table variant_row_position; + +-- A file read as several splits: each split reports positions from the file's own start, so the row groups +-- a later split reads must not be numbered as though its split began the file. +drop table if exists variant_row_position_split; + +CREATE EXTERNAL TABLE variant_row_position_split ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'read.split.target-size'='1024' +); + +INSERT INTO variant_row_position_split +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val; + +SELECT min(variant_row_position_split.ROW__POSITION) AS first_gold_position, + max(variant_row_position_split.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position_split +WHERE variant_get(data, '$.tier', 'string') = 'gold'; + +drop table variant_row_position_split; + +-- A positional delete addresses rows by position, so a position shifted by pruning deletes the wrong row. +-- Here the delete predicate prunes row groups on the read side while the positions are being recorded. +drop table if exists variant_row_position_del; + +CREATE EXTERNAL TABLE variant_row_position_del ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'write.delete.mode'='merge-on-read' +); + +INSERT INTO variant_row_position_del +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val; + +DELETE FROM variant_row_position_del +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205; + +-- exactly ids 200-204 are gone: 395 rows left, and the ids either side of the hole are untouched +SELECT count(*) AS rows_left, min(id) AS lowest_id, max(id) AS highest_id FROM variant_row_position_del; + +SELECT id FROM variant_row_position_del WHERE id BETWEEN 197 AND 208 ORDER BY id; + +drop table variant_row_position_del; diff --git a/iceberg/iceberg-handler/src/test/results/positive/llap/variant_type_row_position.q.out b/iceberg/iceberg-handler/src/test/results/positive/llap/variant_type_row_position.q.out new file mode 100644 index 000000000000..78e8611e5966 --- /dev/null +++ b/iceberg/iceberg-handler/src/test/results/positive/llap/variant_type_row_position.q.out @@ -0,0 +1,247 @@ +PREHOOK: query: drop table if exists variant_row_position +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_row_position +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_row_position ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position +POSTHOOK: query: CREATE EXTERNAL TABLE variant_row_position ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position +PREHOOK: query: INSERT INTO variant_row_position +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_row_position +POSTHOOK: query: INSERT INTO variant_row_position +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_row_position +PREHOOK: query: SELECT id, variant_row_position.ROW__POSITION +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +ORDER BY id +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position +#### A masked pattern was here #### +POSTHOOK: query: SELECT id, variant_row_position.ROW__POSITION +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +ORDER BY id +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position +#### A masked pattern was here #### +200 200 +201 201 +202 202 +203 203 +204 204 +PREHOOK: query: SELECT min(variant_row_position.ROW__POSITION) AS first_gold_position, + max(variant_row_position.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position +#### A masked pattern was here #### +POSTHOOK: query: SELECT min(variant_row_position.ROW__POSITION) AS first_gold_position, + max(variant_row_position.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position +WHERE variant_get(data, '$.tier', 'string') = 'gold' +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position +#### A masked pattern was here #### +200 399 200 +PREHOOK: query: drop table variant_row_position +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_row_position +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position +POSTHOOK: query: drop table variant_row_position +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_row_position +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position +PREHOOK: query: drop table if exists variant_row_position_split +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_row_position_split +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_row_position_split ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'read.split.target-size'='1024' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_split +POSTHOOK: query: CREATE EXTERNAL TABLE variant_row_position_split ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'read.split.target-size'='1024' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_split +PREHOOK: query: INSERT INTO variant_row_position_split +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_row_position_split +POSTHOOK: query: INSERT INTO variant_row_position_split +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_row_position_split +PREHOOK: query: SELECT min(variant_row_position_split.ROW__POSITION) AS first_gold_position, + max(variant_row_position_split.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position_split +WHERE variant_get(data, '$.tier', 'string') = 'gold' +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_split +#### A masked pattern was here #### +POSTHOOK: query: SELECT min(variant_row_position_split.ROW__POSITION) AS first_gold_position, + max(variant_row_position_split.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position_split +WHERE variant_get(data, '$.tier', 'string') = 'gold' +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_split +#### A masked pattern was here #### +200 399 200 +PREHOOK: query: drop table variant_row_position_split +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_row_position_split +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_split +POSTHOOK: query: drop table variant_row_position_split +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_row_position_split +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_split +PREHOOK: query: drop table if exists variant_row_position_del +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_row_position_del +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_row_position_del ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'write.delete.mode'='merge-on-read' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: CREATE EXTERNAL TABLE variant_row_position_del ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'write.delete.mode'='merge-on-read' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_del +PREHOOK: query: INSERT INTO variant_row_position_del +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: INSERT INTO variant_row_position_del +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_row_position_del +PREHOOK: query: DELETE FROM variant_row_position_del +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_del +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: DELETE FROM variant_row_position_del +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_del +POSTHOOK: Output: default@variant_row_position_del +PREHOOK: query: SELECT count(*) AS rows_left, min(id) AS lowest_id, max(id) AS highest_id FROM variant_row_position_del +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_del +#### A masked pattern was here #### +POSTHOOK: query: SELECT count(*) AS rows_left, min(id) AS lowest_id, max(id) AS highest_id FROM variant_row_position_del +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_del +#### A masked pattern was here #### +395 0 399 +PREHOOK: query: SELECT id FROM variant_row_position_del WHERE id BETWEEN 197 AND 208 ORDER BY id +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_del +#### A masked pattern was here #### +POSTHOOK: query: SELECT id FROM variant_row_position_del WHERE id BETWEEN 197 AND 208 ORDER BY id +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_del +#### A masked pattern was here #### +197 +198 +199 +205 +206 +207 +208 +PREHOOK: query: drop table variant_row_position_del +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_row_position_del +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: drop table variant_row_position_del +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_row_position_del +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_del diff --git a/iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out b/iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out index 2f1fac1f4b41..209acc1b2066 100644 --- a/iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out +++ b/iceberg/iceberg-handler/src/test/results/positive/variant_type_row_position.q.out @@ -86,3 +86,162 @@ POSTHOOK: type: DROPTABLE POSTHOOK: Input: default@variant_row_position POSTHOOK: Output: database:default POSTHOOK: Output: default@variant_row_position +PREHOOK: query: drop table if exists variant_row_position_split +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_row_position_split +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_row_position_split ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'read.split.target-size'='1024' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_split +POSTHOOK: query: CREATE EXTERNAL TABLE variant_row_position_split ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'read.split.target-size'='1024' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_split +PREHOOK: query: INSERT INTO variant_row_position_split +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_row_position_split +POSTHOOK: query: INSERT INTO variant_row_position_split +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_row_position_split +PREHOOK: query: SELECT min(variant_row_position_split.ROW__POSITION) AS first_gold_position, + max(variant_row_position_split.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position_split +WHERE variant_get(data, '$.tier', 'string') = 'gold' +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_split +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT min(variant_row_position_split.ROW__POSITION) AS first_gold_position, + max(variant_row_position_split.ROW__POSITION) AS last_gold_position, + count(*) AS gold_rows +FROM variant_row_position_split +WHERE variant_get(data, '$.tier', 'string') = 'gold' +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_split +POSTHOOK: Output: hdfs://### HDFS PATH ### +200 399 200 +PREHOOK: query: drop table variant_row_position_split +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_row_position_split +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_split +POSTHOOK: query: drop table variant_row_position_split +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_row_position_split +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_split +PREHOOK: query: drop table if exists variant_row_position_del +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_row_position_del +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_row_position_del ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'write.delete.mode'='merge-on-read' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: CREATE EXTERNAL TABLE variant_row_position_del ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3', + 'variant.shredding.enabled'='true', + 'write.parquet.row-group-size-bytes'='1024', + 'write.delete.mode'='merge-on-read' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_del +PREHOOK: query: INSERT INTO variant_row_position_del +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: INSERT INTO variant_row_position_del +SELECT pos, parse_json(concat('{"tier": "', if(pos < 200, 'bronze', 'gold'), '", "n": ', pos, '}')) +FROM (SELECT 1) x LATERAL VIEW posexplode(split(space(399), ' ')) e AS pos, val +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_row_position_del +PREHOOK: query: DELETE FROM variant_row_position_del +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_del +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: DELETE FROM variant_row_position_del +WHERE variant_get(data, '$.tier', 'string') = 'gold' AND id < 205 +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_del +POSTHOOK: Output: default@variant_row_position_del +PREHOOK: query: SELECT count(*) AS rows_left, min(id) AS lowest_id, max(id) AS highest_id FROM variant_row_position_del +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_del +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT count(*) AS rows_left, min(id) AS lowest_id, max(id) AS highest_id FROM variant_row_position_del +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_del +POSTHOOK: Output: hdfs://### HDFS PATH ### +395 0 399 +PREHOOK: query: SELECT id FROM variant_row_position_del WHERE id BETWEEN 197 AND 208 ORDER BY id +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_row_position_del +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT id FROM variant_row_position_del WHERE id BETWEEN 197 AND 208 ORDER BY id +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_row_position_del +POSTHOOK: Output: hdfs://### HDFS PATH ### +197 +198 +199 +205 +206 +207 +208 +PREHOOK: query: drop table variant_row_position_del +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_row_position_del +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_row_position_del +POSTHOOK: query: drop table variant_row_position_del +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_row_position_del +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_row_position_del diff --git a/itests/src/test/resources/testconfiguration.properties b/itests/src/test/resources/testconfiguration.properties index 5b077e8f5bf5..4afdfbccd246 100644 --- a/itests/src/test/resources/testconfiguration.properties +++ b/itests/src/test/resources/testconfiguration.properties @@ -413,6 +413,7 @@ iceberg.llap.query.files=\ llap_iceberg_read_orc.q,\ llap_iceberg_read_parquet.q,\ puffin_col_stats_with_time_travel.q,\ + variant_type_row_position.q,\ vectorized_iceberg_read_mixed.q,\ vectorized_iceberg_read_multitable.q,\ vectorized_iceberg_read_orc.q,\ diff --git a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java index a24ed83e8489..c22354898dfb 100644 --- a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java +++ b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/VectorizedParquetInputFormat.java @@ -61,16 +61,18 @@ public RecordReader getRecordReader( inputSplit, jobConf, metadataCache, dataCache, cacheConf, metadata, initialDefaults, includedRowGroups); } - public void setMetadata(ParquetMetadata metadata) throws IOException { - this.metadata = metadata; - } - /** - * Restricts the reader to these of the file's row groups, by footer index, for a caller that has already - * ruled some out with a filter Parquet cannot express. The footer given to {@link #setMetadata} stays the - * file's own, so row positions remain absolute. + * The footer to read the file by, and which of its row groups to read: one flag each by footer index, or + * null for all of them. A caller passes a subset when it has ruled row groups out with a filter Parquet + * cannot express. The footer stays the file's own either way, so row positions remain absolute. */ - public void setIncludedRowGroups(boolean[] includedRowGroups) { + public void setMetadata(ParquetMetadata metadata, boolean[] includedRowGroups) throws IOException { + if (metadata != null && includedRowGroups != null + && includedRowGroups.length != metadata.getBlocks().size()) { + throw new IOException("Got " + includedRowGroups.length + " row group flags for a footer of " + + metadata.getBlocks().size() + " row groups"); + } + this.metadata = metadata; this.includedRowGroups = includedRowGroups; } diff --git a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java index c4d0237cf7be..01a0ce3c7c2f 100644 --- a/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java +++ b/ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java @@ -260,8 +260,14 @@ public void initialize( for (BlockMetaData block : parquetMetadata.getBlocks()) { if (offsets.contains(block.getStartingPos())) { // Parquet records where each row group's first row sits in the file, so the position is read - // rather than counted up, and stays right whatever subset of row groups this split reads. - rowGroupNumToRowPos.put(blockIndex++, block.getRowIndexOffset()); + // rather than counted up, and stays right whatever subset of row groups this split reads. The + // footer readers all fill it in; -1 is Parquet's marker for a footer that did not. + long rowIndexOffset = block.getRowIndexOffset(); + if (rowIndexOffset < 0) { + throw new IOException("Parquet footer for " + filePath + " has no row index offset for the row " + + "group at " + block.getStartingPos() + "; cannot determine row positions"); + } + rowGroupNumToRowPos.put(blockIndex++, rowIndexOffset); blocks.add(block); } } diff --git a/ql/src/test/org/apache/hadoop/hive/ql/io/parquet/TestVectorizedColumnReader.java b/ql/src/test/org/apache/hadoop/hive/ql/io/parquet/TestVectorizedColumnReader.java index e1ed9f1bbb08..d29ed9270c27 100644 --- a/ql/src/test/org/apache/hadoop/hive/ql/io/parquet/TestVectorizedColumnReader.java +++ b/ql/src/test/org/apache/hadoop/hive/ql/io/parquet/TestVectorizedColumnReader.java @@ -26,7 +26,10 @@ import org.apache.hadoop.hive.serde2.ColumnProjectionUtils; import org.apache.hadoop.mapred.FileSplit; import org.apache.hadoop.mapred.JobConf; +import org.apache.hadoop.mapred.Reporter; import org.apache.hadoop.mapreduce.Job; +import org.apache.parquet.format.converter.ParquetMetadataConverter; +import org.apache.parquet.hadoop.ParquetFileReader; import org.apache.parquet.hadoop.ParquetInputFormat; import org.apache.parquet.hadoop.ParquetInputSplit; import org.apache.parquet.hadoop.metadata.ParquetMetadata; @@ -37,6 +40,7 @@ import org.junit.Test; import java.io.IOException; +import java.util.Arrays; import static org.apache.parquet.hadoop.api.ReadSupport.PARQUET_READ_SCHEMA; @@ -185,4 +189,52 @@ public void testNullSplitForParquetReader() throws Exception { TestVectorizedParquetRecordReader testReader = new TestVectorizedParquetRecordReader(fsplit, jobConf); Assert.assertNull("Test should return null split from getSplit() method", testReader.getSplit(null)); } + + /** + * A caller that has already ruled row groups out names the ones left, and the reader reads those and no + * others. Nothing else in the suite would notice if the picks were ignored: the rows returned and their + * positions are the same either way, so only the blocks the reader settled on can show it. + */ + @Test + public void testReaderReadsOnlyThePickedRowGroups() throws Exception { + Configuration conf = newSingleColumnConf(); + Job vectorJob = new Job(conf, "read vector"); + ParquetInputFormat.setInputPaths(vectorJob, file); + initialVectorizedRowBatchCtx(conf); + FileSplit fsplit = getFileSplit(vectorJob); + JobConf jobConf = new JobConf(conf); + ParquetMetadata footer = ParquetFileReader.readFooter(jobConf, file, ParquetMetadataConverter.NO_FILTER); + int rowGroups = footer.getBlocks().size(); + + boolean[] everyRowGroup = new boolean[rowGroups]; + Arrays.fill(everyRowGroup, true); + Assert.assertEquals("Picking every row group should read the whole split", + rowGroups, readerOver(footer, everyRowGroup, fsplit, jobConf).getFilteredBlocks().size()); + + Assert.assertNull("Picking no row group should leave the reader nothing to read", + readerOver(footer, new boolean[rowGroups], fsplit, jobConf).getFilteredBlocks()); + + VectorizedParquetInputFormat inputFormat = new VectorizedParquetInputFormat(); + Assert.assertThrows("Picks belonging to another footer should be refused, not applied by index", + IOException.class, () -> inputFormat.setMetadata(footer, new boolean[rowGroups + 1])); + } + + private static VectorizedParquetRecordReader readerOver(ParquetMetadata footer, boolean[] includedRowGroups, + FileSplit fsplit, JobConf jobConf) throws IOException { + VectorizedParquetInputFormat inputFormat = new VectorizedParquetInputFormat(); + inputFormat.setMetadata(footer, includedRowGroups); + return (VectorizedParquetRecordReader) inputFormat.getRecordReader(fsplit, jobConf, Reporter.NULL); + } + + private static Configuration newSingleColumnConf() { + Configuration conf = new Configuration(); + conf.set(IOConstants.COLUMNS, "int32_field"); + conf.set(IOConstants.COLUMNS_TYPES, "int"); + conf.setBoolean(ColumnProjectionUtils.READ_ALL_COLUMNS, false); + conf.set(ColumnProjectionUtils.READ_COLUMN_IDS_CONF_STR, "0"); + conf.set(PARQUET_READ_SCHEMA, "message test { required int32 int32_field;}"); + HiveConf.setBoolVar(conf, HiveConf.ConfVars.HIVE_VECTORIZATION_ENABLED, true); + HiveConf.setVar(conf, HiveConf.ConfVars.PLAN, "//tmp"); + return conf; + } }