From 7c20ecd33d4c5fb4eccc460fc5ccc095de628df0 Mon Sep 17 00:00:00 2001 From: Krisztian Szucs Date: Thu, 1 Oct 2026 12:27:36 +0200 Subject: [PATCH 1/2] [C++][Parquet] Test that CDC pages don't depend on row group boundaries The column writer owns the content defined chunker and is recreated for every row group, so the chunking state is reset at each row group boundary. Writing the same table into one and into multiple row groups must give the same page boundaries apart from the row group ends. --- cpp/src/parquet/chunker_internal_test.cc | 49 ++++++++++++++++++++++++ 1 file changed, 49 insertions(+) diff --git a/cpp/src/parquet/chunker_internal_test.cc b/cpp/src/parquet/chunker_internal_test.cc index 2469d54afbf1..6700ffae3ce6 100644 --- a/cpp/src/parquet/chunker_internal_test.cc +++ b/cpp/src/parquet/chunker_internal_test.cc @@ -403,6 +403,20 @@ ParquetInfo GetColumnParquetInfo(const std::shared_ptr& 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 GetPageEnds(const std::shared_ptr& data, int column_index) { + std::vector 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; @@ -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(single))->num_row_groups(), 1); + ASSERT_EQ(ReadMetaData(std::make_shared(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 From d5b7b42890b2e2edec33f5c0e64ee4c7ea81714a Mon Sep 17 00:00:00 2001 From: Krisztian Szucs Date: Thu, 1 Oct 2026 20:10:46 +0200 Subject: [PATCH 2/2] [C++][Parquet] Keep the CDC chunking state across row groups The column writer created the chunker for every column chunk, so the chunking state was reset at each row group. The file writer now keeps one chunker per leaf column and passes it to the column writers of all the row groups. --- cpp/src/parquet/chunker_internal.cc | 13 +++++- cpp/src/parquet/chunker_internal.h | 18 ++++++-- cpp/src/parquet/column_writer.cc | 72 +++++++++++++++++++---------- cpp/src/parquet/column_writer.h | 21 +++++++++ cpp/src/parquet/file_writer.cc | 43 +++++++++++++---- 5 files changed, 126 insertions(+), 41 deletions(-) diff --git a/cpp/src/parquet/chunker_internal.cc b/cpp/src/parquet/chunker_internal.cc index 794075b73379..fa80948f8d73 100644 --- a/cpp/src/parquet/chunker_internal.cc +++ b/cpp/src/parquet/chunker_internal.cc @@ -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. @@ -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 ContentDefinedChunker::GetChunks(const int16_t* def_levels, const int16_t* rep_levels, int64_t num_levels, diff --git a/cpp/src/parquet/chunker_internal.h b/cpp/src/parquet/chunker_internal.h index 070b5f6c0b21..b926f8905e8b 100644 --- a/cpp/src/parquet/chunker_internal.h +++ b/cpp/src/parquet/chunker_internal.h @@ -22,6 +22,7 @@ #include "arrow/array.h" #include "parquet/level_conversion.h" +#include "parquet/properties.h" namespace parquet::internal { @@ -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 @@ -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 diff --git a/cpp/src/parquet/column_writer.cc b/cpp/src/parquet/column_writer.cc index 3296af62f0c4..d55440c1bf81 100644 --- a/cpp/src/parquet/column_writer.cc +++ b/cpp/src/parquet/column_writer.cc @@ -744,10 +744,12 @@ class ColumnWriterImpl { public: ColumnWriterImpl(ColumnChunkMetaDataBuilder* metadata, std::unique_ptr 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), @@ -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(AllocateBuffer(allocator_, 0)); repetition_levels_rle_ = @@ -775,11 +778,11 @@ class ColumnWriterImpl { compressor_temp_buffer_ = std::static_pointer_cast(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."); } } @@ -912,7 +915,9 @@ class ColumnWriterImpl { std::vector> data_pages_; - std::optional 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() { @@ -1286,9 +1291,10 @@ class TypedColumnWriterImpl : public ColumnWriterImpl, TypedColumnWriterImpl(ColumnChunkMetaDataBuilder* metadata, std::unique_ptr 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 @@ -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++) { @@ -2700,10 +2705,11 @@ Status TypedColumnWriterImpl::WriteArrowDense( // ---------------------------------------------------------------------- // Dynamic column writer constructor -std::shared_ptr ColumnWriter::Make(ColumnChunkMetaDataBuilder* metadata, - std::unique_ptr pager, - const WriterProperties* properties, - BloomFilter* bloom_filter) { +std::shared_ptr ColumnWriter::Make( + ColumnChunkMetaDataBuilder* metadata, std::unique_ptr 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; @@ -2725,29 +2731,36 @@ std::shared_ptr ColumnWriter::Make(ColumnChunkMetaDataBuilder* met } return std::make_shared>( 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>( - 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>( - 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>( - 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>( - 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>( - 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>( - 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>( - 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())); @@ -2756,4 +2769,13 @@ std::shared_ptr ColumnWriter::Make(ColumnChunkMetaDataBuilder* met return std::shared_ptr(nullptr); } +std::shared_ptr ColumnWriter::Make(ColumnChunkMetaDataBuilder* metadata, + std::unique_ptr 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); +} + } // namespace parquet diff --git a/cpp/src/parquet/column_writer.h b/cpp/src/parquet/column_writer.h index 5ad58c5ecf21..3385a221385a 100644 --- a/cpp/src/parquet/column_writer.h +++ b/cpp/src/parquet/column_writer.h @@ -54,6 +54,11 @@ class Encryptor; class OffsetIndexBuilder; class WriterProperties; +namespace internal { +class ContentDefinedChunker; +struct LevelInfo; +} // namespace internal + class PARQUET_EXPORT LevelEncoder { public: LevelEncoder(); @@ -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 Make( + ColumnChunkMetaDataBuilder*, std::unique_ptr, + 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. diff --git a/cpp/src/parquet/file_writer.cc b/cpp/src/parquet/file_writer.cc index ec303408f363..4436b00b15b2 100644 --- a/cpp/src/parquet/file_writer.cc +++ b/cpp/src/parquet/file_writer.cc @@ -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" @@ -93,12 +94,13 @@ 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 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 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, + const std::vector& level_infos, + std::vector& content_defined_chunkers) : sink_(std::move(sink)), metadata_(metadata), properties_(properties), @@ -111,7 +113,9 @@ 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), + level_infos_(level_infos), + content_defined_chunkers_(content_defined_chunkers) { if (buffered_row_group) { InitColumns(); } else { @@ -139,6 +143,9 @@ class RowGroupSerializer : public RowGroupWriter::Contents { // Throws an error if more columns are being written auto col_meta = metadata_->NextColumnChunk(); + // Keep the ordinal in step with the metadata even if closing the previous column + // writer throws + const int32_t column_ordinal = next_column_index_++; if (column_writers_[0]) { total_bytes_written_ += column_writers_[0]->Close(); @@ -146,7 +153,6 @@ class RowGroupSerializer : public RowGroupWriter::Contents { column_writers_[0]->total_compressed_bytes_written(); } - const int32_t column_ordinal = next_column_index_++; column_writers_[0] = CreateColumnWriterForColumn(col_meta, column_ordinal); return column_writers_[0].get(); } @@ -255,6 +261,8 @@ class RowGroupSerializer : public RowGroupWriter::Contents { InternalFileEncryptor* file_encryptor_; PageIndexBuilder* page_index_builder_; BloomFilterBuilder* bloom_filter_builder_; + const std::vector& level_infos_; + std::vector& content_defined_chunkers_; void CheckRowsWritten() const { // verify when only one column is written at a time @@ -308,6 +316,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; @@ -321,7 +333,8 @@ class RowGroupSerializer : public RowGroupWriter::Contents { static_cast(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, + level_infos_[column_ordinal], content_defined_chunker); } // If buffered_row_group_ is false, only column_writers_[0] is used as current writer. @@ -415,7 +428,8 @@ class FileSerializer : public ParquetFileWriter::Contents { } std::unique_ptr 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(), + level_infos_, content_defined_chunkers_)); row_group_writer_ = std::make_unique(std::move(contents)); return row_group_writer_.get(); } @@ -458,6 +472,13 @@ class FileSerializer : public ParquetFileWriter::Contents { } else { throw ParquetException("Appending to file not implemented."); } + for (int i = 0; i < num_columns(); i++) { + level_infos_.push_back(internal::LevelInfo::ComputeLevelInfo(schema_.Column(i))); + if (properties_->content_defined_chunking_enabled()) { + content_defined_chunkers_.push_back(internal::ContentDefinedChunker::Make( + level_infos_[i], properties_->content_defined_chunking_options())); + } + } } void CloseEncryptedFile(FileEncryptionProperties* file_encryption_properties) { @@ -524,6 +545,8 @@ class FileSerializer : public ParquetFileWriter::Contents { std::unique_ptr page_index_builder_; std::unique_ptr file_encryptor_; std::unique_ptr bloom_filter_builder_; + std::vector level_infos_; + std::vector content_defined_chunkers_; void StartFile() { auto file_encryption_properties = properties_->file_encryption_properties();