diff --git a/be/src/format/parquet/vparquet_column_chunk_reader.cpp b/be/src/format/parquet/vparquet_column_chunk_reader.cpp index ca9d290943adc6..5398aeb2e0a18b 100644 --- a/be/src/format/parquet/vparquet_column_chunk_reader.cpp +++ b/be/src/format/parquet/vparquet_column_chunk_reader.cpp @@ -43,6 +43,35 @@ namespace cctz { class time_zone; } // namespace cctz namespace doris { + +namespace { + +class EmptyValueSectionDecoder final : public Decoder { +public: + Status decode_values(MutableColumnPtr& doris_column, DataTypePtr&, + ColumnSelectVector& select_vector, bool) override { + const size_t physical_values = select_vector.num_values() - select_vector.num_nulls(); + if (UNLIKELY(physical_values != 0)) { + return Status::Corruption( + "Parquet definition levels require {} values from an empty value section", + physical_values); + } + doris_column->insert_many_defaults(select_vector.num_values() - + select_vector.num_filtered()); + return Status::OK(); + } + + Status skip_values(size_t num_values) override { + if (UNLIKELY(num_values != 0)) { + return Status::Corruption( + "Parquet definition levels require {} values from an empty value section", + num_values); + } + return Status::OK(); + } +}; + +} // namespace namespace io { class BufferedStreamReader; struct IOContext; @@ -358,17 +387,29 @@ Status ColumnChunkReader::load_page_data() { } // Reuse page decoder + Decoder* encoding_decoder = nullptr; if (_decoders.find(static_cast(encoding)) != _decoders.end()) { - _page_decoder = _decoders[static_cast(encoding)].get(); + encoding_decoder = _decoders[static_cast(encoding)].get(); } else { std::unique_ptr page_decoder; RETURN_IF_ERROR(Decoder::get_decoder(_metadata.type, encoding, page_decoder)); // Set type length page_decoder->set_type_length(_get_type_length()); _decoders[static_cast(encoding)] = std::move(page_decoder); - _page_decoder = _decoders[static_cast(encoding)].get(); + encoding_decoder = _decoders[static_cast(encoding)].get(); + } + _empty_value_section = _page_data.empty() && _max_def_level > 0; + if (_empty_value_section) { + // Nullable all-NULL pages legally contain only definition levels. Use a dedicated decoder + // so a later non-NULL definition level cannot consume stale state from the previous page. + if (_empty_value_decoder == nullptr) { + _empty_value_decoder = std::make_unique(); + } + _page_decoder = _empty_value_decoder.get(); + } else { + _page_decoder = encoding_decoder; + RETURN_IF_ERROR(_page_decoder->set_data(&_page_data)); } - RETURN_IF_ERROR(_page_decoder->set_data(&_page_data)); _state = DATA_LOADED; return Status::OK(); @@ -546,13 +587,13 @@ Status ColumnChunkReader::skip_values(size_t num_va return Status::IOError("Skip too many values in current page. {} vs. {}", _remaining_num_values, num_values); } - _remaining_num_values -= num_values; if (skip_data) { SCOPED_RAW_TIMER(&_chunk_statistics.decode_value_time); - return _page_decoder->skip_values(num_values); - } else { - return Status::OK(); + RETURN_IF_ERROR(_page_decoder->skip_values(num_values)); } + // Commit logical page progress only after the physical decoder accepted the whole request. + _remaining_num_values -= num_values; + return Status::OK(); } template @@ -569,8 +610,10 @@ Status ColumnChunkReader::decode_values( if (UNLIKELY(_remaining_num_values < select_vector.num_values())) { return Status::IOError("Decode too many values in current page"); } + RETURN_IF_ERROR( + _page_decoder->decode_values(doris_column, data_type, select_vector, is_dict_filter)); _remaining_num_values -= select_vector.num_values(); - return _page_decoder->decode_values(doris_column, data_type, select_vector, is_dict_filter); + return Status::OK(); } template diff --git a/be/src/format/parquet/vparquet_column_chunk_reader.h b/be/src/format/parquet/vparquet_column_chunk_reader.h index 064d28fc115941..7a1908d71ba730 100644 --- a/be/src/format/parquet/vparquet_column_chunk_reader.h +++ b/be/src/format/parquet/vparquet_column_chunk_reader.h @@ -277,7 +277,9 @@ class ColumnChunkReader { Slice _v2_def_levels; bool _dict_checked = false; bool _has_dict = false; + bool _empty_value_section = false; Decoder* _page_decoder = nullptr; + std::unique_ptr _empty_value_decoder; // Map: encoding -> Decoder // Plain or Dictionary encoding. If the dictionary grows too big, the encoding will fall back to the plain encoding std::unordered_map> _decoders; diff --git a/be/test/format/parquet/parquet_column_chunk_reader_test.cpp b/be/test/format/parquet/parquet_column_chunk_reader_test.cpp index be9616c523f638..f1326020f96671 100644 --- a/be/test/format/parquet/parquet_column_chunk_reader_test.cpp +++ b/be/test/format/parquet/parquet_column_chunk_reader_test.cpp @@ -27,13 +27,17 @@ #include "core/assert_cast.h" #include "core/column/column_string.h" +#include "core/column/column_vector.h" +#include "core/data_type/data_type_number.h" #include "format/parquet/schema_desc.h" #include "format/parquet/vparquet_column_chunk_reader.h" #include "format/parquet/vparquet_column_reader.h" #include "io/fs/buffered_reader.h" #include "io/fs/file_reader.h" #include "runtime/runtime_state.h" +#include "util/block_compression.h" #include "util/coding.h" +#include "util/faststring.h" #include "util/thrift_util.h" namespace doris { @@ -251,6 +255,52 @@ Status make_plain_fixture(ColumnChunkFixture* fixture, int page_count = 1) { return Status::OK(); } +Status make_empty_boolean_v2_fixture(bool is_null, ColumnChunkFixture* fixture) { + BlockCompressionCodec* codec = nullptr; + RETURN_IF_ERROR(get_block_compression_codec(segment_v2::CompressionTypePB::SNAPPY, &codec)); + const uint8_t unused = 0; + faststring compressed_values; + RETURN_IF_ERROR(codec->compress(Slice(&unused, 0), &compressed_values)); + + const std::vector definition_levels {2, static_cast(is_null ? 0 : 1)}; + std::vector payload = definition_levels; + if (compressed_values.size() != 0) { + payload.insert(payload.end(), compressed_values.data(), + compressed_values.data() + compressed_values.size()); + } + + tparquet::DataPageHeaderV2 data_header; + data_header.__set_num_values(1); + data_header.__set_num_nulls(is_null ? 1 : 0); + data_header.__set_num_rows(1); + data_header.__set_encoding(tparquet::Encoding::PLAIN); + data_header.__set_definition_levels_byte_length(definition_levels.size()); + data_header.__set_repetition_levels_byte_length(0); + data_header.__set_is_compressed(true); + + tparquet::PageHeader header; + header.type = tparquet::PageType::DATA_PAGE_V2; + header.__set_compressed_page_size(cast_set(payload.size())); + header.__set_uncompressed_page_size(cast_set(definition_levels.size())); + header.__set_data_page_header_v2(data_header); + + int64_t data_page_offset = 0; + int32_t data_page_size = 0; + RETURN_IF_ERROR( + append_page(&header, payload, fixture->data, &data_page_offset, &data_page_size)); + + auto& metadata = fixture->chunk.meta_data; + metadata.__set_type(tparquet::Type::BOOLEAN); + metadata.__set_codec(tparquet::CompressionCodec::SNAPPY); + metadata.__set_num_values(1); + metadata.__set_data_page_offset(data_page_offset); + metadata.__set_total_compressed_size(data_page_size); + + fixture->field_schema.physical_type = tparquet::Type::BOOLEAN; + fixture->field_schema.definition_level = 1; + return Status::OK(); +} + void expect_dictionary_values(ColumnChunkReader* reader) { MutableColumnPtr column = ColumnString::create(); ASSERT_TRUE(reader->read_dict_values_to_column(column).ok()); @@ -386,6 +436,59 @@ TEST(ParquetColumnChunkReaderTest, FailedDictionaryCheckCanBeRetried) { } } +TEST(ParquetColumnChunkReaderTest, CompressedV2AllNullBooleanAcceptsEmptyValueSection) { + ColumnChunkFixture fixture; + ASSERT_TRUE(make_empty_boolean_v2_fixture(true, &fixture).ok()); + CountingBufferedReader buffered_reader(std::move(fixture.data)); + ParquetPageReadContext page_read_ctx(false); + ColumnChunkReader reader(&buffered_reader, &fixture.chunk, &fixture.field_schema, + nullptr, 1, nullptr, page_read_ctx); + + ASSERT_TRUE(reader.init().ok()); + ASSERT_TRUE(reader.parse_page_header().ok()); + ASSERT_TRUE(reader.load_page_data().ok()); + EXPECT_TRUE(reader.get_page_data().empty()); + EXPECT_TRUE(reader.skip_values(0).ok()); + + FilterMap filter_map; + ASSERT_TRUE(filter_map.init(nullptr, 0, false).ok()); + ColumnSelectVector select_vector; + const std::vector null_map {0, 1}; + ASSERT_TRUE(select_vector.init(null_map, 1, nullptr, &filter_map, 0).ok()); + MutableColumnPtr column = ColumnUInt8::create(); + DataTypePtr data_type = std::make_shared(); + ASSERT_TRUE(reader.decode_values(column, data_type, select_vector, false).ok()); + ASSERT_EQ(column->size(), 1); + EXPECT_EQ(assert_cast(*column).get_data()[0], 0); +} + +TEST(ParquetColumnChunkReaderTest, CompressedV2NonNullBooleanRejectsEmptyValueSection) { + ColumnChunkFixture fixture; + ASSERT_TRUE(make_empty_boolean_v2_fixture(false, &fixture).ok()); + CountingBufferedReader buffered_reader(std::move(fixture.data)); + ParquetPageReadContext page_read_ctx(false); + ColumnChunkReader reader(&buffered_reader, &fixture.chunk, &fixture.field_schema, + nullptr, 1, nullptr, page_read_ctx); + + ASSERT_TRUE(reader.init().ok()); + ASSERT_TRUE(reader.parse_page_header().ok()); + ASSERT_TRUE(reader.load_page_data().ok()); + EXPECT_TRUE(reader.skip_values(1).is()); + EXPECT_EQ(reader.remaining_num_values(), 1); + + FilterMap filter_map; + ASSERT_TRUE(filter_map.init(nullptr, 0, false).ok()); + ColumnSelectVector select_vector; + const std::vector null_map {1}; + ASSERT_TRUE(select_vector.init(null_map, 1, nullptr, &filter_map, 0).ok()); + MutableColumnPtr column = ColumnUInt8::create(); + DataTypePtr data_type = std::make_shared(); + EXPECT_TRUE(reader.decode_values(column, data_type, select_vector, false) + .is()); + EXPECT_EQ(reader.remaining_num_values(), 1); + EXPECT_TRUE(column->empty()); +} + TEST(ParquetColumnChunkReaderTest, ScalarDictionaryReadUsesExplicitProbe) { ColumnChunkFixture fixture; ASSERT_TRUE(make_dictionary_fixture(&fixture).ok());