Skip to content

Commit fd0d0f1

Browse files
committed
[fix](parquet) Handle empty value sections
### What problem does this PR solve? Issue Number: close #66430 Related PR: None Problem Summary: A compressed Parquet DataPageV2 can legally contain definition levels but no physical values when every nullable value is NULL. The legacy reader passed the resulting null, zero-length slice to the BOOLEAN PLAIN decoder and aborted the BE in BatchedBitReader::Reset(). Use a dedicated empty-value-section decoder for these pages. Accept the page when definition levels require no physical values, and return Corruption when values are required but absent. Update page progress only after decoding succeeds so failures leave reader state unchanged. ### Release note Fix a BE crash when reading Parquet pages with empty value sections. ### Check List (For Author) - Test <!-- At least one of them must be included. --> - [ ] Regression test - [x] Unit Test - ./run-be-ut.sh --run --filter=ParquetColumnChunkReaderTest.* -j 2 - All 11 tests passed with ASAN. - [x] Manual test (add detailed scripts or steps below) - ./build.sh --be -j 2 - BE build succeeded. - [ ] No need to test or manual test. Explain why: - [ ] This is a refactor/code format and no logic has been changed. - [ ] Previous test can cover this change. - [ ] No code files have been changed. - [ ] Other reason <!-- Add your reason? --> - Behavior changed: - [ ] No. - [x] Yes. Valid all-NULL Parquet pages are accepted. Pages whose definition levels require missing physical values return Corruption instead of aborting the BE. - Does this need documentation? - [x] No. - [ ] Yes. <!-- Add document PR link here. eg: apache/doris-website#1214 --> ### Check List (For Reviewer who merge this PR) - [ ] Confirm the release note - [ ] Confirm test cases - [ ] Confirm document - [ ] Add branch pick label <!-- Add branch pick label that this PR should merge into -->
1 parent 30a699a commit fd0d0f1

3 files changed

Lines changed: 164 additions & 8 deletions

File tree

be/src/format/parquet/vparquet_column_chunk_reader.cpp

