github-actions[bot] commented on code in PR #65674:
URL: https://github.com/apache/doris/pull/65674#discussion_r3607216546


##########
be/src/format_v2/parquet/parquet_statistics.cpp:
##########
@@ -1550,4 +1632,238 @@ Status select_row_group_ranges_by_page_index(
     return Status::OK();
 }
 
+namespace {
+
+template <typename ValueType>
+bool set_native_page_scalar_min_max(const tparquet::ColumnIndex& column_index,
+                                    const ParquetColumnSchema& column_schema, 
size_t page_idx,
+                                    DecodedValueKind kind, 
ParquetColumnStatistics* page_statistics,
+                                    const cctz::time_zone* timezone) {
+    if (page_idx >= column_index.min_values.size() || page_idx >= 
column_index.max_values.size() ||
+        column_index.min_values[page_idx].size() != sizeof(ValueType) ||
+        column_index.max_values[page_idx].size() != sizeof(ValueType)) {
+        return false;
+    }
+    const auto min_value = 
unaligned_load<ValueType>(column_index.min_values[page_idx].data());
+    const auto max_value = 
unaligned_load<ValueType>(column_index.max_values[page_idx].data());
+    if constexpr (std::is_same_v<ValueType, int64_t>) {
+        if (!timestamp_min_max_is_safe(column_schema, min_value, max_value, 
timezone)) {
+            return false;
+        }
+    }
+    if (!valid_min_max(min_value, max_value)) {
+        return true;
+    }
+    if (!set_decoded_field(column_schema, kind, min_value, 
&page_statistics->min_value, timezone) ||
+        !set_decoded_field(column_schema, kind, max_value, 
&page_statistics->max_value, timezone)) {
+        return false;
+    }
+    if (decoded_min_max_is_ordered(*page_statistics)) {
+        page_statistics->has_min_max = true;
+    }
+    return true;
+}
+
+bool build_native_page_statistics(const tparquet::ColumnIndex& column_index,
+                                  const ParquetColumnSchema& column_schema, 
size_t page_idx,
+                                  ParquetColumnStatistics* page_statistics,
+                                  const cctz::time_zone* timezone) {
+    DORIS_CHECK(page_statistics != nullptr);
+    *page_statistics = {};
+    if (!column_index.__isset.null_counts || page_idx >= 
column_index.null_pages.size() ||
+        page_idx >= column_index.null_counts.size()) {
+        return false;
+    }
+    page_statistics->has_null_count = true;
+    page_statistics->has_null = column_index.null_counts[page_idx] > 0;

Review Comment:
   [P1] Reject contradictory ColumnIndex null metadata before pruning
   
   The native loader checks the `null_pages` vector length but not these counts 
or their consistency. For a positive-row page with `null_pages[i] = true` and 
`null_counts[i] = 0` (or a negative count), this produces a zone map with both 
`has_null` and `has_not_null` false. `IS NULL` then evaluates to `kNoMatch`, so 
an all-null page can be removed without evaluating its data. Please 
invalidate/fall back from the optional ColumnIndex on negative or contradictory 
null metadata, and add native PageIndex tests for `IS NULL`/`IS NOT NULL` with 
those cases.



##########
be/src/format_v2/parquet/parquet_file_context.cpp:
##########
@@ -568,6 +753,219 @@ Status ParquetFileContext::open(io::FileReaderSPtr 
input_file_reader, io::IOCont
     return Status::OK();
 }
 
+Status ParquetFileContext::load_native_offset_indexes(
+        int row_group_id, const std::unordered_set<int>& leaf_column_ids,
+        std::unordered_map<int, tparquet::OffsetIndex>* offset_indexes) const {
+    DORIS_CHECK(offset_indexes != nullptr);
+    offset_indexes->clear();
+    if (leaf_column_ids.empty()) {
+        return Status::OK();
+    }
+    const auto& thrift_metadata = native_metadata->to_thrift();
+    if (row_group_id < 0 || row_group_id >= 
static_cast<int>(thrift_metadata.row_groups.size())) {
+        return Status::Corruption("Invalid Parquet row group {} for 
OffsetIndex", row_group_id);
+    }
+    const auto& native_row_group = thrift_metadata.row_groups[row_group_id];
+    const auto compat = native::parquet_reader_compat(
+            thrift_metadata.__isset.created_by ? thrift_metadata.created_by : 
"");
+    try {
+        for (const int leaf_column_id : leaf_column_ids) {
+            if (leaf_column_id < 0 ||
+                leaf_column_id >= 
static_cast<int>(native_row_group.columns.size())) {
+                return Status::Corruption("Invalid Parquet leaf {} for 
OffsetIndex",
+                                          leaf_column_id);
+            }
+            const auto& column_chunk = 
native_row_group.columns[leaf_column_id];
+            if (!column_chunk.__isset.offset_index_offset ||
+                !column_chunk.__isset.offset_index_length ||
+                column_chunk.offset_index_length <= 0) {
+                continue;
+            }
+            const int64_t index_offset = column_chunk.offset_index_offset;
+            const int64_t index_length = column_chunk.offset_index_length;
+            if (index_offset < 0 || index_length <= 0 || index_offset > 
native_file->size() ||
+                index_length > native_file->size() - index_offset ||
+                index_length > std::numeric_limits<uint32_t>::max()) {
+                // OffsetIndex is optional. A malformed range must not 
allocate from untrusted
+                // footer values or redirect the native reader outside the 
file.
+                continue;
+            }
+            std::vector<uint8_t> 
serialized_index(static_cast<size_t>(index_length));
+            Slice index_slice(serialized_index.data(), 
serialized_index.size());
+            size_t bytes_read = 0;
+            if (!native_file->read_at(index_offset, index_slice, &bytes_read, 
native_io_ctx).ok() ||
+                bytes_read != serialized_index.size()) {
+                continue;
+            }
+            uint32_t thrift_length = 
static_cast<uint32_t>(serialized_index.size());
+            tparquet::OffsetIndex native_index;
+            if (!deserialize_thrift_msg(serialized_index.data(), 
&thrift_length, true,
+                                        &native_index)
+                         .ok() ||
+                native_index.page_locations.empty()) {
+                continue;
+            }
+            native::ColumnChunkRange chunk_range;
+            RETURN_IF_ERROR(native::compute_column_chunk_range(
+                    native_row_group.columns[leaf_column_id].meta_data, 
native_file->size(),
+                    compat.parquet_816_padding, &chunk_range));
+            if (!native::validate_offset_index(
+                        native_index, chunk_range,
+                        
native_row_group.columns[leaf_column_id].meta_data.data_page_offset,
+                        native_row_group.num_rows)) {
+                // OffsetIndex is optional. Reject the complete index instead 
of letting one bad
+                // location redirect an indexed reader outside its owning 
column chunk.
+                continue;
+            }
+            offset_indexes->emplace(leaf_column_id, std::move(native_index));
+        }
+    } catch (const std::exception&) {
+        // OffsetIndex is optional. Selected logical ranges still enforce 
correctness, while the
+        // native reader conservatively falls back to sequential page 
traversal.
+        offset_indexes->clear();
+    }
+    return Status::OK();
+}
+
+Status ParquetFileContext::load_native_page_indexes(
+        int row_group_id, const std::unordered_set<int>& leaf_column_ids,
+        std::unordered_map<int, NativeParquetPageIndex>* page_indexes, 
int64_t* read_time,
+        int64_t* parse_time) const {
+    DORIS_CHECK(page_indexes != nullptr);
+    page_indexes->clear();
+    if (leaf_column_ids.empty()) {
+        return Status::OK();
+    }
+    const auto& thrift_metadata = native_metadata->to_thrift();
+    if (row_group_id < 0 || row_group_id >= 
static_cast<int>(thrift_metadata.row_groups.size())) {
+        return Status::Corruption("Invalid Parquet row group {} for 
PageIndex", row_group_id);
+    }
+    const auto& row_group = thrift_metadata.row_groups[row_group_id];
+    const auto compat = native::parquet_reader_compat(
+            thrift_metadata.__isset.created_by ? thrift_metadata.created_by : 
"");
+
+    struct SerializedIndexRange {
+        int leaf_column_id;
+        int64_t offset;
+        int64_t length;
+    };
+    struct PendingPageIndex {
+        NativeParquetPageIndex indexes;
+        bool has_column_index = false;
+        bool has_offset_index = false;
+    };
+    std::vector<SerializedIndexRange> column_index_ranges;
+    std::vector<SerializedIndexRange> offset_index_ranges;
+    std::unordered_map<int, PendingPageIndex> pending_indexes;
+
+    auto valid_index_range = [&](int64_t offset, int64_t length) {
+        if (offset < 0 || length <= 0 || offset > native_file->size() ||
+            length > native_file->size() - offset ||
+            length > std::numeric_limits<uint32_t>::max()) {
+            return false;
+        }
+        return true;
+    };
+
+    for (const int leaf_column_id : leaf_column_ids) {
+        if (leaf_column_id < 0 || leaf_column_id >= 
static_cast<int>(row_group.columns.size())) {
+            return Status::Corruption("Invalid Parquet leaf {} for PageIndex", 
leaf_column_id);
+        }
+        const auto& chunk = row_group.columns[leaf_column_id];
+        if (!chunk.__isset.column_index_offset || 
!chunk.__isset.column_index_length ||
+            !chunk.__isset.offset_index_offset || 
!chunk.__isset.offset_index_length) {
+            continue;
+        }
+        if (!valid_index_range(chunk.column_index_offset, 
chunk.column_index_length) ||
+            !valid_index_range(chunk.offset_index_offset, 
chunk.offset_index_length)) {
+            continue;
+        }
+        column_index_ranges.push_back(
+                {leaf_column_id, chunk.column_index_offset, 
chunk.column_index_length});
+        offset_index_ranges.push_back(
+                {leaf_column_id, chunk.offset_index_offset, 
chunk.offset_index_length});
+        pending_indexes.try_emplace(leaf_column_id);
+    }
+
+    auto read_coalesced_indexes = [&](std::vector<SerializedIndexRange>* 
ranges,
+                                      bool column_index) {
+        std::sort(ranges->begin(), ranges->end(),
+                  [](const auto& lhs, const auto& rhs) { return lhs.offset < 
rhs.offset; });
+        size_t range_begin = 0;
+        while (range_begin < ranges->size()) {
+            size_t range_end = range_begin + 1;
+            int64_t span_end = (*ranges)[range_begin].offset + 
(*ranges)[range_begin].length;
+            while (range_end < ranges->size() && (*ranges)[range_end].offset 
<= span_end) {
+                span_end = std::max(span_end,
+                                    (*ranges)[range_end].offset + 
(*ranges)[range_end].length);
+                ++range_end;
+            }
+
+            const int64_t span_offset = (*ranges)[range_begin].offset;
+            const int64_t span_length = span_end - span_offset;
+            std::vector<uint8_t> serialized(static_cast<size_t>(span_length));

Review Comment:
   [P1] Bound optional index reads before allocating them
   
   Each individual footer range is accepted up to the integer limit when it 
merely fits inside the file, and adjacent ranges across projected leaves are 
merged here without an aggregate budget. A large object can therefore advertise 
one roughly-2-GiB index (or several adjacent ones) and make this optional 
pruning path allocate a multi-gigabyte `std::vector` before Thrift parsing can 
reject it; the OffsetIndex-only loader has the same whole-range allocation at 
line 793. Please cap individual and coalesced serialized-index bytes (using 
tracked storage where applicable) and conservatively skip the optional index 
above that budget, with large-length/adjacent-range fallback tests.



##########
be/src/format_v2/parquet/reader/native/column_chunk_reader.cpp:
##########
@@ -0,0 +1,1159 @@
+// 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.
+
+#include "format_v2/parquet/reader/native/column_chunk_reader.h"
+
+#include <gen_cpp/parquet_types.h>
+#include <glog/logging.h>
+#include <parquet/metadata.h>
+#include <string.h>
+
+#include <algorithm>
+#include <cstdint>
+#include <limits>
+#include <memory>
+#include <utility>
+
+#include "common/compiler_util.h" // IWYU pragma: keep
+#include "core/column/column.h"
+#include "core/custom_allocator.h"
+#include "core/data_type_serde/data_type_serde.h"
+#include "format/parquet/schema_desc.h"
+#include "format_v2/parquet/reader/native/decoder.h"
+#include "format_v2/parquet/reader/native/level_decoder.h"
+#include "format_v2/parquet/reader/native/page_reader.h"
+#include "io/fs/buffered_reader.h"
+#include "runtime/runtime_profile.h"
+#include "storage/cache/page_cache.h"
+#include "util/bit_util.h"
+#include "util/block_compression.h"
+
+namespace cctz {
+class time_zone;
+} // namespace cctz
+namespace doris {
+namespace io {
+class BufferedStreamReader;
+struct IOContext;
+} // namespace io
+} // namespace doris
+
+namespace doris::format::parquet::native {
+
+ParquetReaderCompat parquet_reader_compat(const std::string& created_by) {
+    if (created_by.empty()) {
+        return {};
+    }
+    const ::parquet::ApplicationVersion version(created_by);
+    return {.parquet_816_padding =
+                    
version.VersionLt(::parquet::ApplicationVersion::PARQUET_816_FIXED_VERSION()),
+            .data_page_v2_always_compressed = version.VersionLt(
+                    
::parquet::ApplicationVersion::PARQUET_CPP_10353_FIXED_VERSION())};
+}
+
+Status compute_column_chunk_range(const tparquet::ColumnMetaData& metadata, 
size_t file_size,
+                                  bool parquet_816_padding, ColumnChunkRange* 
range) {
+    DORIS_CHECK(range != nullptr);
+    int64_t start = metadata.data_page_offset;
+    if (metadata.__isset.dictionary_page_offset && 
metadata.dictionary_page_offset >= 0 &&
+        metadata.dictionary_page_offset < start) {
+        start = metadata.dictionary_page_offset;
+    }
+    const int64_t length = metadata.total_compressed_size;
+    if (UNLIKELY(start < 0 || length < 0)) {
+        return Status::Corruption("Parquet column chunk has a negative offset 
or length");
+    }
+    const uint64_t unsigned_start = static_cast<uint64_t>(start);
+    const uint64_t unsigned_length = static_cast<uint64_t>(length);
+    if (UNLIKELY(unsigned_start > file_size || unsigned_length > file_size - 
unsigned_start)) {
+        // Thrift range fields are signed and untrusted; validate before 
converting them to the
+        // unsigned stream-reader coordinates so overflow cannot wrap back 
into the file.
+        return Status::Corruption("Parquet column chunk [{}, {}) exceeds file 
size {}", start,
+                                  unsigned_start + unsigned_length, file_size);
+    }
+    size_t bounded_length = static_cast<size_t>(unsigned_length);
+    if (parquet_816_padding) {
+        // parquet-mr before PARQUET-816 under-reported the chunk by up to 100 
bytes. Padding stays
+        // file-bounded and is only enabled for the affected writer versions.
+        bounded_length += std::min<size_t>(100, file_size - unsigned_start - 
unsigned_length);
+    }
+    range->offset = static_cast<size_t>(unsigned_start);
+    range->length = bounded_length;
+    return Status::OK();
+}
+
+bool validate_offset_index(const tparquet::OffsetIndex& index, const 
ColumnChunkRange& chunk_range,
+                           int64_t data_page_offset, int64_t row_count) {
+    if (index.page_locations.empty() || data_page_offset < 0 || row_count < 0 
||
+        index.page_locations.front().first_row_index != 0 ||
+        index.page_locations.front().offset != data_page_offset ||
+        chunk_range.length > std::numeric_limits<size_t>::max() - 
chunk_range.offset) {
+        return false;
+    }
+    // Row indexes alone cannot detect a uniformly shifted OffsetIndex. Anchor 
its first location
+    // to the owning metadata so page-to-row mapping cannot silently move by 
one physical page.
+    const uint64_t chunk_begin = chunk_range.offset;
+    const uint64_t chunk_end = chunk_begin + chunk_range.length;
+    uint64_t previous_end = chunk_begin;
+    int64_t previous_row = -1;
+    for (const auto& location : index.page_locations) {
+        if (location.first_row_index <= previous_row || 
location.first_row_index >= row_count ||
+            location.offset < 0 || location.compressed_page_size <= 0) {
+            return false;
+        }
+        const uint64_t begin = static_cast<uint64_t>(location.offset);
+        const uint64_t size = 
static_cast<uint64_t>(location.compressed_page_size);
+        if (begin < chunk_begin || begin < previous_end || begin > chunk_end ||
+            size > chunk_end - begin) {
+            return false;
+        }
+        previous_row = location.first_row_index;
+        previous_end = begin + size;
+    }
+    return true;
+}
+
+namespace {
+
+Status translate_value_encoding(tparquet::Encoding::type encoding,
+                                ParquetValueEncoding* translated) {
+    DORIS_CHECK(translated != nullptr);
+    switch (encoding) {
+    case tparquet::Encoding::PLAIN:
+        *translated = ParquetValueEncoding::PLAIN;
+        return Status::OK();
+    case tparquet::Encoding::RLE_DICTIONARY:
+    case tparquet::Encoding::PLAIN_DICTIONARY:
+        *translated = ParquetValueEncoding::DICTIONARY;
+        return Status::OK();
+    case tparquet::Encoding::RLE:
+        *translated = ParquetValueEncoding::RLE;
+        return Status::OK();
+    case tparquet::Encoding::BIT_PACKED:
+        *translated = ParquetValueEncoding::BIT_PACKED;
+        return Status::OK();
+    case tparquet::Encoding::DELTA_BINARY_PACKED:
+        *translated = ParquetValueEncoding::DELTA_BINARY_PACKED;
+        return Status::OK();
+    case tparquet::Encoding::DELTA_LENGTH_BYTE_ARRAY:
+        *translated = ParquetValueEncoding::DELTA_LENGTH_BYTE_ARRAY;
+        return Status::OK();
+    case tparquet::Encoding::DELTA_BYTE_ARRAY:
+        *translated = ParquetValueEncoding::DELTA_BYTE_ARRAY;
+        return Status::OK();
+    case tparquet::Encoding::BYTE_STREAM_SPLIT:
+        *translated = ParquetValueEncoding::BYTE_STREAM_SPLIT;
+        return Status::OK();
+    default:
+        return Status::NotSupported("Unsupported Parquet encoding {}",
+                                    tparquet::to_string(encoding));
+    }
+}
+
+template <bool HAS_FILTER>
+Status decode_selected_values(IColumn& column, const DataTypeSerDe& serde, 
Decoder& decoder,
+                              const ParquetDecodeContext& context,
+                              ParquetMaterializationState& state, 
ColumnSelectVector& select_vector,
+                              int64_t* materialization_time) {
+    SCOPED_RAW_TIMER(materialization_time);
+    ColumnSelectVector::DataReadType read_type;
+    while (const size_t run_length = 
select_vector.get_next_run<HAS_FILTER>(&read_type)) {
+        switch (read_type) {
+        case ColumnSelectVector::CONTENT:
+            RETURN_IF_ERROR(
+                    serde.read_column_from_parquet(column, decoder, context, 
run_length, state));
+            break;
+        case ColumnSelectVector::NULL_DATA:
+            column.insert_many_defaults(run_length);
+            break;
+        case ColumnSelectVector::FILTERED_CONTENT:
+            RETURN_IF_ERROR(decoder.skip_values(run_length));
+            break;
+        case ColumnSelectVector::FILTERED_NULL:
+            break;
+        }
+    }
+    return Status::OK();
+}
+
+// Presents one sparse page request as an ordinary sequential source to 
DataTypeSerDe. SerDe is
+// entered once per page fragment; the concrete decoder decides whether to 
gather selected spans,
+// batch-decode and compact, or use the cursor-preserving range fallback.
+class SelectedDecodeSource final : public ParquetDecodeSource {
+public:
+    SelectedDecodeSource(Decoder& decoder, const ParquetSelection& selection)
+            : _decoder(decoder), _selection(selection) {}
+
+    Status decode_fixed_values(size_t num_values, ParquetFixedValueConsumer& 
consumer) override {
+        DORIS_CHECK_EQ(num_values, _selection.selected_values);
+        return _decoder.decode_selected_fixed_values(_selection, consumer);
+    }
+
+    Status decode_binary_values(size_t num_values, ParquetBinaryValueConsumer& 
consumer) override {
+        DORIS_CHECK_EQ(num_values, _selection.selected_values);
+        return _decoder.decode_selected_binary_values(_selection, consumer);
+    }
+
+    Status skip_values(size_t num_values) override {
+        return Status::NotSupported("Selected Parquet source cannot be 
skipped, values={}",
+                                    num_values);
+    }
+
+    bool has_dictionary() const override { return _decoder.has_dictionary(); }
+    uint64_t dictionary_generation() const override { return 
_decoder.dictionary_generation(); }
+    size_t dictionary_size() const override { return 
_decoder.dictionary_size(); }
+
+    Status decode_dictionary(ParquetFixedValueConsumer& fixed_consumer,
+                             ParquetBinaryValueConsumer& binary_consumer) 
override {
+        return _decoder.decode_dictionary(fixed_consumer, binary_consumer);
+    }
+
+    Status decode_dictionary_indices(size_t num_values, std::vector<uint32_t>* 
indices) override {
+        DORIS_CHECK_EQ(num_values, _selection.selected_values);
+        return _decoder.decode_selected_dictionary_indices(_selection, 
indices);
+    }
+
+private:
+    Decoder& _decoder;
+    const ParquetSelection& _selection;
+};
+
+Status decode_selected_non_null_values(IColumn& column, const DataTypeSerDe& 
serde,
+                                       Decoder& decoder, const 
ParquetDecodeContext& context,
+                                       ParquetMaterializationState& state,
+                                       ColumnSelectVector& select_vector,
+                                       int64_t* materialization_time) {
+    auto& selection = state.selection;
+    selection.ranges.clear();
+    selection.total_values = select_vector.num_values();
+    selection.selected_values = 0;
+
+    size_t cursor = 0;
+    ColumnSelectVector::DataReadType read_type;
+    while (const size_t run_length = 
select_vector.get_next_run<true>(&read_type)) {
+        DORIS_CHECK(read_type == ColumnSelectVector::CONTENT ||
+                    read_type == ColumnSelectVector::FILTERED_CONTENT);
+        if (read_type == ColumnSelectVector::CONTENT) {
+            selection.ranges.push_back({.first = cursor, .count = run_length});
+            selection.selected_values += run_length;
+        }
+        cursor += run_length;
+    }
+    DORIS_CHECK_EQ(cursor, selection.total_values);
+    if (selection.selected_values == 0) {
+        return decoder.skip_values(selection.total_values);
+    }
+
+    SCOPED_RAW_TIMER(materialization_time);
+    SelectedDecodeSource selected_source(decoder, selection);
+    return serde.read_column_from_parquet(column, selected_source, context,
+                                          selection.selected_values, state);
+}
+
+} // namespace
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::ColumnChunkReader(
+        io::BufferedStreamReader* reader, tparquet::ColumnChunk* column_chunk,
+        FieldSchema* field_schema, const tparquet::OffsetIndex* offset_index, 
size_t total_rows,
+        io::IOContext* io_ctx, const ParquetPageReadContext& page_read_ctx,
+        const ColumnChunkRange* chunk_range)
+        : _field_schema(field_schema),
+          _max_rep_level(field_schema->repetition_level),
+          _max_def_level(field_schema->definition_level),
+          _stream_reader(reader),
+          _metadata(column_chunk->meta_data),
+          _offset_index(offset_index),
+          _total_rows(total_rows),
+          _io_ctx(io_ctx),
+          _page_read_ctx(page_read_ctx) {
+    if (chunk_range != nullptr) {
+        _chunk_range = *chunk_range;
+        _has_validated_chunk_range = true;
+    }
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::init() {
+    size_t start_offset = _has_validated_chunk_range
+                                  ? _chunk_range.offset
+                                  : (has_dict_page(_metadata) ? 
_metadata.dictionary_page_offset
+                                                              : 
_metadata.data_page_offset);
+    size_t chunk_size =
+            _has_validated_chunk_range ? _chunk_range.length : 
_metadata.total_compressed_size;
+    // create page reader
+    _page_reader = create_page_reader<IN_COLLECTION, OFFSET_INDEX>(
+            _stream_reader, _io_ctx, start_offset, chunk_size, _total_rows, 
_metadata,
+            _page_read_ctx, _offset_index);
+    // get the block compression codec
+    RETURN_IF_ERROR(get_block_compression_codec(_metadata.codec, 
&_block_compress_codec));
+    _state = INITIALIZED;
+    RETURN_IF_ERROR(_parse_first_page_header());
+    return Status::OK();
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::skip_nested_values(
+        const std::vector<level_t>& def_levels, size_t start_index) {
+    size_t no_value_cnt = 0;
+    size_t value_cnt = 0;
+
+    DORIS_CHECK(start_index <= def_levels.size());
+    for (size_t idx = start_index; idx < def_levels.size(); idx++) {
+        level_t def_level = def_levels[idx];
+        if (IN_COLLECTION && def_level < 
_field_schema->repeated_parent_def_level) {
+            no_value_cnt++;
+        } else if (def_level < _field_schema->definition_level) {
+            no_value_cnt++;
+        } else {
+            value_cnt++;
+        }
+    }
+
+    RETURN_IF_ERROR(skip_values(value_cnt, true));
+    RETURN_IF_ERROR(skip_values(no_value_cnt, false));
+    return Status::OK();
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::read_levels(
+        size_t num_values, std::vector<level_t>* rep_levels, 
std::vector<level_t>* def_levels) {
+    DORIS_CHECK(rep_levels != nullptr);
+    DORIS_CHECK(def_levels != nullptr);
+    if (_remaining_num_values < num_values || _remaining_rep_nums < num_values 
||
+        _remaining_def_nums < num_values) {
+        return Status::Corruption(
+                "Parquet level reader requested {} slots with only {}/{}/{} 
remaining", num_values,
+                _remaining_num_values, _remaining_rep_nums, 
_remaining_def_nums);
+    }
+
+    const size_t start_index = def_levels->size();
+    rep_levels->resize(rep_levels->size() + num_values, 0);
+    def_levels->resize(def_levels->size() + num_values, 0);
+    if (_max_rep_level > 0) {
+        const size_t decoded = _rep_level_decoder.get_levels(
+                rep_levels->data() + rep_levels->size() - num_values, 
num_values);
+        if (decoded != num_values) {
+            return Status::Corruption("Parquet repetition level stream ended 
after {} of {} slots",
+                                      decoded, num_values);
+        }
+    }
+    if (_max_def_level > 0) {
+        const size_t decoded = _def_level_decoder.get_levels(
+                def_levels->data() + def_levels->size() - num_values, 
num_values);
+        if (decoded != num_values) {
+            return Status::Corruption("Parquet definition level stream ended 
after {} of {} slots",
+                                      decoded, num_values);
+        }
+    }
+    _remaining_rep_nums -= num_values;
+    _remaining_def_nums -= num_values;
+    return skip_nested_values(*def_levels, start_index);
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, 
OFFSET_INDEX>::_parse_first_page_header() {
+    while (true) {
+        RETURN_IF_ERROR(_page_reader->parse_page_header());
+        const tparquet::PageHeader* header = nullptr;
+        RETURN_IF_ERROR(_page_reader->get_page_header(&header));
+        if (header->type == tparquet::PageType::DATA_PAGE ||
+            header->type == tparquet::PageType::DATA_PAGE_V2) {
+            _state = INITIALIZED;
+            return parse_page_header();
+        }
+        if (header->type != tparquet::PageType::DICTIONARY_PAGE) {
+            RETURN_IF_ERROR(_page_reader->skip_auxiliary_page());
+            _state = INITIALIZED;
+            continue;
+        }
+        // the first page maybe directory page even if 
_metadata.__isset.dictionary_page_offset == false,
+        // so we should parse the directory page in next_page()
+        RETURN_IF_ERROR(_decode_dict_page());
+        // parse the real first data page
+        RETURN_IF_ERROR(_page_reader->dict_next_page());
+        _state = INITIALIZED;
+        // A dictionary is the only non-data page with decoder state. Any 
following index or
+        // extension pages are skipped by the same pre-data loop.
+    }
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::parse_page_header() {
+    if (_state == HEADER_PARSED || _state == DATA_LOADED) {
+        return Status::OK();
+    }
+    const tparquet::PageHeader* header = nullptr;
+    while (true) {
+        RETURN_IF_ERROR(_page_reader->parse_page_header());
+        RETURN_IF_ERROR(_page_reader->get_page_header(&header));
+        if (header->type == tparquet::PageType::DATA_PAGE ||
+            header->type == tparquet::PageType::DATA_PAGE_V2) {
+            break;
+        }
+        if (header->type == tparquet::PageType::DICTIONARY_PAGE) {
+            return Status::Corruption("Parquet dictionary page appears after 
data pages");
+        }
+        RETURN_IF_ERROR(_page_reader->skip_auxiliary_page());
+    }
+    int32_t page_num_values = _page_reader->is_header_v2() ? 
header->data_page_header_v2.num_values
+                                                           : 
header->data_page_header.num_values;
+    if (page_num_values < 0 || page_num_values > _metadata.num_values ||
+        (!OFFSET_INDEX &&
+         static_cast<uint64_t>(page_num_values) >
+                 static_cast<uint64_t>(_metadata.num_values) - 
_chunk_parsed_values)) {
+        // Page counts are untrusted and feed both level decoders and scratch 
sizing. Bound each
+        // page by the column metadata before converting to unsigned counters.
+        return Status::Corruption("Parquet data page value count {} exceeds 
column total {}",
+                                  page_num_values, _metadata.num_values);
+    }
+    if constexpr (!IN_COLLECTION) {
+        const size_t page_start_row = _page_reader->start_row();
+        const size_t page_end_row = _page_reader->end_row();
+        if (UNLIKELY(page_end_row < page_start_row ||
+                     static_cast<size_t>(page_num_values) != page_end_row - 
page_start_row)) {
+            // Flat columns have exactly one physical value slot per logical 
row. Rejecting a
+            // divergent header/OffsetIndex span prevents every later page 
from shifting rows.
+            return Status::Corruption(
+                    "Parquet flat data page has {} values for logical row 
range [{}, {})",
+                    page_num_values, page_start_row, page_end_row);
+        }
+    }
+    _remaining_rep_nums = page_num_values;
+    _remaining_def_nums = page_num_values;
+    _remaining_num_values = page_num_values;
+
+    // no offset will parse all header.
+    if constexpr (OFFSET_INDEX == false) {
+        _chunk_parsed_values += _remaining_num_values;
+    }
+    _state = HEADER_PARSED;
+    return Status::OK();
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::next_page() {
+    _state = INITIALIZED;
+    RETURN_IF_ERROR(_page_reader->next_page());
+    return Status::OK();
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, 
OFFSET_INDEX>::_get_uncompressed_levels(
+        const tparquet::DataPageHeaderV2& page_v2, Slice& page_data) {
+    const size_t rl = page_v2.repetition_levels_byte_length;
+    const size_t dl = page_v2.definition_levels_byte_length;
+    if (UNLIKELY(rl > page_data.size || dl > page_data.size - rl)) {
+        // Validate the physical slice again because a cached entry may itself 
be truncated.
+        return Status::Corruption("Parquet data page v2 level bytes exceed 
available payload");
+    }
+    _v2_rep_levels = Slice(page_data.data, rl);
+    _v2_def_levels = Slice(page_data.data + rl, dl);
+    page_data.data += dl + rl;
+    page_data.size -= dl + rl;
+    return Status::OK();
+}
+
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::load_page_data() {
+    if (_state == DATA_LOADED) {
+        return Status::OK();
+    }
+    if (UNLIKELY(_state != HEADER_PARSED)) {
+        return Status::Corruption("Should parse page header");
+    }
+
+    const tparquet::PageHeader* header = nullptr;
+    RETURN_IF_ERROR(_page_reader->get_page_header(&header));
+    int32_t uncompressed_size = header->uncompressed_page_size;
+    bool page_loaded = false;
+
+    // First, try to reuse a cache handle previously discovered by PageReader
+    // (header-only lookup) to avoid a second lookup here.
+    if (_page_read_ctx.enable_parquet_file_page_cache && 
!config::disable_storage_page_cache &&
+        StoragePageCache::instance() != nullptr) {
+        if (_page_reader->has_page_cache_handle()) {
+            const PageCacheHandle& handle = _page_reader->page_cache_handle();
+            Slice cached = handle.data();
+            size_t header_size = _page_reader->header_bytes().size();
+            size_t levels_size = 0;
+            if (header->__isset.data_page_header_v2) {
+                const tparquet::DataPageHeaderV2& header_v2 = 
header->data_page_header_v2;
+                size_t rl = header_v2.repetition_levels_byte_length;
+                size_t dl = header_v2.definition_levels_byte_length;
+                levels_size = rl + dl;
+                if (UNLIKELY(header_size > cached.size ||
+                             levels_size > cached.size - header_size)) {
+                    return Status::Corruption("Cached Parquet page is shorter 
than its v2 levels");
+                }
+                _v2_rep_levels =
+                        Slice(reinterpret_cast<const uint8_t*>(cached.data) + 
header_size, rl);
+                _v2_def_levels =
+                        Slice(reinterpret_cast<const uint8_t*>(cached.data) + 
header_size + rl, dl);
+            }
+            // payload_slice points to the bytes after header and levels
+            if (UNLIKELY(header_size + levels_size > cached.size)) {
+                return Status::Corruption("Cached Parquet page is shorter than 
its header");
+            }
+            Slice payload_slice(cached.data + header_size + levels_size,
+                                cached.size - header_size - levels_size);
+
+            bool cache_payload_is_decompressed = 
_page_reader->is_cache_payload_decompressed();
+            const size_t expected_payload_size =
+                    cache_payload_is_decompressed
+                            ? 
static_cast<size_t>(header->uncompressed_page_size) - levels_size
+                            : 
static_cast<size_t>(header->compressed_page_size) - levels_size;
+            if (UNLIKELY(payload_slice.size != expected_payload_size)) {
+                return Status::Corruption("Cached Parquet page payload has 
size {}, expected {}",
+                                          payload_slice.size, 
expected_payload_size);
+            }
+
+            if (cache_payload_is_decompressed) {
+                // Cached payload is already uncompressed
+                _page_data = payload_slice;
+            } else {
+                CHECK(_block_compress_codec);
+                // Decompress cached payload into _decompress_buf for decoding
+                size_t uncompressed_payload_size =
+                        header->__isset.data_page_header_v2
+                                ? 
static_cast<size_t>(header->uncompressed_page_size) - levels_size
+                                : 
static_cast<size_t>(header->uncompressed_page_size);
+                _reserve_decompress_buf(uncompressed_payload_size);
+                _page_data = Slice(_decompress_buf.get(), 
uncompressed_payload_size);
+                SCOPED_RAW_TIMER(&_chunk_statistics.decompress_time);
+                _chunk_statistics.decompress_cnt++;
+                
RETURN_IF_ERROR(_block_compress_codec->decompress(payload_slice, &_page_data));
+                if (UNLIKELY(_page_data.size != uncompressed_payload_size)) {
+                    return Status::Corruption("Parquet page decompressed to {} 
bytes, expected {}",
+                                              _page_data.size, 
uncompressed_payload_size);
+                }
+            }
+            // page cache counters were incremented when PageReader did the 
header-only
+            // cache lookup. Do not increment again to avoid double-counting.
+            page_loaded = true;
+        }
+    }
+
+    if (!page_loaded) {
+        if (_block_compress_codec != nullptr) {
+            Slice compressed_data;
+            RETURN_IF_ERROR(_page_reader->get_page_data(compressed_data));
+            std::vector<uint8_t> level_bytes;

Review Comment:
   [P2] Avoid copying V2 levels when the page cannot be cached
   
   This allocates and copies every repetition/definition byte before checking 
whether the session cache is enabled, storage caching is globally disabled, or 
a cache instance exists. The live decoders use 
`_v2_rep_levels`/`_v2_def_levels` slices into the page buffer; `level_bytes` is 
consumed only by `_insert_page_into_cache`. Thus cache-disabled nested V2 scans 
still pay one allocation, memcpy, and free per page. Please build this buffer 
only inside cache admission (or reuse bounded scratch only for an actual 
insertion), with a cache-disabled coverage hook/benchmark.



-- 
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