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
5 changes: 3 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,6 +422,7 @@ 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() = default;

std::vector<Chunk> ContentDefinedChunker::GetChunks(const int16_t* def_levels,
Expand Down
9 changes: 5 additions & 4 deletions cpp/src/parquet/chunker_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -72,10 +72,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,6 +118,7 @@ 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();

/// Get the chunk boundaries for the given column data
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
59 changes: 35 additions & 24 deletions cpp/src/parquet/column_writer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -744,7 +744,8 @@ 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,
internal::ContentDefinedChunker* content_defined_chunker)
: metadata_(metadata),
descr_(metadata->descr()),
level_info_(internal::LevelInfo::ComputeLevelInfo(metadata->descr())),
Expand All @@ -763,7 +764,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 +777,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 requires a content defined chunker, use "
"ParquetFileWriter instead.");
}
}

Expand Down Expand Up @@ -912,7 +914,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 +1290,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,
internal::ContentDefinedChunker* content_defined_chunker)
: ColumnWriterImpl(metadata, std::move(pager), use_dictionary, encoding, properties,
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 +1449,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 +2704,10 @@ 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,
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 +2729,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, 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,
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,
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,
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,
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,
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,
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,
content_defined_chunker);
default:
ParquetException::NYI("Column writer not implemented for type: " +
TypeToString(descr->physical_type()));
Expand Down
12 changes: 8 additions & 4 deletions cpp/src/parquet/column_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,10 @@ class Encryptor;
class OffsetIndexBuilder;
class WriterProperties;

namespace internal {
class ContentDefinedChunker;
} // namespace internal

