diff --git a/cpp/src/parquet/arrow/reader_writer_benchmark.cc b/cpp/src/parquet/arrow/reader_writer_benchmark.cc index 2f288fd2eb0f..a5e84243b7ff 100644 --- a/cpp/src/parquet/arrow/reader_writer_benchmark.cc +++ b/cpp/src/parquet/arrow/reader_writer_benchmark.cc @@ -273,6 +273,48 @@ BENCHMARK(BM_WriteBinaryColumn) ->Args({50, kInfiniteUniqueValues}) ->Args({99, kInfiniteUniqueValues}); +// Writes int32 columns with content-defined chunking into a single row group, with +// use_threads the columns are chunked by different threads in parallel +static void BM_WriteContentDefinedChunking(::benchmark::State& state) { + const auto num_columns = static_cast(state.range(0)); + const bool use_threads = state.range(1) != 0; + constexpr int64_t kNumRows = 1024 * 1024; + + ::arrow::random::RandomArrayGenerator generator(/*seed=*/500); + ::arrow::FieldVector fields; + ::arrow::ArrayVector columns; + for (int i = 0; i < num_columns; i++) { + fields.push_back(::arrow::field("column" + std::to_string(i), ::arrow::int32(), + /*nullable=*/false)); + columns.push_back(generator.Int32(kNumRows, /*min=*/0, /*max=*/1 << 30)); + } + auto schema = ::arrow::schema(fields); + auto batch = ::arrow::RecordBatch::Make(schema, kNumRows, columns); + auto properties = WriterProperties::Builder() + .enable_content_defined_chunking() + ->disable_dictionary() + ->build(); + auto arrow_properties = + ArrowWriterProperties::Builder().set_use_threads(use_threads)->build(); + + for (auto _ : state) { + auto output = CreateOutputStream(); + auto writer = arrow::FileWriter::Open(*schema, ::arrow::default_memory_pool(), output, + properties, arrow_properties); + EXIT_NOT_OK(writer.status()); + EXIT_NOT_OK((*writer)->WriteRecordBatch(*batch)); + EXIT_NOT_OK((*writer)->Close()); + } + state.SetBytesProcessed(state.iterations() * num_columns * kNumRows * + static_cast(sizeof(int32_t))); +} + +BENCHMARK(BM_WriteContentDefinedChunking) + ->ArgNames({"columns", "use_threads"}) + ->Args({16, 0}) + ->Args({16, 1}) + ->UseRealTime(); + template struct Examples { static constexpr std::array values() { return {127, 128}; } diff --git a/cpp/src/parquet/chunker_internal.cc b/cpp/src/parquet/chunker_internal.cc index 794075b73379..d8b3c0356f30 100644 --- a/cpp/src/parquet/chunker_internal.cc +++ b/cpp/src/parquet/chunker_internal.cc @@ -112,14 +112,11 @@ uint64_t CalculateMask(int64_t min_chunk_size, int64_t max_chunk_size, int norm_ } } -} // namespace - -class ContentDefinedChunker::Impl { +/// The gear hash of a column together with the state deciding the chunk boundaries +class GearHash { public: - Impl(const LevelInfo& level_info, int64_t min_chunk_size, int64_t max_chunk_size, - int norm_level) - : level_info_(level_info), - min_chunk_size_(min_chunk_size), + GearHash(int64_t min_chunk_size, int64_t max_chunk_size, int norm_level) + : min_chunk_size_(min_chunk_size), max_chunk_size_(max_chunk_size), rolling_hash_mask_(CalculateMask(min_chunk_size, max_chunk_size, norm_level)) {} @@ -201,6 +198,40 @@ class ContentDefinedChunker::Impl { return false; } + private: + // 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. + int64_t min_chunk_size_; + int64_t max_chunk_size_; + // The mask to match the rolling hash against to determine if a new chunk should be + // created. The mask is calculated based on min/max chunk size and the normalization + // level. + uint64_t rolling_hash_mask_; + + // Whether the rolling hash has matched the mask since the last chunk creation. This + // flag is set true by the Roll() function when the mask is matched and reset to false + // by NeedNewChunk() method. + bool has_matched_ = false; + // The current run of the rolling hash, used to normalize the chunk size distribution + // by requiring multiple consecutive matches to create a new chunk. + int8_t nth_run_ = 0; + // Current chunk size in bytes, reset to 0 when a new chunk is created. + int64_t chunk_size_ = 0; + // Rolling hash state, never reset only initialized once for the entire column. + uint64_t rolling_hash_ = 0; +}; + +} // namespace + +class ContentDefinedChunker::Impl { + public: + Impl(const LevelInfo& level_info, int64_t min_chunk_size, int64_t max_chunk_size, + int norm_level) + : level_info_(level_info), gearhash_(min_chunk_size, max_chunk_size, norm_level) {} + + uint64_t GetRollingHashMask() const { return gearhash_.GetRollingHashMask(); } + void ValidateChunks(const std::vector& chunks, int64_t num_levels) const { // chunks must be non-empty and monotonic increasing ARROW_DCHECK(!chunks.empty()); @@ -248,6 +279,9 @@ class ContentDefinedChunker::Impl { // requirements, we create a new chunk. See the `NeedNewChunk()` method for more // details. std::vector chunks; + // Roll a local copy of the hash so the compiler can keep it in registers instead of + // storing it to memory for every hashed byte + GearHash gearhash = gearhash_; int64_t offset; int64_t prev_offset = 0; int64_t prev_value_offset = 0; @@ -257,8 +291,8 @@ class ContentDefinedChunker::Impl { if (!has_rep_levels && !has_def_levels) { // fastest path for non-nested non-null data for (offset = 0; offset < num_levels; ++offset) { - RollValue(offset); - if (NeedNewChunk()) { + RollValue(gearhash, offset); + if (gearhash.NeedNewChunk()) { chunks.push_back({prev_offset, prev_offset, offset - prev_offset}); prev_offset = offset; } @@ -271,11 +305,11 @@ class ContentDefinedChunker::Impl { for (int64_t offset = 0; offset < num_levels; ++offset) { def_level = def_levels[offset]; - Roll(&def_level); + gearhash.Roll(&def_level); if (def_level == level_info_.def_level) { - RollValue(offset); + RollValue(gearhash, offset); } - if (NeedNewChunk()) { + if (gearhash.NeedNewChunk()) { chunks.push_back({prev_offset, prev_offset, offset - prev_offset}); prev_offset = offset; } @@ -292,13 +326,13 @@ class ContentDefinedChunker::Impl { def_level = def_levels[offset]; rep_level = rep_levels[offset]; - Roll(&def_level); - Roll(&rep_level); + gearhash.Roll(&def_level); + gearhash.Roll(&rep_level); if (def_level == level_info_.def_level) { - RollValue(value_offset); + RollValue(gearhash, value_offset); } - if (rep_level == 0 && NeedNewChunk()) { + if (rep_level == 0 && gearhash.NeedNewChunk()) { // if we are at a record boundary and need a new chunk, we create a new chunk auto levels_to_write = offset - prev_offset; if (levels_to_write > 0) { @@ -313,6 +347,7 @@ class ContentDefinedChunker::Impl { } } } + gearhash_ = gearhash; // add the last chunk if we have any levels left if (prev_offset < num_levels) { @@ -332,9 +367,10 @@ class ContentDefinedChunker::Impl { const uint8_t* raw_values = values.data()->GetValues(/*i=*/1, /*absolute_offset=*/0) + values.offset() * kByteWidth; - return Calculate(def_levels, rep_levels, num_levels, [&](int64_t i) { - return Roll(&raw_values[i * kByteWidth]); - }); + return Calculate(def_levels, rep_levels, num_levels, + [&](GearHash& gearhash, int64_t i) { + return gearhash.Roll(&raw_values[i * kByteWidth]); + }); } template @@ -342,11 +378,12 @@ class ContentDefinedChunker::Impl { const int16_t* rep_levels, int64_t num_levels, const ::arrow::Array& values) { const auto& array = checked_cast(values); - return Calculate(def_levels, rep_levels, num_levels, [&](int64_t i) { - typename ArrayType::offset_type length; - const uint8_t* value = array.GetValue(i, &length); - Roll(value, length); - }); + return Calculate(def_levels, rep_levels, num_levels, + [&](GearHash& gearhash, int64_t i) { + typename ArrayType::offset_type length; + const uint8_t* value = array.GetValue(i, &length); + gearhash.Roll(value, length); + }); } std::vector GetChunks(const int16_t* def_levels, const int16_t* rep_levels, @@ -354,16 +391,19 @@ class ContentDefinedChunker::Impl { auto handle_type = [&](auto&& type) -> std::vector { using ArrowType = std::decay_t; if constexpr (ArrowType::type_id == ::arrow::Type::NA) { - return Calculate(def_levels, rep_levels, num_levels, [](int64_t) {}); + return Calculate(def_levels, rep_levels, num_levels, [](GearHash&, int64_t) {}); } else if constexpr (ArrowType::type_id == ::arrow::Type::BOOL) { const auto& array = static_cast(values); - return Calculate(def_levels, rep_levels, num_levels, - [&](int64_t i) { Roll(array.Value(i)); }); + return Calculate( + def_levels, rep_levels, num_levels, + [&](GearHash& gearhash, int64_t i) { gearhash.Roll(array.Value(i)); }); } else if constexpr (ArrowType::type_id == ::arrow::Type::FIXED_SIZE_BINARY) { const auto& array = static_cast(values); const auto byte_width = array.byte_width(); return Calculate(def_levels, rep_levels, num_levels, - [&](int64_t i) { Roll(array.GetValue(i), byte_width); }); + [&](GearHash& gearhash, int64_t i) { + gearhash.Roll(array.GetValue(i), byte_width); + }); } else if constexpr (ArrowType::type_id == ::arrow::Type::EXTENSION) { const auto& array = static_cast(values); return GetChunks(def_levels, rep_levels, num_levels, *array.storage()); @@ -392,29 +432,11 @@ class ContentDefinedChunker::Impl { } private: - // Reference to 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. - const int64_t min_chunk_size_; - const int64_t max_chunk_size_; - // The mask to match the rolling hash against to determine if a new chunk should be - // created. The mask is calculated based on min/max chunk size and the normalization - // level. - const uint64_t rolling_hash_mask_; - - // Whether the rolling hash has matched the mask since the last chunk creation. This - // flag is set true by the Roll() function when the mask is matched and reset to false - // by NeedNewChunk() method. - bool has_matched_ = false; - // The current run of the rolling hash, used to normalize the chunk size distribution - // by requiring multiple consecutive matches to create a new chunk. - int8_t nth_run_ = 0; - // Current chunk size in bytes, reset to 0 when a new chunk is created. - int64_t chunk_size_ = 0; - // Rolling hash state, never reset only initialized once for the entire column. - uint64_t rolling_hash_ = 0; + // The column's level information + const internal::LevelInfo level_info_; + // The gear hash carried over between the calls, so the chunking continues across + // the pages and the row groups + GearHash gearhash_; }; ContentDefinedChunker::ContentDefinedChunker(const LevelInfo& level_info, @@ -422,8 +444,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/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 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();