diff --git a/cpp/src/parquet/chunker_internal.cc b/cpp/src/parquet/chunker_internal.cc index 794075b7337..02658d361a5 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,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 ContentDefinedChunker::GetChunks(const int16_t* def_levels, diff --git a/cpp/src/parquet/chunker_internal.h b/cpp/src/parquet/chunker_internal.h index 070b5f6c0b2..52de91b04dd 100644 --- a/cpp/src/parquet/chunker_internal.h +++ b/cpp/src/parquet/chunker_internal.h @@ -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 @@ -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 diff --git a/cpp/src/parquet/chunker_internal_test.cc b/cpp/src/parquet/chunker_internal_test.cc index 2469d54afbf..6700ffae3ce 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 diff --git a/cpp/src/parquet/column_writer.cc b/cpp/src/parquet/column_writer.cc index 3296af62f0c..a5c7375835a 100644 --- a/cpp/src/parquet/column_writer.cc +++ b/cpp/src/parquet/column_writer.cc @@ -744,7 +744,8 @@ 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, + internal::ContentDefinedChunker* content_defined_chunker) : metadata_(metadata), descr_(metadata->descr()), level_info_(internal::LevelInfo::ComputeLevelInfo(metadata->descr())), @@ -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(AllocateBuffer(allocator_, 0)); repetition_levels_rle_ = @@ -775,11 +777,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 requires a content defined chunker, use " + "ParquetFileWriter instead."); } } @@ -912,7 +914,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 +1290,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, + 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 @@ -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++) { @@ -2700,10 +2704,10 @@ 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, + 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 +2729,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, 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, + 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, + 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, + 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, + 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, + 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, + 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, + content_defined_chunker); default: ParquetException::NYI("Column writer not implemented for type: " + TypeToString(descr->physical_type())); diff --git a/cpp/src/parquet/column_writer.h b/cpp/src/parquet/column_writer.h index 5ad58c5ecf2..069b09abaab 100644 --- a/cpp/src/parquet/column_writer.h +++ b/cpp/src/parquet/column_writer.h @@ -54,6 +54,10 @@ class Encryptor; class OffsetIndexBuilder; class WriterProperties; +namespace internal { +class ContentDefinedChunker; +} // namespace internal + class PARQUET_EXPORT LevelEncoder { public: LevelEncoder(); @@ -127,10 +131,10 @@ class PARQUET_EXPORT ColumnWriter { public: virtual ~ColumnWriter() = default; - static std::shared_ptr Make(ColumnChunkMetaDataBuilder*, - std::unique_ptr, - const WriterProperties* properties, - BloomFilter* bloom_filter = NULLPTR); + static std::shared_ptr Make( + ColumnChunkMetaDataBuilder*, std::unique_ptr, + const WriterProperties* properties, BloomFilter* bloom_filter = NULLPTR, + internal::ContentDefinedChunker* content_defined_chunker = NULLPTR); /// \brief Closes the ColumnWriter, commits any buffered values to pages. /// \return Total size of the column in bytes diff --git a/cpp/src/parquet/file_writer.cc b/cpp/src/parquet/file_writer.cc index ec303408f36..6d164881f66 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,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 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, + std::vector& content_defined_chunkers) : sink_(std::move(sink)), metadata_(metadata), properties_(properties), @@ -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 { @@ -255,6 +257,7 @@ class RowGroupSerializer : public RowGroupWriter::Contents { InternalFileEncryptor* file_encryptor_; PageIndexBuilder* page_index_builder_; BloomFilterBuilder* bloom_filter_builder_; + std::vector& content_defined_chunkers_; void CheckRowsWritten() const { // verify when only one column is written at a time @@ -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; @@ -321,7 +328,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, + content_defined_chunker); } // If buffered_row_group_ is false, only column_writers_[0] is used as current writer. @@ -415,7 +423,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(), + content_defined_chunkers_)); row_group_writer_ = std::make_unique(std::move(contents)); return row_group_writer_.get(); } @@ -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) { @@ -524,6 +541,7 @@ class FileSerializer : public ParquetFileWriter::Contents { std::unique_ptr page_index_builder_; std::unique_ptr file_encryptor_; std::unique_ptr bloom_filter_builder_; + std::vector content_defined_chunkers_; void StartFile() { auto file_encryption_properties = properties_->file_encryption_properties();