github-actions[bot] commented on code in PR #68780:
URL: https://github.com/apache/doris/pull/68780#discussion_r4220001448
##########
be/src/format_v2/parquet/reader/native/column_chunk_reader.cpp:
##########
@@ -1142,19 +1140,58 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::parse_page_header() {
_remaining_def_nums = page_num_values;
_remaining_num_values = page_num_values;
- // no offset will parse all header.
- if (!active_offset_index) {
- _chunk_parsed_values += _remaining_num_values;
- }
+ // Count each parsed data header exactly once, including the indexed
prefix. A safe fallback
+ // uses this physical value count for EOF; nested logical rows are not
interchangeable with it.
+ _chunk_parsed_values += _remaining_num_values;
_first_data_page_parsed = true;
_state = HEADER_PARSED;
return Status::OK();
}
+template <bool IN_COLLECTION, bool OFFSET_INDEX>
+Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_finish_indexed_nested_page() {
+ if constexpr (IN_COLLECTION && OFFSET_INDEX) {
+ if (_state != INITIALIZED && _page_reader->has_active_offset_index() &&
+ !_page_reader->is_header_v2()) {
+ // A V1 header has no logical row count. Loading only a prefix
does not validate its
+ // indexed span: drain the remaining level runs before trusting
the next indexed row.
+ // No value decoding or per-value scratch is needed for a page we
are leaving.
+ RETURN_IF_ERROR(load_page_data_idempotent());
+ SCOPED_RAW_TIMER(&_chunk_statistics.decode_level_time);
+ const size_t end_row = _page_reader->end_row();
+ while (_remaining_rep_nums > 0) {
+ level_t level = -1;
+ const size_t run = _rep_level_decoder.get_next_run(&level,
_remaining_rep_nums);
+ if (run == 0) {
+ return Status::Corruption(
+ "Parquet repetition level stream ended while
leaving an indexed page");
+ }
+ if (level == 0) {
+ if (_current_row > end_row || run > end_row -
_current_row) {
+ return Status::Corruption(
+ "Parquet nested page exceeds its indexed row
span");
+ }
+ _current_row += run;
+ }
+ _remaining_rep_nums -= run;
+ }
+ if (_current_row != end_row) {
+ return Status::Corruption(
+ "Parquet nested page does not match its indexed row
span");
+ }
+ }
+ }
+ return Status::OK();
+}
+
template <bool IN_COLLECTION, bool OFFSET_INDEX>
Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::next_page() {
if constexpr (OFFSET_INDEX) {
RETURN_IF_ERROR(ensure_first_data_page_parsed());
+ RETURN_IF_ERROR(_finish_indexed_nested_page());
+ if (_state == INITIALIZED && _page_reader->has_active_offset_index()) {
Review Comment:
[P1] Verify skipped page spans before trusting a later active OffsetIndex
coordinate. The new flag only rejects a subsequent index *discard*. With
physical flat page counts `[1,3,1,2]`, an index advertising starts `[0,1,3,4]`
and correct byte sizes passes structural validation. Selecting only row 3 skips
the three-row second page unparsed; the next page's one value matches its local
indexed `[3,4)` span, so it returns physical row 4 as logical row 3. Nested
V1/V2 reads have the same skipped-prefix gap. This differs from the existing
skip-then-fallback and parsed-V1 threads: the index remains active here.
Validate the skipped span or fail closed before returning selected rows, and
cover this active-index sparse seek.
##########
be/test/format_v2/parquet/native_decoder_test.cpp:
##########
@@ -3454,6 +3457,81 @@ TEST(ParquetV2NativeDecoderTest,
LazyFixedWidthFilterUsesReconciledFirstPageRang
EXPECT_EQ(row_filter, (IColumn::Filter {1, 1}));
}
+// Exercise the adapter and native reader together: an empty fragment must not
change the
+// predicate domain, and an index fallback must still refresh every subsequent
page's row bounds.
+void check_fixed_width_filter_across_pages(bool shifted_index) {
+ const std::vector<std::vector<uint8_t>> pages
{serialize_plain_int32_page({10, 11}),
+
serialize_plain_int32_page({20, 21}),
+
serialize_plain_int32_page({30, 31})};
+ std::vector<uint8_t> bytes;
+ tparquet::OffsetIndex offset_index;
+ for (size_t page = 0; page < pages.size(); ++page) {
+ tparquet::PageLocation location;
+ location.__set_offset(bytes.size() + (shifted_index && page > 0 ? 1 :
0));
+ location.__set_compressed_page_size(pages[page].size() +
+ (shifted_index && page == 0 ? 1 :
0));
+ location.__set_first_row_index(page * 2);
+ offset_index.page_locations.push_back(location);
+ bytes.insert(bytes.end(), pages[page].begin(), pages[page].end());
+ }
+ bytes.push_back(0);
+ tparquet::ColumnChunk chunk;
+ chunk.meta_data.__set_type(tparquet::Type::INT32);
+ chunk.meta_data.__set_codec(tparquet::CompressionCodec::UNCOMPRESSED);
+ chunk.meta_data.__set_num_values(6);
+ chunk.meta_data.__set_total_compressed_size(bytes.size());
+ chunk.meta_data.__set_data_page_offset(0);
+ chunk.meta_data.__set_encodings({tparquet::Encoding::PLAIN});
+ NativeFieldSchema field;
+ field.physical_type = tparquet::Type::INT32;
+ field.data_type = std::make_shared<DataTypeInt32>();
+ field.parquet_schema.__set_type(tparquet::Type::INT32);
+
field.parquet_schema.__set_repetition_type(tparquet::FieldRepetitionType::REQUIRED);
+ auto ranges = ::doris::RowRanges::create_single(0, shifted_index ? 6 : 2);
+ if (!shifted_index) {
+ ranges.add({4, 6});
+ }
+ auto file = std::make_shared<NativeDecoderMemoryFileReader>(bytes);
+ auto native = std::make_unique<ScalarColumnReader<false, true>>(ranges, 6,
chunk, &offset_index,
+ nullptr,
nullptr);
+ ASSERT_TRUE(native->init(file, &field, bytes.size(), nullptr, "",
ParquetReaderCompat {}, true)
+ .ok());
+ ParquetColumnSchema schema;
+ schema.name = "value";
+ schema.type = field.data_type;
+ NativeColumnReader reader(schema, schema.type, schema.type, nullptr, {});
Review Comment:
[P1] Compile the new adapter test through an accessible API. This helper
constructs `NativeColumnReader` here, then assigns `_native_reader` and
`_row_group_rows` and calls `read_with_fixed_width_filter()`; all are declared
after `private:` in `native_column_reader.h`. `be/test/CMakeLists.txt` includes
this file, so C++ access checking fails before the two new regressions can run.
Use the public factory and `select_with_fixed_width_filter()`, or introduce a
narrow test seam.
##########
be/test/format_v2/parquet/native_decoder_test.cpp:
##########
@@ -3559,6 +3643,846 @@ TEST(ParquetV2NativeDecoderTest,
LazyNestedV2SeekValidatesFirstPageRowRange) {
EXPECT_TRUE(status.is<ErrorCode::CORRUPTION>()) << status;
}
+TEST(ParquetV2NativeDecoderTest,
SelectedLaterPageReconcilesRangesAfterIndexFallback) {
+ for (const bool shrinking : {false, true}) {
+ for (const int mode : {0, 1, 2}) {
+ for (const size_t batch_size : {1, 2, 3}) {
+ for (const bool cache_hit : {false, true}) {
+ SCOPED_TRACE(testing::Message()
+ << "mode=" << mode << ", batch=" << batch_size
+ << ", cache=" << cache_hit << ", shrinking="
<< shrinking);
+ const std::vector<std::vector<int32_t>> values {
+ {10, 11}, {20, 21, 22}, {30, 31}};
+ std::vector<uint8_t> bytes(1, 0);
+ auto dictionary = ColumnInt32::create();
+ dictionary->get_data() = {10, 11, 20, 21, 22, 30, 31};
+ if (mode == 2) {
+ std::vector<uint8_t> payload(dictionary->size() *
sizeof(int32_t));
+ memcpy(payload.data(), dictionary->get_data().data(),
payload.size());
+ tparquet::PageHeader header;
+ header.__set_type(tparquet::PageType::DICTIONARY_PAGE);
+ header.__set_compressed_page_size(payload.size());
+ header.__set_uncompressed_page_size(payload.size());
+ tparquet::DictionaryPageHeader data;
+ data.__set_num_values(dictionary->size());
+ data.__set_encoding(tparquet::Encoding::PLAIN);
+ header.__set_dictionary_page_header(data);
+ auto page = serialize_page(header, payload);
+ bytes.insert(bytes.end(), page.begin(), page.end());
+ }
+ const size_t data_offset = bytes.size();
+ tparquet::OffsetIndex index;
+ std::vector<std::vector<uint8_t>> pages;
+ uint32_t next_id = 0;
+ for (size_t i = 0; i < values.size(); ++i) {
+ if (mode == 2) {
+ faststring ids;
+ RleEncoder<uint32_t> encoder(&ids, 3);
+ for (size_t j = 0; j < values[i].size(); ++j) {
+ encoder.Put(next_id++);
+ }
+ for (size_t j = values[i].size(); j < 8; ++j) {
+ encoder.Put(0);
+ }
+ encoder.Flush();
+ std::vector<uint8_t> payload(ids.size() + 1, 3);
+ memcpy(payload.data() + 1, ids.data(), ids.size());
+ tparquet::PageHeader header;
+ header.__set_type(tparquet::PageType::DATA_PAGE);
+ header.__set_compressed_page_size(payload.size());
+
header.__set_uncompressed_page_size(payload.size());
+ tparquet::DataPageHeader data;
+ data.__set_num_values(values[i].size());
+
data.__set_encoding(tparquet::Encoding::RLE_DICTIONARY);
+
data.__set_definition_level_encoding(tparquet::Encoding::RLE);
+
data.__set_repetition_level_encoding(tparquet::Encoding::RLE);
+ header.__set_data_page_header(data);
+ pages.push_back(serialize_page(header, payload));
+ } else {
+
pages.push_back(serialize_plain_int32_page(values[i]));
+ }
+ tparquet::PageLocation location;
+ location.__set_offset(bytes.size());
+
location.__set_compressed_page_size(pages.back().size() - (i == 1 ? 1 : 0));
+ // Reconciliation can expand a selected range or
remove it entirely.
+ // Neither case may advance using the stale indexed
page boundary.
+ location.__set_first_row_index(shrinking && i == 2 ? 6
: i * 2);
+ index.page_locations.push_back(location);
+ bytes.insert(bytes.end(), pages.back().begin(),
pages.back().end());
+ }
+ const std::string cache_key =
+ fmt::format("selected-late-{}-{}-{}", mode,
batch_size, shrinking);
+ if (cache_hit) {
+ auto* page = new DataPage(pages[1].size(), true,
segment_v2::DATA_PAGE);
+ memcpy(page->data(), pages[1].data(), pages[1].size());
+ page->reset_size(pages[1].size());
+ PageCacheHandle handle;
+ StoragePageCache::instance()->insert(
+ StoragePageCache::CacheKey(cache_key,
bytes.size(),
+
index.page_locations[1].offset),
+ page, &handle, segment_v2::DATA_PAGE);
+ }
+ tparquet::ColumnChunk chunk;
+ chunk.meta_data.__set_type(tparquet::Type::INT32);
+
chunk.meta_data.__set_codec(tparquet::CompressionCodec::UNCOMPRESSED);
+ chunk.meta_data.__set_num_values(7);
+ chunk.meta_data.__set_total_compressed_size(bytes.size() -
1);
+ chunk.meta_data.__set_data_page_offset(data_offset);
+
chunk.meta_data.__set_encodings({tparquet::Encoding::PLAIN});
+ if (mode == 2) {
+ chunk.meta_data.__set_dictionary_page_offset(1);
+ chunk.meta_data.__set_encodings(
+ {tparquet::Encoding::PLAIN,
tparquet::Encoding::RLE_DICTIONARY});
+ }
+ NativeFieldSchema field;
+ field.physical_type = tparquet::Type::INT32;
+ field.data_type = std::make_shared<DataTypeInt32>();
+ field.parquet_schema.__set_type(tparquet::Type::INT32);
+ field.parquet_schema.__set_repetition_type(
+ tparquet::FieldRepetitionType::REQUIRED);
+ const auto ranges =
+ ::doris::RowRanges::create_single(shrinking ? 5 :
2, shrinking ? 6 : 5);
+ const size_t expected_rows = shrinking ? 1 : 3;
+ auto file =
std::make_shared<NativeDecoderMemoryFileReader>(bytes);
+ ScalarColumnReader<false, true> reader(ranges, 7, chunk,
&index, nullptr,
+ nullptr);
+ ASSERT_TRUE(reader.init(file, &field, bytes.size(),
nullptr,
+ cache_hit ? cache_key : "",
ParquetReaderCompat {},
+ true)
+ .ok());
+ FilterMap filter;
+ ASSERT_TRUE(filter.init(nullptr, expected_rows,
false).ok());
+ ColumnPtr ordinary = ColumnInt32::create();
+ auto projected = ColumnInt32::create();
+ size_t consumed = 0;
+ bool eof = false;
+ for (size_t call = 0; consumed < expected_rows && !eof;
++call) {
+ ASSERT_LT(call, 8);
+ size_t rows = 0;
+ Status status;
+ IColumn::Filter row_filter;
+ bool used = false;
+ const size_t count = std::min(batch_size,
expected_rows - consumed);
+ if (mode == 0) {
+ status = reader.read_column_data(ordinary,
field.data_type, nullptr,
+ filter, count,
&rows, &eof, false);
+ } else if (mode == 1) {
+ DirectPredicateExecutionKind kind =
DirectPredicateExecutionKind::NONE;
+ status = reader.read_fixed_width_filter(
+ {create_int32_raw_comparison(0, "gt",
TExprOpcode::GT, 0)}, 0,
+ filter, count, projected.get(),
&row_filter, &rows, &eof, &used,
+ &kind);
+ } else {
+ size_t survivors = 0;
+ bool direct = false;
+ status = reader.read_dictionary_filter(
+ IColumn::Filter(7, 1), filter, count,
dictionary.get(),
+ projected.get(), nullptr, &row_filter,
&survivors, &rows, &eof,
+ &direct, &used);
+ EXPECT_EQ(survivors, rows);
+ }
+ ASSERT_TRUE(status.ok()) << status;
+ if (mode != 0) {
+ EXPECT_TRUE(used);
+ EXPECT_EQ(row_filter, IColumn::Filter(rows, 1));
+ }
+ consumed += rows;
+ }
+ EXPECT_EQ(consumed, expected_rows);
+ const auto& actual =
+ mode == 0 ? assert_cast<const
ColumnInt32&>(*ordinary) : *projected;
+ EXPECT_EQ(actual.get_data(), (shrinking ?
ColumnInt32::Container {30}
+ :
ColumnInt32::Container {20, 21, 22}));
+
EXPECT_EQ(reader.column_statistics().page_cache_hit_counter > 0, cache_hit);
+ }
+ }
+ }
+ }
+}
+
+// Serialize independent page versions and row shapes so navigation tests do
not inherit a writer's
+// fixed page layout. Expected values are derived from the physical row
sequence, never the index.
+struct NestedNavigationFile {
+ std::vector<uint8_t> bytes {0};
+ tparquet::ColumnChunk chunk;
+ tparquet::OffsetIndex index;
+ NativeFieldSchema field;
+ size_t rows = 0;
+ std::vector<std::vector<int32_t>> row_values;
+ std::vector<std::vector<uint8_t>> pages;
+
+ NestedNavigationFile(const std::vector<bool>& v2, const
std::vector<size_t>& page_rows,
+ bool repeated = false, bool dictionary = false, bool
nullable = false) {
+ field.physical_type = tparquet::Type::INT32;
+ field.data_type = std::make_shared<DataTypeInt32>();
+ field.repetition_level = 1;
+ field.definition_level = nullable ? 1 : 0;
+ if (nullable) {
+ field.data_type = make_nullable(field.data_type);
+ }
+ field.parquet_schema.__set_type(tparquet::Type::INT32);
+ field.parquet_schema.__set_repetition_type(
+ nullable ? tparquet::FieldRepetitionType::OPTIONAL
+ : tparquet::FieldRepetitionType::REQUIRED);
+ const size_t row_count = std::accumulate(page_rows.begin(),
page_rows.end(), size_t {0});
+ const size_t value_count = row_count * (repeated ? 2 : 1);
+ if (dictionary) {
+ std::vector<int32_t> dictionary_values(value_count);
+ std::iota(dictionary_values.begin(), dictionary_values.end(), 100);
+ std::vector<uint8_t> payload(value_count * sizeof(int32_t));
+ memcpy(payload.data(), dictionary_values.data(), payload.size());
+ tparquet::PageHeader header;
+ header.__set_type(tparquet::PageType::DICTIONARY_PAGE);
+ header.__set_compressed_page_size(payload.size());
+ header.__set_uncompressed_page_size(payload.size());
+ tparquet::DictionaryPageHeader data;
+ data.__set_num_values(value_count);
+ data.__set_encoding(tparquet::Encoding::PLAIN);
+ header.__set_dictionary_page_header(data);
+ const auto page = serialize_page(header, payload);
+ bytes.insert(bytes.end(), page.begin(), page.end());
+ chunk.meta_data.__set_dictionary_page_offset(1);
+ }
+ chunk.meta_data.__set_data_page_offset(bytes.size());
+ size_t value_index = 0;
+ for (size_t i = 0; i < v2.size(); ++i) {
+ std::vector<int32_t> values;
+ std::vector<uint32_t> levels;
+ std::vector<uint32_t> definitions;
+ for (size_t row = 0; row < page_rows[i]; ++row) {
+ row_values.emplace_back();
+ for (size_t element = 0; element < (repeated ? 2 : 1);
++element) {
+ levels.push_back(element == 0 ? 0 : 1);
+ const bool is_null = nullable && value_index % 3 == 1;
+ const int32_t value = 100 + value_index++;
+ definitions.push_back(is_null ? 0 : 1);
+ if (!is_null) {
+ values.push_back(value);
+ }
+ row_values.back().push_back(is_null ?
std::numeric_limits<int32_t>::min()
+ : value);
+ }
+ }
+ faststring encoded_levels;
+ RleEncoder<uint32_t> level_encoder(&encoded_levels, 1);
+ for (const auto level : levels) {
+ level_encoder.Put(level);
+ }
+ for (size_t padding = levels.size(); padding % 8 != 0; ++padding) {
+ level_encoder.Put(0);
+ }
+ level_encoder.Flush();
+ std::vector<uint8_t> payload;
+ if (!v2[i]) {
+ const uint32_t size = encoded_levels.size();
+ const auto* raw = reinterpret_cast<const uint8_t*>(&size);
+ payload.insert(payload.end(), raw, raw + sizeof(size));
+ }
+ payload.insert(payload.end(), encoded_levels.data(),
+ encoded_levels.data() + encoded_levels.size());
+ faststring encoded_definitions;
+ if (nullable) {
+ RleEncoder<uint32_t> def_encoder(&encoded_definitions, 1);
+ for (const auto level : definitions) {
+ def_encoder.Put(level);
+ }
+ for (size_t padding = definitions.size(); padding % 8 != 0;
++padding) {
+ def_encoder.Put(0);
+ }
+ def_encoder.Flush();
+ if (!v2[i]) {
+ const uint32_t size = encoded_definitions.size();
+ const auto* raw = reinterpret_cast<const uint8_t*>(&size);
+ payload.insert(payload.end(), raw, raw + sizeof(size));
+ }
+ payload.insert(payload.end(), encoded_definitions.data(),
+ encoded_definitions.data() +
encoded_definitions.size());
+ }
+ if (dictionary) {
+ faststring ids;
+ RleEncoder<uint32_t> encoder(&ids, 8);
+ for (const auto value : values) {
+ encoder.Put(value - 100);
+ }
+ for (size_t padding = values.size(); padding % 8 != 0;
++padding) {
+ encoder.Put(0);
+ }
+ encoder.Flush();
+ payload.push_back(8);
+ payload.insert(payload.end(), ids.data(), ids.data() +
ids.size());
+ } else {
+ const auto* raw = reinterpret_cast<const
uint8_t*>(values.data());
+ payload.insert(payload.end(), raw, raw + values.size() *
sizeof(int32_t));
+ }
+ const auto encoding =
+ dictionary ? tparquet::Encoding::RLE_DICTIONARY :
tparquet::Encoding::PLAIN;
+ tparquet::PageHeader header;
+ header.__set_type(v2[i] ? tparquet::PageType::DATA_PAGE_V2
+ : tparquet::PageType::DATA_PAGE);
+ header.__set_compressed_page_size(payload.size());
+ header.__set_uncompressed_page_size(payload.size());
+ if (v2[i]) {
+ tparquet::DataPageHeaderV2 data;
+ data.__set_num_values(levels.size());
+ data.__set_num_rows(page_rows[i]);
+ data.__set_num_nulls(levels.size() - values.size());
+ data.__set_encoding(encoding);
+
data.__set_repetition_levels_byte_length(encoded_levels.size());
+
data.__set_definition_levels_byte_length(encoded_definitions.size());
+ data.__set_is_compressed(false);
+ header.__set_data_page_header_v2(data);
+ } else {
+ tparquet::DataPageHeader data;
+ data.__set_num_values(levels.size());
+ data.__set_encoding(encoding);
+ data.__set_repetition_level_encoding(tparquet::Encoding::RLE);
+ data.__set_definition_level_encoding(tparquet::Encoding::RLE);
+ header.__set_data_page_header(data);
+ }
+ pages.push_back(serialize_page(header, payload));
+ tparquet::PageLocation location;
+ location.__set_offset(bytes.size());
+ location.__set_compressed_page_size(pages.back().size());
+ location.__set_first_row_index(rows);
+ index.page_locations.push_back(location);
+ rows += page_rows[i];
+ bytes.insert(bytes.end(), pages.back().begin(),
pages.back().end());
+ }
+ chunk.meta_data.__set_type(tparquet::Type::INT32);
+ chunk.meta_data.__set_codec(tparquet::CompressionCodec::UNCOMPRESSED);
+ chunk.meta_data.__set_num_values(value_count);
+ chunk.meta_data.__set_total_compressed_size(bytes.size() - 1);
+ chunk.meta_data.__set_encodings(
+ dictionary
+ ? std::vector<tparquet::Encoding::type>
{tparquet::Encoding::PLAIN,
+
tparquet::Encoding::RLE_DICTIONARY}
+ : std::vector<tparquet::Encoding::type>
{tparquet::Encoding::PLAIN});
+ }
+
+ Status read(const RowRanges& ranges, size_t batch_size, bool indexed, bool
cache,
+ std::vector<int32_t>* actual, bool levels_only = false) {
+ static size_t cache_id = 0;
+ const std::string cache_key = cache ?
fmt::format("nested-navigation-{}", ++cache_id) : "";
+ if (cache) {
+ for (size_t i = 0; i < pages.size(); ++i) {
+ auto* page = new DataPage(pages[i].size(), true,
segment_v2::DATA_PAGE);
+ memcpy(page->data(), pages[i].data(), pages[i].size());
+ page->reset_size(pages[i].size());
+ PageCacheHandle handle;
+ StoragePageCache::instance()->insert(
+ StoragePageCache::CacheKey(cache_key, bytes.size(),
+
index.page_locations[i].offset),
+ page, &handle, segment_v2::DATA_PAGE);
+ }
+ }
+ auto file = std::make_shared<NativeDecoderMemoryFileReader>(bytes);
+ ScalarColumnReader<true, true> reader(ranges, rows, chunk, indexed ?
&index : nullptr,
+ nullptr, nullptr);
+ RETURN_IF_ERROR(reader.init(file, &field, bytes.size(), nullptr,
cache_key,
+ ParquetReaderCompat {}, true));
+ reader.set_column_in_nested();
+ FilterMap filter;
+ RETURN_IF_ERROR(filter.init(nullptr, rows, false));
+ bool eof = false;
+ size_t calls = 0;
+ while (!eof) {
+ if (++calls > rows + pages.size() + 1) {
+ return Status::InternalError("Nested navigation did not reach
EOF");
+ }
+ ColumnPtr output = field.data_type->create_column();
+ size_t read_rows = 0;
+ if (levels_only) {
+ RETURN_IF_ERROR(reader.read_column_levels(filter, batch_size,
&read_rows, &eof));
+ actual->insert(actual->end(), read_rows, 0);
+ } else {
+ RETURN_IF_ERROR(reader.read_column_data(output,
field.data_type, nullptr, filter,
+ batch_size,
&read_rows, &eof, false));
+ if (field.data_type->is_nullable()) {
+ const auto& nullable_output = assert_cast<const
ColumnNullable&>(*output);
+ const auto& values =
+ assert_cast<const
ColumnInt32&>(nullable_output.get_nested_column())
+ .get_data();
+ for (size_t i = 0; i < values.size(); ++i) {
+ actual->push_back(nullable_output.is_null_at(i)
+ ?
std::numeric_limits<int32_t>::min()
+ : values[i]);
+ }
+ } else {
+ const auto& values = assert_cast<const
ColumnInt32&>(*output).get_data();
+ actual->insert(actual->end(), values.begin(),
values.end());
+ }
+ }
+ }
+ return Status::OK();
+ }
+
+ std::vector<int32_t> expected(const RowRanges& ranges) const {
+ std::vector<int32_t> result;
+ for (size_t i = 0; i < ranges.range_size(); ++i) {
+ for (size_t row = ranges.get_range_from(i); row <
ranges.get_range_to(i); ++row) {
+ result.insert(result.end(), row_values[row].begin(),
row_values[row].end());
+ }
+ }
+ return result;
+ }
+};
+
+TEST(ParquetV2NativeDecoderTest,
PartialNestedPageCannotRelabelLaterIndexedRows) {
+ for (const bool cache : {false, true}) {
+ for (const bool levels_only : {false, true}) {
+ NestedNavigationFile file({true, false, false}, {1, 3, 2}, true);
+ file.index.page_locations[2].first_row_index = 3;
+ auto ranges = RowRanges::create_single(1, 2);
+ ranges.add({3, 4});
+ std::vector<int32_t> actual;
+ auto status = file.read(ranges, 1, true, cache, &actual,
levels_only);
+ EXPECT_TRUE(status.is<ErrorCode::CORRUPTION>()) << status;
+ }
+ }
+}
+
+TEST(ParquetV2NativeDecoderTest, FinalNestedFallbackPreservesColumnValueCount)
{
+ for (const bool repeated : {false, true}) {
+ for (const bool cache : {false, true}) {
+ NestedNavigationFile file({true, false}, {1, 1}, repeated);
+ --file.index.page_locations.back().compressed_page_size;
+ const auto ranges = RowRanges::create_single(0, 2);
+ std::vector<int32_t> actual;
+ const auto status = file.read(ranges, 1, true, cache, &actual);
+ ASSERT_TRUE(status.ok()) << status;
+ EXPECT_EQ(actual, file.expected(ranges));
+ }
+ }
+}
+
+TEST(ParquetV2NativeDecoderTest,
NestedNavigationPreservesRowsAcrossPageVersionTransitions) {
+ for (int versions = 0; versions < 8; ++versions) {
+ for (const bool indexed : {false, true}) {
+ for (const bool dictionary : {false, true}) {
+ for (const bool nullable : {false, true}) {
+ for (const bool cache : {false, true}) {
+ for (const size_t batch : {1, 2, 7}) {
+ for (int selection = 0; selection < 3;
++selection) {
+ SCOPED_TRACE(testing::Message()
+ << "versions=" << versions << ",
indexed=" << indexed
+ << ", dictionary=" << dictionary
+ << ", nullable=" << nullable <<
", cache=" << cache
+ << ", batch=" << batch << ",
selection=" << selection);
+ NestedNavigationFile file({bool(versions & 1),
bool(versions & 2),
+ bool(versions & 4)},
+ {2, 3, 2}, true,
dictionary, nullable);
+ auto ranges = RowRanges::create_single(0, 7);
+ if (selection == 1) {
+ ranges = RowRanges::create_single(1, 2);
+ ranges.add({3, 4});
+ ranges.add({6, 7});
+ } else if (selection == 2) {
+ ranges = RowRanges::create_single(6, 7);
+ }
+ for (const bool levels_only : {false, true}) {
+ std::vector<int32_t> actual;
+ const auto status = file.read(ranges,
batch, indexed, cache,
+ &actual,
levels_only);
+ ASSERT_TRUE(status.ok()) << status;
+ EXPECT_EQ(actual,
+ levels_only ?
std::vector<int32_t>(ranges.count(), 0)
+ :
file.expected(ranges));
+ }
+ }
+ }
+ }
+ }
+ }
+ }
+ }
+}
+
+TEST(ParquetV2NativeDecoderTest, NestedNavigationFallbackMatrix) {
+ for (int versions = 0; versions < 8; ++versions) {
+ for (int stale_page = 0; stale_page < 3; ++stale_page) {
+ for (const bool cache : {false, true}) {
+ for (const bool dictionary : {false, true}) {
+ for (const size_t batch : {1, 2, 7}) {
+ for (int selection = 0; selection < 3; ++selection) {
+ SCOPED_TRACE(testing::Message()
+ << "versions=" << versions << ",
stale=" << stale_page
+ << ", cache=" << cache << ",
dictionary=" << dictionary
+ << ", batch=" << batch << ",
selection=" << selection);
+ NestedNavigationFile file(
+ {bool(versions & 1), bool(versions & 2),
bool(versions & 4)},
+ {2, 3, 2}, true, dictionary, true);
+
--file.index.page_locations[stale_page].compressed_page_size;
+ auto ranges = RowRanges::create_single(0, 7);
+ if (selection == 1) {
+ ranges = RowRanges::create_single(1, 2);
+ ranges.add({3, 4});
+ ranges.add({6, 7});
+ } else if (selection == 2) {
+ ranges = RowRanges::create_single(6, 7);
+ }
+ for (const bool levels_only : {false, true}) {
+ std::vector<int32_t> actual;
+ const auto status =
+ file.read(ranges, batch, true, cache,
&actual, levels_only);
+ // Only a final-page fallback after an
unparsed middle page lacks a
+ // complete prefix. All other transitions must
preserve the exact rows.
+ if ((versions & 1) && stale_page == 2 &&
selection == 2) {
+
EXPECT_TRUE(status.is<ErrorCode::CORRUPTION>()) << status;
+ } else {
+ ASSERT_TRUE(status.ok()) << status;
+ EXPECT_EQ(actual,
+ levels_only ?
std::vector<int32_t>(ranges.count(), 0)
+ :
file.expected(ranges));
+ }
+ }
+ }
+ }
+ }
+ }
+ }
+ }
+}
+
+TEST(ParquetV2NativeDecoderTest, NestedFallbackReevaluatesSelectedPageRange) {
+ NestedNavigationFile file({true, true, false}, {2, 1, 3}, true);
+ file.index.page_locations[2].first_row_index = 4;
+ --file.index.page_locations[1].compressed_page_size;
+ const auto ranges = RowRanges::create_single(3, 4);
+ std::vector<int32_t> actual;
+ const auto status = file.read(ranges, 1, true, false, &actual);
+ ASSERT_TRUE(status.ok()) << status;
+ EXPECT_EQ(actual, file.expected(ranges));
+}
+
+TEST(ParquetV2NativeDecoderTest, NestedNavigationAllRowSelections) {
+ for (int versions = 0; versions < 8; ++versions) {
+ for (unsigned mask = 1; mask < 128; ++mask) {
+ RowRanges ranges;
+ for (int64_t row = 0; row < 7; ++row) {
+ if (mask & (1U << row)) {
+ ranges.add({row, row + 1});
+ }
+ }
+ // Every nonempty row subset is checked against physical values.
Vary encoding, nulls
+ // and cache independently of the navigation policy, including
partial and whole skips.
+ for (int index_mode = -2; index_mode < 3; ++index_mode) {
+ NestedNavigationFile file(
+ {bool(versions & 1), bool(versions & 2), bool(versions
& 4)}, {2, 3, 2},
+ true, mask & 1, mask & 2);
+ if (index_mode >= 0) {
+
--file.index.page_locations[index_mode].compressed_page_size;
+ }
+ for (const size_t batch : {1, 2, 7}) {
+ for (const bool levels_only : {false, true}) {
+ SCOPED_TRACE(testing::Message()
+ << "versions=" << versions << ", mask="
<< mask
+ << ", index_mode=" << index_mode << ",
batch=" << batch
+ << ", levels_only=" << levels_only);
+ std::vector<int32_t> actual;
+ const auto status = file.read(ranges, batch,
index_mode != -2, mask & 4,
+ &actual, levels_only);
+ const bool skipped_middle = (mask & 0b0011100) == 0;
+ const bool selected_last = (mask & 0b1100000) != 0;
+ if ((versions & 1) && index_mode == 2 &&
skipped_middle && selected_last) {
+ EXPECT_TRUE(status.is<ErrorCode::CORRUPTION>()) <<
status;
+ } else {
+ ASSERT_TRUE(status.ok()) << status;
+ EXPECT_EQ(actual, levels_only ?
std::vector<int32_t>(ranges.count(), 0)
+ :
file.expected(ranges));
+ }
+ }
+ }
+ }
+ }
+ }
+}
+
+TEST(ParquetV2NativeDecoderTest,
IndexedNestedNavigationValidatesParsedPagesBeforeAdvancing) {
+ for (const bool bad_span : {false, true}) {
+ for (const bool partial : {false, true}) {
+ SCOPED_TRACE(testing::Message() << "bad_span=" << bad_span << ",
partial=" << partial);
+ NestedNavigationFile file({true, false, false}, {1, 3, 2}, true);
+ if (bad_span) {
+ file.index.page_locations[2].first_row_index = 3;
+ }
+ MemoryBufferedReader stream(file.bytes);
+ ColumnChunkReader<true, true> reader(&stream, &file.chunk,
&file.field, &file.index,
+ file.rows, nullptr,
ParquetPageReadContext());
+ ASSERT_TRUE(reader.init().ok());
+ ASSERT_TRUE(reader.parse_page_header().ok());
+ ASSERT_TRUE(reader.next_page().ok());
+ ASSERT_TRUE(reader.parse_page_header().ok());
+ ASSERT_TRUE(reader.parse_page_header().ok());
+ if (partial) {
+ ASSERT_TRUE(reader.load_page_data().ok());
+ std::vector<level_t> levels;
+ size_t rows = 0;
+ bool cross_page = false;
+ ASSERT_TRUE(reader.load_page_nested_rows(levels, 1, &rows,
&cross_page).ok());
+ EXPECT_EQ(rows, 1);
+ }
+ const auto status = reader.next_page();
+ if (bad_span) {
+ EXPECT_TRUE(status.is<ErrorCode::CORRUPTION>()) << status;
+ } else {
+ ASSERT_TRUE(status.ok()) << status;
+ EXPECT_EQ(reader.page_start_row(), 4);
+ ASSERT_TRUE(reader.parse_page_header().ok());
+ EXPECT_EQ(reader._chunk_parsed_values,
file.chunk.meta_data.num_values);
Review Comment:
[P1] Assert parsed-value accounting without reading a private field.
`_chunk_parsed_values` is declared below `private:` in `ColumnChunkReader`, so
this added `EXPECT_EQ` is a separate compile error even after the adapter
helper at line 3502 is fixed. Check the public read/EOF result or expose a
narrow test statistic instead, so the BE unit-test translation unit can compile.
--
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]