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]

Reply via email to