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
13 changes: 11 additions & 2 deletions cpp/src/parquet/chunker_internal.cc
Original file line number Diff line number Diff line change
Expand Up @@ -392,8 +392,8 @@ class ContentDefinedChunker::Impl {
}

private:
// Reference to the column's level information
const internal::LevelInfo& level_info_;
// The column's level information
const internal::LevelInfo level_info_;
// Minimum chunk size in bytes, the rolling hash will not be updated until this size is
// reached for each chunk. Note that all data sent through the hash function is counted
// towards the chunk size, including definition and repetition levels.
Expand Down Expand Up @@ -422,8 +422,17 @@ ContentDefinedChunker::ContentDefinedChunker(const LevelInfo& level_info,
int64_t max_chunk_size, int norm_level)
: impl_(new Impl(level_info, min_chunk_size, max_chunk_size, norm_level)) {}

ContentDefinedChunker::ContentDefinedChunker(ContentDefinedChunker&&) noexcept = default;
ContentDefinedChunker& ContentDefinedChunker::operator=(
ContentDefinedChunker&&) noexcept = default;
ContentDefinedChunker::~ContentDefinedChunker() = default;

ContentDefinedChunker ContentDefinedChunker::Make(const LevelInfo& level_info,
const CdcOptions& options) {
return ContentDefinedChunker(level_info, options.min_chunk_size, options.max_chunk_size,
options.norm_level);
}

