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]