Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 51 additions & 8 deletions be/src/format/parquet/vparquet_column_chunk_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -358,17 +387,29 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::load_page_data() {
}

// Reuse page decoder
Decoder* encoding_decoder = nullptr;
if (_decoders.find(static_cast<int>(encoding)) != _decoders.end()) {
_page_decoder = _decoders[static_cast<int>(encoding)].get();
encoding_decoder = _decoders[static_cast<int>(encoding)].get();
} else {
std::unique_ptr<Decoder> 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<int>(encoding)] = std::move(page_decoder);
_page_decoder = _decoders[static_cast<int>(encoding)].get();
encoding_decoder = _decoders[static_cast<int>(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<EmptyValueSectionDecoder>();
}
_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();
Expand Down Expand Up @@ -546,13 +587,13 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::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 <bool IN_COLLECTION, bool OFFSET_INDEX>
Expand All @@ -569,8 +610,10 @@ Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::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 <bool IN_COLLECTION, bool OFFSET_INDEX>
Expand Down
2 changes: 2 additions & 0 deletions be/src/format/parquet/vparquet_column_chunk_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<Decoder> _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<int, std::unique_ptr<Decoder>> _decoders;
Expand Down
103 changes: 103 additions & 0 deletions be/test/format/parquet/parquet_column_chunk_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<uint8_t> definition_levels {2, static_cast<uint8_t>(is_null ? 0 : 1)};
std::vector<uint8_t> 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<int32_t>(payload.size()));
header.__set_uncompressed_page_size(cast_set<int32_t>(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<false, false>* reader) {
MutableColumnPtr column = ColumnString::create();
ASSERT_TRUE(reader->read_dict_values_to_column(column).ok());
Expand Down Expand Up @@ -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<false, false> 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<uint16_t> 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<DataTypeUInt8>();
ASSERT_TRUE(reader.decode_values(column, data_type, select_vector, false).ok());
ASSERT_EQ(column->size(), 1);
EXPECT_EQ(assert_cast<const ColumnUInt8&>(*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<false, false> 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<ErrorCode::CORRUPTION>());
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<uint16_t> 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<DataTypeUInt8>();
EXPECT_TRUE(reader.decode_values(column, data_type, select_vector, false)
.is<ErrorCode::CORRUPTION>());
EXPECT_EQ(reader.remaining_num_values(), 1);
EXPECT_TRUE(column->empty());
}

TEST(ParquetColumnChunkReaderTest, ScalarDictionaryReadUsesExplicitProbe) {
ColumnChunkFixture fixture;
ASSERT_TRUE(make_dictionary_fixture(&fixture).ok());
Expand Down
Loading