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..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 @@ -266,9 +266,10 @@ 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 alongside it; null means read them all. + inputFormat.setMetadata(parquetMetadata, + VariantParquetFilters.pickRowGroups(fileSchema, residual, parquetMetadata.getBlocks())); MessageType typeWithIds = null; Schema expectedSchema = task.spec().schema(); @@ -287,7 +288,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..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; @@ -31,6 +32,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 +54,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 "; @@ -137,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(); @@ -154,6 +157,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..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 @@ -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,11 @@ private static ResolvedVariantFilter resolveVariantFilter(MessageType schema, Ex return new ResolvedVariantFilter(predicate, visitor.fallbackValueColumns()); } - public static boolean[] variantRowGroupMayMatch( + /** + * 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()) { return null; @@ -135,54 +138,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..6b640e7e5234 --- /dev/null +++ b/iceberg/iceberg-handler/src/test/queries/positive/variant_type_row_position.q @@ -0,0 +1,95 @@ +-- 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; + +-- 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 new file mode 100644 index 000000000000..209acc1b2066 --- /dev/null +++ b/iceberg/iceberg-handler/src/test/results/positive/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 +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 +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/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..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 @@ -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,11 +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 { + /** + * 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 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; } public void seInitialColumnDefaults(Map 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..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 @@ -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,20 @@ 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. 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); } - allRowsInFile += block.getRowCount(); } // verify we found them all if (blocks.size() != rowGroupOffsets.length) { 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; + } }