Lines changed: 62 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,35 @@ namespace cctz {
4343
class time_zone;
4444
} // namespace cctz
4545
namespace doris {
46+
47+
namespace {
48+
49+
class EmptyValueSectionDecoder final : public Decoder {
50+
public:
51+
Status decode_values(MutableColumnPtr& doris_column, DataTypePtr&,
52+
ColumnSelectVector& select_vector, bool) override {
53+
const size_t physical_values = select_vector.num_values() - select_vector.num_nulls();
54+
if (UNLIKELY(physical_values != 0)) {
55+
return Status::Corruption(
56+
"Parquet definition levels require {} values from an empty value section",
57+
physical_values);
58+
}
59+
doris_column->insert_many_defaults(select_vector.num_values() -
60+
select_vector.num_filtered());
61+
return Status::OK();
62+
}
63+
64+
Status skip_values(size_t num_values) override {
65+
if (UNLIKELY(num_values != 0)) {
66+
return Status::Corruption(
67+
"Parquet definition levels require {} values from an empty value section",
68+
num_values);
69+
}
70+
return Status::OK();
71+
}
72+
};
73+
74+
} // namespace
4675
namespace io {
4776
class BufferedStreamReader;
4877
struct IOContext;
@@ -358,17 +387,29 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::load_page_data() {
358387
}
359388

360389
// Reuse page decoder
390+
Decoder* encoding_decoder = nullptr;
361391
if (_decoders.find(static_cast<int>(encoding)) != _decoders.end()) {
362-
_page_decoder = _decoders[static_cast<int>(encoding)].get();
392+
encoding_decoder = _decoders[static_cast<int>(encoding)].get();
363393
} else {
364394
std::unique_ptr<Decoder> page_decoder;
365395
RETURN_IF_ERROR(Decoder::get_decoder(_metadata.type, encoding, page_decoder));
366396
// Set type length
367397
page_decoder->set_type_length(_get_type_length());
368398
_decoders[static_cast<int>(encoding)] = std::move(page_decoder);
369-
_page_decoder = _decoders[static_cast<int>(encoding)].get();
399+
encoding_decoder = _decoders[static_cast<int>(encoding)].get();
400+
}
401+
_empty_value_section = _page_data.empty() && _max_def_level > 0;
402+
if (_empty_value_section) {
403+
// Nullable all-NULL pages legally contain only definition levels. Use a dedicated decoder
404+
// so a later non-NULL definition level cannot consume stale state from the previous page.
405+
if (_empty_value_decoder == nullptr) {
406+
_empty_value_decoder = std::make_unique<EmptyValueSectionDecoder>();
407+
}
408+
_page_decoder = _empty_value_decoder.get();
409+
} else {
410+
_page_decoder = encoding_decoder;
411+
RETURN_IF_ERROR(_page_decoder->set_data(&_page_data));
370412
}
371-
RETURN_IF_ERROR(_page_decoder->set_data(&_page_data));
372413

373414
_state = DATA_LOADED;
374415
return Status::OK();
@@ -546,13 +587,18 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::skip_values(size_t num_va
546587
return Status::IOError("Skip too many values in current page. {} vs. {}",
547588
_remaining_num_values, num_values);
548589
}
549-
_remaining_num_values -= num_values;
550590
if (skip_data) {
591+
if (UNLIKELY(_empty_value_section && num_values != 0)) {
592+
return Status::Corruption(
593+
"Parquet definition levels require {} values from an empty value section",
594+
num_values);
595+
}
551596
SCOPED_RAW_TIMER(&_chunk_statistics.decode_value_time);
552-
return _page_decoder->skip_values(num_values);
553-
} else {
554-
return Status::OK();
597+
RETURN_IF_ERROR(_page_decoder->skip_values(num_values));
555598
}
599+
// Commit logical page progress only after the physical decoder accepted the whole request.
600+
_remaining_num_values -= num_values;
601+
return Status::OK();
556602
}
557603

558604
template <bool IN_COLLECTION, bool OFFSET_INDEX>
@@ -569,8 +615,16 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::decode_values(
569615
if (UNLIKELY(_remaining_num_values < select_vector.num_values())) {
570616
return Status::IOError("Decode too many values in current page");
571617
}
618+
const size_t physical_values = select_vector.num_values() - select_vector.num_nulls();
619+
if (UNLIKELY(_empty_value_section && physical_values != 0)) {
620+
return Status::Corruption(
621+
"Parquet definition levels require {} values from an empty value section",
622+
physical_values);
623+
}
624+
RETURN_IF_ERROR(
625+
_page_decoder->decode_values(doris_column, data_type, select_vector, is_dict_filter));
572626
_remaining_num_values -= select_vector.num_values();
573-
return _page_decoder->decode_values(doris_column, data_type, select_vector, is_dict_filter);
627+
return Status::OK();
574628
}
575629

576630
template <bool IN_COLLECTION, bool OFFSET_INDEX>

be/src/format/parquet/vparquet_column_chunk_reader.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -277,7 +277,9 @@ class ColumnChunkReader {
277277
Slice _v2_def_levels;
278278
bool _dict_checked = false;
279279
bool _has_dict = false;
280+
bool _empty_value_section = false;
280281
Decoder* _page_decoder = nullptr;
282+
std::unique_ptr<Decoder> _empty_value_decoder;
281283
// Map: encoding -> Decoder
282284
// Plain or Dictionary encoding. If the dictionary grows too big, the encoding will fall back to the plain encoding
283285
std::unordered_map<int, std::unique_ptr<Decoder>> _decoders;

be/test/format/parquet/parquet_column_chunk_reader_test.cpp

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,17 @@
2727

2828
#include "core/assert_cast.h"
2929
#include "core/column/column_string.h"
30+
#include "core/column/column_vector.h"
31+
#include "core/data_type/data_type_number.h"
3032
#include "format/parquet/schema_desc.h"
3133
#include "format/parquet/vparquet_column_chunk_reader.h"
3234
#include "format/parquet/vparquet_column_reader.h"
3335
#include "io/fs/buffered_reader.h"
3436
#include "io/fs/file_reader.h"
3537
#include "runtime/runtime_state.h"
38+
#include "util/block_compression.h"
3639
#include "util/coding.h"
40+
#include "util/faststring.h"
3741
#include "util/thrift_util.h"
3842

3943
namespace doris {
@@ -251,6 +255,52 @@ Status make_plain_fixture(ColumnChunkFixture* fixture, int page_count = 1) {
251255
return Status::OK();
252256
}
253257

258+
Status make_empty_boolean_v2_fixture(bool is_null, ColumnChunkFixture* fixture) {
259+
BlockCompressionCodec* codec = nullptr;
260+
RETURN_IF_ERROR(get_block_compression_codec(segment_v2::CompressionTypePB::SNAPPY, &codec));
261+
const uint8_t unused = 0;
262+
faststring compressed_values;
263+
RETURN_IF_ERROR(codec->compress(Slice(&unused, 0), &compressed_values));
264+
265+
const std::vector<uint8_t> definition_levels {2, static_cast<uint8_t>(is_null ? 0 : 1)};
266+
std::vector<uint8_t> payload = definition_levels;
267+
if (compressed_values.size() != 0) {
268+
payload.insert(payload.end(), compressed_values.data(),
269+
compressed_values.data() + compressed_values.size());
270+
}
271+
272+
tparquet::DataPageHeaderV2 data_header;
273+
data_header.__set_num_values(1);
274+
data_header.__set_num_nulls(is_null ? 1 : 0);
275+
data_header.__set_num_rows(1);
276+
data_header.__set_encoding(tparquet::Encoding::PLAIN);
277+
data_header.__set_definition_levels_byte_length(definition_levels.size());
278+
data_header.__set_repetition_levels_byte_length(0);
279+
data_header.__set_is_compressed(true);
280+
281+
tparquet::PageHeader header;
282+
header.type = tparquet::PageType::DATA_PAGE_V2;
283+
header.__set_compressed_page_size(cast_set<int32_t>(payload.size()));
284+
header.__set_uncompressed_page_size(cast_set<int32_t>(definition_levels.size()));
285+
header.__set_data_page_header_v2(data_header);
286+
287+
int64_t data_page_offset = 0;
288+
int32_t data_page_size = 0;
289+
RETURN_IF_ERROR(
290+
append_page(&header, payload, fixture->data, &data_page_offset, &data_page_size));
291+
292+
auto& metadata = fixture->chunk.meta_data;
293+
metadata.__set_type(tparquet::Type::BOOLEAN);
294+
metadata.__set_codec(tparquet::CompressionCodec::SNAPPY);
295+
metadata.__set_num_values(1);
296+
metadata.__set_data_page_offset(data_page_offset);
297+
metadata.__set_total_compressed_size(data_page_size);
298+
299+
fixture->field_schema.physical_type = tparquet::Type::BOOLEAN;
300+
fixture->field_schema.definition_level = 1;
301+
return Status::OK();
302+
}
303+
254304
void expect_dictionary_values(ColumnChunkReader<false, false>* reader) {
255305
MutableColumnPtr column = ColumnString::create();
256306
ASSERT_TRUE(reader->read_dict_values_to_column(column).ok());
@@ -386,6 +436,56 @@ TEST(ParquetColumnChunkReaderTest, FailedDictionaryCheckCanBeRetried) {
386436
}
387437
}
388438

439+
TEST(ParquetColumnChunkReaderTest, CompressedV2AllNullBooleanAcceptsEmptyValueSection) {
440+
ColumnChunkFixture fixture;
441+
ASSERT_TRUE(make_empty_boolean_v2_fixture(true, &fixture).ok());
442+
CountingBufferedReader buffered_reader(std::move(fixture.data));
443+
ParquetPageReadContext page_read_ctx(false);
444+
ColumnChunkReader<false, false> reader(&buffered_reader, &fixture.chunk, &fixture.field_schema,
445+
nullptr, 1, nullptr, page_read_ctx);
446+
447+
ASSERT_TRUE(reader.init().ok());
448+
ASSERT_TRUE(reader.parse_page_header().ok());
449+
ASSERT_TRUE(reader.load_page_data().ok());
450+
EXPECT_TRUE(reader.get_page_data().empty());
451+
452+
FilterMap filter_map;
453+
ASSERT_TRUE(filter_map.init(nullptr, 0, false).ok());
454+
ColumnSelectVector select_vector;
455+
const std::vector<uint16_t> null_map {0, 1};
456+
ASSERT_TRUE(select_vector.init(null_map, 1, nullptr, &filter_map, 0).ok());
457+
MutableColumnPtr column = ColumnUInt8::create();
458+
DataTypePtr data_type = std::make_shared<DataTypeUInt8>();
459+
ASSERT_TRUE(reader.decode_values(column, data_type, select_vector, false).ok());
460+
ASSERT_EQ(column->size(), 1);
461+
EXPECT_EQ(assert_cast<const ColumnUInt8&>(*column).get_data()[0], 0);
462+
}
463+
464+
TEST(ParquetColumnChunkReaderTest, CompressedV2NonNullBooleanRejectsEmptyValueSection) {
465+
ColumnChunkFixture fixture;
466+
ASSERT_TRUE(make_empty_boolean_v2_fixture(false, &fixture).ok());
467+
CountingBufferedReader buffered_reader(std::move(fixture.data));
468+
ParquetPageReadContext page_read_ctx(false);
469+
ColumnChunkReader<false, false> reader(&buffered_reader, &fixture.chunk, &fixture.field_schema,
470+
nullptr, 1, nullptr, page_read_ctx);
471+
472+
ASSERT_TRUE(reader.init().ok());
473+
ASSERT_TRUE(reader.parse_page_header().ok());
474+
ASSERT_TRUE(reader.load_page_data().ok());
475+
476+
FilterMap filter_map;
477+
ASSERT_TRUE(filter_map.init(nullptr, 0, false).ok());
478+
ColumnSelectVector select_vector;
479+
const std::vector<uint16_t> null_map {1};
480+
ASSERT_TRUE(select_vector.init(null_map, 1, nullptr, &filter_map, 0).ok());
481+
MutableColumnPtr column = ColumnUInt8::create();
482+
DataTypePtr data_type = std::make_shared<DataTypeUInt8>();
483+
EXPECT_TRUE(reader.decode_values(column, data_type, select_vector, false)
484+
.is<ErrorCode::CORRUPTION>());
485+
EXPECT_EQ(reader.remaining_num_values(), 1);
486+
EXPECT_TRUE(column->empty());
487+
}
488+
389489
TEST(ParquetColumnChunkReaderTest, ScalarDictionaryReadUsesExplicitProbe) {
390490
ColumnChunkFixture fixture;
391491
ASSERT_TRUE(make_dictionary_fixture(&fixture).ok());

0 commit comments

Comments
 (0)