std::vector<Chunk> ContentDefinedChunker::GetChunks(const int16_t* def_levels,
const int16_t* rep_levels,
int64_t num_levels,
Expand Down
18 changes: 14 additions & 4 deletions cpp/src/parquet/chunker_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@

#include "arrow/array.h"
#include "parquet/level_conversion.h"
#include "parquet/properties.h"

namespace parquet::internal {

Expand Down Expand Up @@ -72,10 +73,10 @@ struct Chunk {
/// Implementation details:
///
/// Only the parquet writer must be aware of the content defined chunking, the reader
/// doesn't need to know about it. Each parquet column writer holds a
/// ContentDefinedChunker instance depending on the writer's properties. The chunker's
/// state is maintained across the entire column without being reset between pages and row
/// groups.
/// doesn't need to know about it. The parquet file writer holds one
/// ContentDefinedChunker per leaf column depending on the writer's properties, and passes
/// it to the column writers of every row group. The chunker's state is maintained
/// across the entire column without being reset between pages and row groups.
///
/// The chunker receives the record shredded column data (def_levels, rep_levels, values)
/// and goes over the (def_level, rep_level, value) triplets one by one while adjusting
Expand Down Expand Up @@ -118,8 +119,17 @@ class PARQUET_EXPORT ContentDefinedChunker {
/// expense of fragmentation.
ContentDefinedChunker(const LevelInfo& level_info, int64_t min_chunk_size,
int64_t max_chunk_size, int norm_level = 0);
ContentDefinedChunker(ContentDefinedChunker&&) noexcept;
ContentDefinedChunker& operator=(ContentDefinedChunker&&) noexcept;
~ContentDefinedChunker();

/// Create a new ContentDefinedChunker instance using the given chunking options
///
/// @param level_info Information about definition and repetition levels
/// @param options Content defined chunking options
static ContentDefinedChunker Make(const LevelInfo& level_info,
const CdcOptions& options);

/// Get the chunk boundaries for the given column data
///
/// @param def_levels Definition levels
Expand Down
49 changes: 49 additions & 0 deletions cpp/src/parquet/chunker_internal_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -403,6 +403,20 @@ ParquetInfo GetColumnParquetInfo(const std::shared_ptr<Buffer>& data, int column
return result;
}

// The offsets where the data pages of a column end, counted in levels from the start of
// the column across its row groups, which are rows for flat columns
std::vector<int64_t> GetPageEnds(const std::shared_ptr<Buffer>& data, int column_index) {
std::vector<int64_t> page_ends;
int64_t offset = 0;
for (const auto& row_group : GetColumnParquetInfo(data, column_index)) {
for (auto page_length : row_group.page_lengths) {
offset += page_length;
page_ends.push_back(offset);
}
}
return page_ends;
}

// A git-hunk like side-by-side data structure to represent the differences between two
// vectors of uint64_t values.
using ChunkDiff = std::pair<ChunkList, ChunkList>;
Expand Down Expand Up @@ -1700,4 +1714,39 @@ TEST_F(TestCDCMultipleRowGroups, Append) {
}
}

TEST_F(TestCDCMultipleRowGroups, IndependentOfRowGroupBoundaries) {
// Splitting the data into row groups must not move the content defined page
// boundaries, only add one at the end of each row group. For example, if the pages of
// a single row group file end at rows 100, 250 and 400, then with row groups of 200
// rows the pages must end at rows 100, 200, 250 and 400.
ASSERT_OK_AND_ASSIGN(auto table, ConcatAndCombine({part1_, part2_, part3_}));
const int64_t num_rows = table->num_rows();
const int64_t row_group_length = num_rows / 6;
ASSERT_OK_AND_ASSIGN(auto single,
WriteTableToBuffer(table, kMinChunkSize, kMaxChunkSize,
/*row_group_length=*/num_rows));
ASSERT_OK_AND_ASSIGN(auto multi, WriteTableToBuffer(table, kMinChunkSize, kMaxChunkSize,
row_group_length));
ASSERT_EQ(ReadMetaData(std::make_shared<BufferReader>(single))->num_row_groups(), 1);
ASSERT_EQ(ReadMetaData(std::make_shared<BufferReader>(multi))->num_row_groups(), 6);
// compare the page ends column by column
for (int col = 0; col < table->num_columns(); col++) {
// the page ends of the single row group file, e.g. 100, 250 and 400
auto single_page_ends = GetPageEnds(single, col);
// the page ends of the multiple row group file, e.g. 100, 200, 250 and 400
auto multi_page_ends = GetPageEnds(multi, col);
// expect the page ends of the single row group file plus one at each of the 5
// boundaries between the row groups, e.g. 200, the last row group ends with the data
// where the single row group file's last page ends too
auto expected = single_page_ends;
for (int i = 1; i < 6; i++) {
expected.push_back(i * row_group_length);
}
std::sort(expected.begin(), expected.end());
EXPECT_EQ(multi_page_ends.size(), single_page_ends.size() + 5) << "column " << col;
// unlike ASSERT_EQ, ContainerEq prints the page ends that differ
ASSERT_THAT(multi_page_ends, ::testing::ContainerEq(expected)) << "column " << col;
}
}

} // namespace parquet::internal
72 changes: 47 additions & 25 deletions cpp/src/parquet/column_writer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -744,10 +744,12 @@ class ColumnWriterImpl {
public:
ColumnWriterImpl(ColumnChunkMetaDataBuilder* metadata,
std::unique_ptr<PageWriter> pager, const bool use_dictionary,
Encoding::type encoding, const WriterProperties* properties)
Encoding::type encoding, const WriterProperties* properties,
const internal::LevelInfo& level_info,
internal::ContentDefinedChunker* content_defined_chunker)
: metadata_(metadata),
descr_(metadata->descr()),
level_info_(internal::LevelInfo::ComputeLevelInfo(metadata->descr())),
level_info_(level_info),
pager_(std::move(pager)),
has_dictionary_(use_dictionary),
encoding_(encoding),
Expand All @@ -763,7 +765,8 @@ class ColumnWriterImpl {
closed_(false),
fallback_(false),
definition_levels_sink_(allocator_),
repetition_levels_sink_(allocator_) {
repetition_levels_sink_(allocator_),
content_defined_chunker_(content_defined_chunker) {
definition_levels_rle_ =
std::static_pointer_cast<ResizableBuffer>(AllocateBuffer(allocator_, 0));
repetition_levels_rle_ =
Expand All @@ -775,11 +778,11 @@ class ColumnWriterImpl {
compressor_temp_buffer_ =
std::static_pointer_cast<ResizableBuffer>(AllocateBuffer(allocator_, 0));
}
if (properties_->content_defined_chunking_enabled()) {
auto cdc_options = properties_->content_defined_chunking_options();
content_defined_chunker_.emplace(level_info_, cdc_options.min_chunk_size,
cdc_options.max_chunk_size,
cdc_options.norm_level);
if (properties_->content_defined_chunking_enabled() &&
content_defined_chunker_ == nullptr) {
throw ParquetException(
"Content-defined chunking is not supported in ColumnWriter::Make(), use "
"ParquetFileWriter instead.");
}
}

Expand Down Expand Up @@ -912,7 +915,9 @@ class ColumnWriterImpl {

std::vector<std::unique_ptr<DataPage>> data_pages_;

std::optional<internal::ContentDefinedChunker> content_defined_chunker_;
// The chunker of the column owned by the file writer, null without content defined
// chunking
internal::ContentDefinedChunker* content_defined_chunker_;

private:
void InitSinks() {
Expand Down Expand Up @@ -1286,9 +1291,10 @@ class TypedColumnWriterImpl : public ColumnWriterImpl,
TypedColumnWriterImpl(ColumnChunkMetaDataBuilder* metadata,
std::unique_ptr<PageWriter> pager, const bool use_dictionary,
Encoding::type encoding, const WriterProperties* properties,
BloomFilter* bloom_filter)
: ColumnWriterImpl(metadata, std::move(pager), use_dictionary, encoding,
properties) {
BloomFilter* bloom_filter, const internal::LevelInfo& level_info,
internal::ContentDefinedChunker* content_defined_chunker)
: ColumnWriterImpl(metadata, std::move(pager), use_dictionary, encoding, properties,
level_info, content_defined_chunker) {
current_encoder_ = MakeEncoder(ParquetType::type_num, encoding, use_dictionary,
descr_, properties->memory_pool());
// We have to dynamic_cast as some compilers don't want to static_cast
Expand Down Expand Up @@ -1444,7 +1450,6 @@ class TypedColumnWriterImpl : public ColumnWriterImpl,
}

if (ARROW_PREDICT_FALSE(properties_->content_defined_chunking_enabled())) {
DCHECK(content_defined_chunker_.has_value());
auto chunks = content_defined_chunker_->GetChunks(def_levels, rep_levels,
num_levels, leaf_array);
for (size_t i = 0; i < chunks.size(); i++) {
Expand Down Expand Up @@ -2700,10 +2705,11 @@ Status TypedColumnWriterImpl<FLBAType>::WriteArrowDense(
// ----------------------------------------------------------------------
// Dynamic column writer constructor

std::shared_ptr<ColumnWriter> ColumnWriter::Make(ColumnChunkMetaDataBuilder* metadata,
std::unique_ptr<PageWriter> pager,
const WriterProperties* properties,
BloomFilter* bloom_filter) {
std::shared_ptr<ColumnWriter> ColumnWriter::Make(
ColumnChunkMetaDataBuilder* metadata, std::unique_ptr<PageWriter> pager,
const WriterProperties* properties, BloomFilter* bloom_filter,
const internal::LevelInfo& level_info,
internal::ContentDefinedChunker* content_defined_chunker) {
const ColumnDescriptor* descr = metadata->descr();
const bool use_dictionary = properties->dictionary_enabled(descr->path()) &&
descr->physical_type() != Type::BOOLEAN;
Expand All @@ -2725,29 +2731,36 @@ std::shared_ptr<ColumnWriter> ColumnWriter::Make(ColumnChunkMetaDataBuilder* met
}
return std::make_shared<TypedColumnWriterImpl<BooleanType>>(
metadata, std::move(pager), use_dictionary, encoding, properties,
/*bloom_filter=*/nullptr);
/*bloom_filter=*/nullptr, level_info, content_defined_chunker);
}
case Type::INT32:
return std::make_shared<TypedColumnWriterImpl<Int32Type>>(
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter);
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter,
level_info, content_defined_chunker);
case Type::INT64:
return std::make_shared<TypedColumnWriterImpl<Int64Type>>(
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter);
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter,
level_info, content_defined_chunker);
case Type::INT96:
return std::make_shared<TypedColumnWriterImpl<Int96Type>>(
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter);
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter,
level_info, content_defined_chunker);
case Type::FLOAT:
return std::make_shared<TypedColumnWriterImpl<FloatType>>(
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter);
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter,
level_info, content_defined_chunker);
case Type::DOUBLE:
return std::make_shared<TypedColumnWriterImpl<DoubleType>>(
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter);
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter,
level_info, content_defined_chunker);
case Type::BYTE_ARRAY:
return std::make_shared<TypedColumnWriterImpl<ByteArrayType>>(
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter);
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter,
level_info, content_defined_chunker);
case Type::FIXED_LEN_BYTE_ARRAY:
return std::make_shared<TypedColumnWriterImpl<FLBAType>>(
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter);
metadata, std::move(pager), use_dictionary, encoding, properties, bloom_filter,
level_info, content_defined_chunker);
default:
ParquetException::NYI("Column writer not implemented for type: " +
TypeToString(descr->physical_type()));
Expand All @@ -2756,4 +2769,13 @@ std::shared_ptr<ColumnWriter> ColumnWriter::Make(ColumnChunkMetaDataBuilder* met
return std::shared_ptr<ColumnWriter>(nullptr);
}

std::shared_ptr<ColumnWriter> ColumnWriter::Make(ColumnChunkMetaDataBuilder* metadata,
std::unique_ptr<PageWriter> pager,
const WriterProperties* properties,
BloomFilter* bloom_filter) {
return Make(metadata, std::move(pager), properties, bloom_filter,
internal::LevelInfo::ComputeLevelInfo(metadata->descr()),
/*content_defined_chunker=*/nullptr);
Comment on lines +2776 to +2778
}

} // namespace parquet
21 changes: 21 additions & 0 deletions cpp/src/parquet/column_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,11 @@ class Encryptor;
class OffsetIndexBuilder;
class WriterProperties;

namespace internal {
class ContentDefinedChunker;
struct LevelInfo;
} // namespace internal

class PARQUET_EXPORT LevelEncoder {
public:
LevelEncoder();
Expand Down Expand Up @@ -199,6 +204,22 @@ class PARQUET_EXPORT ColumnWriter {
int64_t num_levels, const ::arrow::Array& leaf_array,
ArrowWriteContext* ctx,
bool leaf_field_nullable) = 0;

private:
friend class RowGroupSerializer;

/// \brief Create a column writer using the given level information and content
/// defined chunker
///
/// The file writer computes the level information of a column once and gives the same
/// chunker to the column writers of a column, so the content defined chunking is
/// carried over between the row groups. The chunker is required if the properties
/// enable content defined chunking.
static std::shared_ptr<ColumnWriter> Make(
ColumnChunkMetaDataBuilder*, std::unique_ptr<PageWriter>,
const WriterProperties* properties, BloomFilter* bloom_filter,
const internal::LevelInfo& level_info,
internal::ContentDefinedChunker* content_defined_chunker);
};

// API to write values to a single column. This is the main client facing API.
Expand Down
Loading
Loading