deniskuzZ commented on code in PR #6793: URL: https://github.com/apache/hive/pull/6793#discussion_r4156953624
########## ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/ParquetRowGroupDecoder.java: ########## @@ -0,0 +1,273 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hadoop.hive.ql.io.parquet.vector; + +import java.io.IOException; +import java.time.ZoneId; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.Optional; + +import org.apache.hadoop.hive.serde2.typeinfo.ListTypeInfo; +import org.apache.hadoop.hive.serde2.typeinfo.PrimitiveTypeInfo; +import org.apache.hadoop.hive.serde2.typeinfo.StructTypeInfo; +import org.apache.hadoop.hive.serde2.typeinfo.TypeInfo; +import org.apache.parquet.ParquetRuntimeException; +import org.apache.parquet.column.ColumnDescriptor; +import org.apache.parquet.column.page.PageReadStore; +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.InvalidSchemaException; +import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.Type; + +/** + * Reusable, stateless helper that builds the {@link VectorizedColumnReader} array for one Parquet + * row group, given an already-read {@link PageReadStore}. The logic here was lifted verbatim from + * {@link VectorizedParquetRecordReader} so that both the row-by-split reader and the LLAP + * cache-backed consumer ({@code ParquetEncodedDataConsumer}) can share a single, behavior-preserving + * implementation of the Hive-type-driven Parquet column-reader construction. + * + * <p>The instance only carries the {@code fileSchema} (used for the schema-evolution check). All + * per-row-group state lives in the supplied {@link PageReadStore}. + */ +public class ParquetRowGroupDecoder { + + private static final int MAP_DEFINITION_LEVEL_MAX = 3; + + /** + * Bag of writer-timezone / proleptic / legacy conversion flags forwarded to every primitive + * column reader. Grouped so the reader-construction entry point stays under the parameter cap. + */ + public record TimestampConversionOptions( + boolean skipTimestampConversion, + ZoneId writerTimezone, + boolean skipProlepticConversion, + boolean legacyConversionEnabled) { + } + + private final MessageType fileSchema; + private final Map<String, Object> initialDefaults; + + public ParquetRowGroupDecoder(MessageType fileSchema, Map<String, Object> initialDefaults) { + this.fileSchema = fileSchema; + this.initialDefaults = initialDefaults; + } + + /** + * Builds the per-(requested-)column {@link VectorizedColumnReader} array for a row group. + * + * @param pages the row group's page store (from {@code reader.readRowGroup(..)} + * / {@code readNextRowGroup()}) + * @param requestedSchema the projected Parquet schema being read + * @param columnTypesList the Hive type infos for ALL table columns (indexed by table col id) + * @param colsToInclude the table column ids being read, in requested-schema field order; + * may be empty (e.g. {@code count(*)}), in which case all readers are null + * @param readAllColumns whether projection is "read all columns" + * @param options writer-timezone / proleptic / legacy conversion flags, see + * {@link TimestampConversionOptions} + * @return one reader per requested-schema field (null entries where no reader is needed) + */ + public VectorizedColumnReader[] buildColumnReaders( + PageReadStore pages, + MessageType requestedSchema, + List<TypeInfo> columnTypesList, + List<Integer> colsToInclude, + boolean readAllColumns, + TimestampConversionOptions options) throws IOException { + List<Type> types = requestedSchema.getFields(); + VectorizedColumnReader[] columnReaders = new VectorizedColumnReader[types.size()]; + + if (!readAllColumns) { + // certain queries like select count(*) from table do not have + // any projected columns and still have isReadAllColumns as false + // in such cases columnReaders are not needed + // However, if colsToInclude is not empty we should initialize each columnReader + if (!colsToInclude.isEmpty()) { + for (int i = 0; i < types.size(); ++i) { + columnReaders[i] = buildVectorizedParquetReader( + columnTypesList.get(colsToInclude.get(i)), types.get(i), pages, + requestedSchema.getColumns(), options, 0, 0); + } + } + } else { + for (int i = 0; i < types.size(); ++i) { + columnReaders[i] = buildVectorizedParquetReader(columnTypesList.get(i), + types.get(i), pages, requestedSchema.getColumns(), options, 0, 0); + } + } + return columnReaders; + } + + private static List<ColumnDescriptor> getAllColumnDescriptorByType( + int depth, + Type type, + List<ColumnDescriptor> columns) throws ParquetRuntimeException { + List<ColumnDescriptor> res = new ArrayList<>(); + for (ColumnDescriptor descriptor : columns) { + if (depth >= descriptor.getPath().length) { + throw new InvalidSchemaException("Corrupted Parquet schema"); + } + if (type.getName().equals(descriptor.getPath()[depth])) { + res.add(descriptor); + } + } + return res; + } + + // Nested types are unsupported here on purpose; callers detect that in advance + // (ParquetEncodedDataReader.projectsNestedTypes) and route the query to the non-native reader. + private static PrimitiveType getElementType(Type type) { + if (type.isPrimitive()) { + return type.asPrimitiveType(); + } + if (type.asGroupType().getFields().size() > 1) { + throw new UnsupportedOperationException( + "Current Parquet Vectorization reader doesn't support nested type"); + } + + Type childType = type.asGroupType().getFields().getFirst(); + + // Parquet file generated using thrift may have child type as PrimitiveType + if (childType.isPrimitive()) { + return childType.asPrimitiveType(); + } else { + return childType.asGroupType().getFields().getFirst().asPrimitiveType(); + } + } + + // Build VectorizedParquetColumnReader via Hive typeInfo and Parquet schema + private VectorizedColumnReader buildVectorizedParquetReader( + TypeInfo typeInfo, + Type type, + PageReadStore pages, + List<ColumnDescriptor> columnDescriptors, + TimestampConversionOptions options, + int depth, int currentDefLevel) throws IOException { + int typeDefLevel = currentDefLevel; + if (type.isRepetition(Type.Repetition.OPTIONAL) || type.isRepetition(Type.Repetition.REPEATED)) { + typeDefLevel++; + } + List<ColumnDescriptor> descriptors = + getAllColumnDescriptorByType(depth, type, columnDescriptors); + // Support for schema evolution: if the column from the current + // query schema is not present in the file schema, return a dummy + // reader that produces nulls. This allows queries to proceed even + // when new columns have been added after the file was written. + if (!fileSchema.getColumns().contains(descriptors.getFirst())) { + return new VectorizedDummyColumnReader(Optional.ofNullable(initialDefaults) + .map(defaults -> defaults.getOrDefault(descriptors.getFirst().getPath()[0], null)).orElse(null)); + } + switch (typeInfo.getCategory()) { + case PRIMITIVE: Review Comment: is formatting good? -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
