Gabriel39 commented on code in PR #68780:
URL: https://github.com/apache/doris/pull/68780#discussion_r4220418340
##########
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:
This is not a compile failure in the repository's BE unit-test target.
`be/test/CMakeLists.txt` explicitly sets `COMPILE_FLAGS "-fno-access-control"`
on `doris_be_test`:
https://github.com/apache/doris/blob/65f2abf24066087ea7b040fd593f51afd9af875f/be/test/CMakeLists.txt#L165
. The generated compile command for `native_decoder_test.cpp` contains that
flag.
I recompiled this translation unit and linked `doris_be_test` successfully,
then executed `FixedWidthFilterSkipsEmptyPageFragments` successfully. This is
the existing test-target access policy, not a local source modification or an
uncompiled test. No production API expansion or additional test seam is needed
for this finding.
##########
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:
The BE unit-test executable is compiled with `-fno-access-control`,
explicitly configured in the checked-in `be/test/CMakeLists.txt`:
https://github.com/apache/doris/blob/65f2abf24066087ea7b040fd593f51afd9af875f/be/test/CMakeLists.txt#L165
. I verified that the generated command for this exact translation unit
includes the flag, recompiled and linked the test executable, and ran
`IndexedNestedNavigationValidatesParsedPagesBeforeAdvancing` successfully.
Consequently this private-field assertion is valid under the repository's
unit-test build policy. The separate final-V1 fallback regression also checks
the externally observable exact-value/EOF behavior. No production accessor is
necessary to resolve the reported compilation concern.
--
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]