class PARQUET_EXPORT LevelEncoder {
public:
LevelEncoder();
Expand Down Expand Up @@ -127,10 +131,10 @@ class PARQUET_EXPORT ColumnWriter {
public:
virtual ~ColumnWriter() = default;

static std::shared_ptr<ColumnWriter> Make(ColumnChunkMetaDataBuilder*,
std::unique_ptr<PageWriter>,
const WriterProperties* properties,
BloomFilter* bloom_filter = NULLPTR);
static std::shared_ptr<ColumnWriter> Make(
ColumnChunkMetaDataBuilder*, std::unique_ptr<PageWriter>,
const WriterProperties* properties, BloomFilter* bloom_filter = NULLPTR,
internal::ContentDefinedChunker* content_defined_chunker = NULLPTR);
Comment on lines +134 to +137

/// \brief Closes the ColumnWriter, commits any buffered values to pages.
/// \return Total size of the column in bytes
Expand Down
36 changes: 27 additions & 9 deletions cpp/src/parquet/file_writer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
#include "arrow/util/key_value_metadata.h"
#include "arrow/util/logging_internal.h"
#include "parquet/bloom_filter_writer.h"
#include "parquet/chunker_internal.h"
#include "parquet/column_writer.h"
#include "parquet/encryption/encryption_internal.h"
#include "parquet/encryption/internal_file_encryptor.h"
Expand Down Expand Up @@ -93,12 +94,12 @@ inline void ThrowRowsMisMatchError(int col, int64_t prev, int64_t curr) {
// RowGroupWriter::Contents implementation for the Parquet file specification
class RowGroupSerializer : public RowGroupWriter::Contents {
public:
RowGroupSerializer(std::shared_ptr<ArrowOutputStream> sink,
RowGroupMetaDataBuilder* metadata, int16_t row_group_ordinal,
const WriterProperties* properties, bool buffered_row_group = false,
InternalFileEncryptor* file_encryptor = nullptr,
PageIndexBuilder* page_index_builder = nullptr,
BloomFilterBuilder* bloom_filter_builder = nullptr)
RowGroupSerializer(
std::shared_ptr<ArrowOutputStream> sink, RowGroupMetaDataBuilder* metadata,
int16_t row_group_ordinal, const WriterProperties* properties,
bool buffered_row_group, InternalFileEncryptor* file_encryptor,
PageIndexBuilder* page_index_builder, BloomFilterBuilder* bloom_filter_builder,
std::vector<internal::ContentDefinedChunker>& content_defined_chunkers)
: sink_(std::move(sink)),
metadata_(metadata),
properties_(properties),
Expand All @@ -111,7 +112,8 @@ class RowGroupSerializer : public RowGroupWriter::Contents {
buffered_row_group_(buffered_row_group),
file_encryptor_(file_encryptor),
page_index_builder_(page_index_builder),
bloom_filter_builder_(bloom_filter_builder) {
bloom_filter_builder_(bloom_filter_builder),
content_defined_chunkers_(content_defined_chunkers) {
if (buffered_row_group) {
InitColumns();
} else {
Expand Down Expand Up @@ -255,6 +257,7 @@ class RowGroupSerializer : public RowGroupWriter::Contents {
InternalFileEncryptor* file_encryptor_;
PageIndexBuilder* page_index_builder_;
BloomFilterBuilder* bloom_filter_builder_;
std::vector<internal::ContentDefinedChunker>& content_defined_chunkers_;

void CheckRowsWritten() const {
// verify when only one column is written at a time
Expand Down Expand Up @@ -308,6 +311,10 @@ class RowGroupSerializer : public RowGroupWriter::Contents {
if (bloom_filter_builder_) {
bloom_filter = bloom_filter_builder_->CreateBloomFilter(column_ordinal);
}
internal::ContentDefinedChunker* content_defined_chunker = nullptr;
if (properties_->content_defined_chunking_enabled()) {
content_defined_chunker = &content_defined_chunkers_[column_ordinal];
}
const CodecOptions* codec_options = column_properties.codec_options()
? column_properties.codec_options().get()
: nullptr;
Expand All @@ -321,7 +328,8 @@ class RowGroupSerializer : public RowGroupWriter::Contents {
static_cast<int16_t>(column_ordinal), properties_->memory_pool(),
buffered_row_group_, meta_encryptor, data_encryptor,
properties_->page_checksum_enabled(), ci_builder, oi_builder, *codec_options);
return ColumnWriter::Make(col_meta, std::move(pager), properties_, bloom_filter);
return ColumnWriter::Make(col_meta, std::move(pager), properties_, bloom_filter,
content_defined_chunker);
}

// If buffered_row_group_ is false, only column_writers_[0] is used as current writer.
Expand Down Expand Up @@ -415,7 +423,8 @@ class FileSerializer : public ParquetFileWriter::Contents {
}
std::unique_ptr<RowGroupWriter::Contents> contents(new RowGroupSerializer(
sink_, rg_metadata, row_group_ordinal, properties_.get(), buffered_row_group,
file_encryptor_.get(), page_index_builder_.get(), bloom_filter_builder_.get()));
file_encryptor_.get(), page_index_builder_.get(), bloom_filter_builder_.get(),
content_defined_chunkers_));
row_group_writer_ = std::make_unique<RowGroupWriter>(std::move(contents));
return row_group_writer_.get();
}
Expand Down Expand Up @@ -458,6 +467,14 @@ class FileSerializer : public ParquetFileWriter::Contents {
} else {
throw ParquetException("Appending to file not implemented.");
}
if (properties_->content_defined_chunking_enabled()) {
const auto& options = properties_->content_defined_chunking_options();
for (int i = 0; i < num_columns(); i++) {
content_defined_chunkers_.emplace_back(
internal::LevelInfo::ComputeLevelInfo(schema_.Column(i)),
options.min_chunk_size, options.max_chunk_size, options.norm_level);
}
}
}

void CloseEncryptedFile(FileEncryptionProperties* file_encryption_properties) {
Expand Down Expand Up @@ -524,6 +541,7 @@ class FileSerializer : public ParquetFileWriter::Contents {
std::unique_ptr<PageIndexBuilder> page_index_builder_;
std::unique_ptr<InternalFileEncryptor> file_encryptor_;
std::unique_ptr<BloomFilterBuilder> bloom_filter_builder_;
std::vector<internal::ContentDefinedChunker> content_defined_chunkers_;

void StartFile() {
auto file_encryption_properties = properties_->file_encryption_properties();
Expand Down
Loading