This is an automated email from the ASF dual-hosted git repository.
maplefu pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow.git
The following commit(s) were added to refs/heads/main by this push:
new c58ec9cc31 GH-45366: [C++][Parquet] Set is_compressed to false when
data page v2 is not compressed (#45367)
c58ec9cc31 is described below
commit c58ec9cc31f65cef424e518f59283cb582a7adf4
Author: mwish <[email protected]>
AuthorDate: Wed Mar 5 18:11:54 2025 +0800
GH-45366: [C++][Parquet] Set is_compressed to false when data page v2 is
not compressed (#45367)
### Rationale for this change
Currently, if data page v2 is enabled, is_compressed is always set if Page
has compression, however, this can be eliminated if:
1. Page is empty
2. Compression makes page event larger.
### What changes are included in this PR?
Enhancement is_compressed setting in column_writer.cc
### Are these changes tested?
Yes
### Are there any user-facing changes?
No
* GitHub Issue: #45366
Authored-by: mwish <[email protected]>
Signed-off-by: mwish <[email protected]>
---
cpp/src/parquet/column_writer.cc | 16 ++++++------
cpp/src/parquet/column_writer_test.cc | 47 ++++++++++++++++++++++++++++++++++-
2 files changed, 55 insertions(+), 8 deletions(-)
diff --git a/cpp/src/parquet/column_writer.cc b/cpp/src/parquet/column_writer.cc
index 4998e6f301..393bf8a162 100644
--- a/cpp/src/parquet/column_writer.cc
+++ b/cpp/src/parquet/column_writer.cc
@@ -1032,13 +1032,15 @@ void ColumnWriterImpl::BuildDataPageV2(int64_t
definition_levels_rle_size,
const std::shared_ptr<Buffer>& values) {
// Compress the values if needed. Repetition and definition levels are
uncompressed in
// V2.
- std::shared_ptr<Buffer> compressed_values;
- if (pager_->has_compressor()) {
+ bool page_is_compressed = false;
+ if (pager_->has_compressor() && values->size() > 0) {
pager_->Compress(*values, compressor_temp_buffer_.get());
- compressed_values = compressor_temp_buffer_;
- } else {
- compressed_values = values;
+ if (compressor_temp_buffer_->size() < values->size()) {
+ page_is_compressed = true;
+ }
}
+ std::shared_ptr<Buffer> compressed_values =
+ (page_is_compressed ? compressor_temp_buffer_ : values);
// Concatenate uncompressed levels and the possibly compressed values
int64_t combined_size =
@@ -1071,14 +1073,14 @@ void ColumnWriterImpl::BuildDataPageV2(int64_t
definition_levels_rle_size,
combined->CopySlice(0, combined->size(),
allocator_));
std::unique_ptr<DataPage> page_ptr = std::make_unique<DataPageV2>(
combined, num_values, null_count, num_rows, encoding_,
def_levels_byte_length,
- rep_levels_byte_length, uncompressed_size, pager_->has_compressor(),
+ rep_levels_byte_length, uncompressed_size, page_is_compressed,
std::move(page_stats), first_row_index, std::move(page_size_stats));
total_compressed_bytes_ += page_ptr->size() + sizeof(format::PageHeader);
data_pages_.push_back(std::move(page_ptr));
} else {
DataPageV2 page(combined, num_values, null_count, num_rows, encoding_,
def_levels_byte_length, rep_levels_byte_length,
uncompressed_size,
- pager_->has_compressor(), std::move(page_stats),
first_row_index,
+ page_is_compressed, std::move(page_stats), first_row_index,
std::move(page_size_stats));
WriteDataPage(page);
}
diff --git a/cpp/src/parquet/column_writer_test.cc
b/cpp/src/parquet/column_writer_test.cc
index 744859cf0f..41c4a3223e 100644
--- a/cpp/src/parquet/column_writer_test.cc
+++ b/cpp/src/parquet/column_writer_test.cc
@@ -410,7 +410,7 @@ class TestPrimitiveWriter : public
PrimitiveTypedTest<TestType> {
const ColumnDescriptor* descr_;
- private:
+ protected:
std::unique_ptr<ColumnChunkMetaDataBuilder> metadata_;
std::shared_ptr<::arrow::io::BufferOutputStream> sink_;
std::shared_ptr<WriterProperties> writer_properties_;
@@ -1807,5 +1807,50 @@ TEST_F(TestValuesWriterInt32Type,
AllNullsCompressionInPageV2) {
}
}
+#ifdef ARROW_WITH_ZSTD
+
+TEST_F(TestValuesWriterInt32Type, AvoidCompressedInDataPageV2) {
+ Compression::type compression = Compression::ZSTD;
+ auto verify_only_one_uncompressed_page = [&](int total_num_values) {
+ ColumnProperties column_properties;
+ column_properties.set_compression(compression);
+
+ auto writer =
+ this->BuildWriter(SMALL_SIZE, column_properties,
ParquetVersion::PARQUET_2_LATEST,
+ ParquetDataPageVersion::V2);
+ writer->WriteBatch(static_cast<int64_t>(values_.size()),
this->def_levels_.data(),
+ nullptr, this->values_ptr_);
+ writer->Close();
+ ASSERT_OK_AND_ASSIGN(auto buffer, this->sink_->Finish());
+ auto source = std::make_shared<::arrow::io::BufferReader>(buffer);
+ ReaderProperties readerProperties;
+ std::unique_ptr<PageReader> page_reader = PageReader::Open(
+ std::move(source), total_num_values, compression, readerProperties);
+ auto data_page =
std::static_pointer_cast<DataPageV2>(page_reader->NextPage());
+ ASSERT_TRUE(data_page != nullptr);
+ ASSERT_FALSE(data_page->is_compressed());
+ ASSERT_TRUE(page_reader->NextPage() == nullptr);
+ };
+ {
+ // zero-sized data buffer should be handled correctly.
+ this->SetUpSchema(Repetition::OPTIONAL);
+ this->GenerateData(SMALL_SIZE);
+ std::fill(this->def_levels_.begin(), this->def_levels_.end(), 0);
+ verify_only_one_uncompressed_page(SMALL_SIZE);
+ }
+ {
+ // When only compress little data, the compressed size would even be
+ // larger than the original size. In this case, the `is_compressed` flag
+ // should be set to false.
+ this->SetUpSchema(Repetition::OPTIONAL);
+ this->GenerateData(1);
+ std::fill(this->def_levels_.begin(), this->def_levels_.end(), 1);
+ values_[0] = 142857;
+ verify_only_one_uncompressed_page(/*total_num_values=*/1);
+ }
+}
+
+#endif
+
} // namespace test
} // namespace parquet