Skip to content
Draft
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
13 changes: 11 additions & 2 deletions cpp/src/parquet/column_writer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -1057,8 +1057,17 @@ void ColumnWriterImpl::BuildDataPageV2(int64_t definition_levels_rle_size,
bool page_is_compressed = false;
if (pager_->has_compressor() && values->size() > 0) {
pager_->Compress(*values, compressor_temp_buffer_.get());
if (compressor_temp_buffer_->size() < values->size()) {
page_is_compressed = true;

if (const auto minimum = properties_->min_space_savings()) {

//Checks if compressed size meets the minimum threshold
const double savings =
1.0 - static_cast<double>(compressor_temp_buffer_->size()) / values->size();
page_is_compressed = savings >= *minimum;
} else {
// Keeps original behavior
// Allows users keep min_space_savings unset
page_is_compressed = compressor_temp_buffer_->size() < values->size();
}
}
std::shared_ptr<Buffer> compressed_values =
Expand Down
55 changes: 54 additions & 1 deletion cpp/src/parquet/column_writer_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@
// specific language governing permissions and limitations
// under the License.

#include <cmath>
#include <memory>
#include <optional>
#include <utility>
#include <vector>

Expand Down Expand Up @@ -112,7 +114,8 @@ class TestPrimitiveWriter : public PrimitiveTypedTest<TestType> {
const ParquetVersion::type version = ParquetVersion::PARQUET_1_0,
const ParquetDataPageVersion data_page_version = ParquetDataPageVersion::V1,
bool enable_checksum = false, int64_t page_size = kDefaultDataPageSize,
int64_t max_rows_per_page = kDefaultMaxRowsPerPage) {
int64_t max_rows_per_page = kDefaultMaxRowsPerPage,
std::optional<double> min_space_savings = {}) {
sink_ = CreateOutputStream();
WriterProperties::Builder wp_builder;
wp_builder.version(version)->data_page_version(data_page_version);
Expand All @@ -130,6 +133,7 @@ class TestPrimitiveWriter : public PrimitiveTypedTest<TestType> {
wp_builder.max_statistics_size(column_properties.max_statistics_size());
wp_builder.data_pagesize(page_size);
wp_builder.max_rows_per_page(max_rows_per_page);
wp_builder.min_space_savings(min_space_savings);
writer_properties_ = wp_builder.build();

metadata_ = ColumnChunkMetaDataBuilder::Make(writer_properties_, this->descr_);
Expand Down Expand Up @@ -2084,6 +2088,55 @@ TEST_F(TestValuesWriterInt32Type, AvoidCompressedInDataPageV2) {
verify_only_one_uncompressed_page(/*total_num_values=*/1);
}
}

// Tests if compression works when:
// - min_space_savings is unset
// - min_space_savings = 0.1
// - min_space_savings = 1
TEST_F(TestValuesWriterInt32Type, MinSpaceSavingsDataPageV2) {

// Used zeros because it compresses well
this->SetUpSchema(Repetition::OPTIONAL);
this->descr_ = this->schema_.Column(0);
this->GenerateData(SMALL_SIZE);
std::fill(this->values_.begin(), this->values_.end(), 0);

ColumnProperties column_properties;
column_properties.set_compression(Compression::ZSTD);

const std::vector<std::pair<std::optional<double>, bool>> cases = {
{std::nullopt, true},
{0.1, true},
{1.0, false},
};

for (const auto& [minimum, expected_compressed] : cases) {
auto writer = this->BuildWriter(
SMALL_SIZE, column_properties, ParquetVersion::PARQUET_2_LATEST,
ParquetDataPageVersion::V2, false, kDefaultDataPageSize,
kDefaultMaxRowsPerPage, minimum);

writer->WriteBatch(SMALL_SIZE, this->def_levels_.data(), nullptr,
this->values_ptr_);
writer->Close();

ASSERT_OK_AND_ASSIGN(auto buffer, this->sink_->Finish());

auto page_reader = PageReader::Open(
std::make_shared<::arrow::io::BufferReader>(buffer), SMALL_SIZE,
Compression::ZSTD, default_reader_properties(), *this->descr_);

auto page = page_reader->NextPage();

ASSERT_NE(page, nullptr);
ASSERT_EQ(PageType::DATA_PAGE_V2, page->type());

auto data_page = std::static_pointer_cast<DataPageV2>(page);

ASSERT_EQ(expected_compressed, data_page->is_compressed());
}
}

#endif

// Test writing and reading geometry columns
Expand Down
24 changes: 22 additions & 2 deletions cpp/src/parquet/properties.h
Original file line number Diff line number Diff line change
Expand Up @@ -403,6 +403,7 @@ class PARQUET_EXPORT WriterProperties {
max_rows_per_page_(properties.max_rows_per_page()),
version_(properties.version()),
data_page_version_(properties.data_page_version()),
min_space_savings_(properties.min_space_savings()),
created_by_(properties.created_by()),
store_decimal_as_integer_(properties.store_decimal_as_integer()),
page_checksum_enabled_(properties.page_checksum_enabled()),
Expand Down Expand Up @@ -529,6 +530,13 @@ class PARQUET_EXPORT WriterProperties {
return this;
}

/// Used to set minimum threshold for V2 page data values compression
/// If min_space_savings unset, keep compression if compressed value is smaller.
Builder* min_space_savings(std::optional<double> min_space_savings) {
min_space_savings_ = min_space_savings;
return this;
}

/// Specify the Parquet file version.
/// Default PARQUET_2_6.
Builder* version(ParquetVersion::type version) {
Expand Down Expand Up @@ -883,6 +891,11 @@ class PARQUET_EXPORT WriterProperties {
/// \brief Build the WriterProperties with the builder parameters.
/// \return The WriterProperties defined by the builder.
std::shared_ptr<WriterProperties> build() {
// checks for invalid input of min_space_savings
if (min_space_savings_ &&
!(*min_space_savings_ >= 0.0 && *min_space_savings_ <= 1.0)) {
throw ParquetException("min_space_savings must be in range [0, 1]");
}
std::unordered_map<std::string, ColumnProperties> column_properties;
auto get = [&](const std::string& key) -> ColumnProperties& {
auto it = column_properties.find(key);
Expand Down Expand Up @@ -918,7 +931,7 @@ class PARQUET_EXPORT WriterProperties {
pagesize_, max_rows_per_page_, version_, created_by_, page_checksum_enabled_,
size_statistics_level_, std::move(file_encryption_properties_),
default_column_properties_, column_properties, data_page_version_,
store_decimal_as_integer_, std::move(sorting_columns_),
min_space_savings_, store_decimal_as_integer_, std::move(sorting_columns_),
content_defined_chunking_enabled_, content_defined_chunking_options_));
}

Expand All @@ -933,6 +946,7 @@ class PARQUET_EXPORT WriterProperties {
int64_t max_rows_per_page_;
ParquetVersion::type version_;
ParquetDataPageVersion data_page_version_;
std::optional<double> min_space_savings_ = {};
std::string created_by_;
bool store_decimal_as_integer_;
bool page_checksum_enabled_;
Expand Down Expand Up @@ -973,6 +987,9 @@ class PARQUET_EXPORT WriterProperties {
return parquet_data_page_version_;
}


std::optional<double> min_space_savings() const { return min_space_savings_; }

inline ParquetVersion::type version() const { return parquet_version_; }

inline std::string created_by() const { return parquet_created_by_; }
Expand Down Expand Up @@ -1100,7 +1117,8 @@ class PARQUET_EXPORT WriterProperties {
std::shared_ptr<FileEncryptionProperties> file_encryption_properties,
const ColumnProperties& default_column_properties,
const std::unordered_map<std::string, ColumnProperties>& column_properties,
ParquetDataPageVersion data_page_version, bool store_short_decimal_as_integer,
ParquetDataPageVersion data_page_version, std::optional<double> min_space_savings,
bool store_short_decimal_as_integer,
std::vector<SortingColumn> sorting_columns, bool content_defined_chunking_enabled,
CdcOptions content_defined_chunking_options)
: pool_(pool),
Expand All @@ -1110,6 +1128,7 @@ class PARQUET_EXPORT WriterProperties {
pagesize_(pagesize),
max_rows_per_page_(max_rows_per_page),
parquet_data_page_version_(data_page_version),
min_space_savings_(min_space_savings),
parquet_version_(version),
parquet_created_by_(created_by),
store_decimal_as_integer_(store_short_decimal_as_integer),
Expand All @@ -1129,6 +1148,7 @@ class PARQUET_EXPORT WriterProperties {
int64_t pagesize_;
int64_t max_rows_per_page_;
ParquetDataPageVersion parquet_data_page_version_;
std::optional<double> min_space_savings_;
ParquetVersion::type parquet_version_;
std::string parquet_created_by_;
bool store_decimal_as_integer_;
Expand Down
27 changes: 27 additions & 0 deletions cpp/src/parquet/properties_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,33 @@ TEST(TestWriterProperties, DefaultCompression) {
::arrow::util::kUseDefaultCompressionLevel);
}

// checks if WriterProperties correctly store, copy, and clear min_space_savings
TEST(TestWriterProperties, MinSpaceSavings) {
WriterProperties::Builder builder;
ASSERT_FALSE(builder.build()->min_space_savings().has_value());

for (double savings : {0.0, 0.1, 1.0}) {
auto properties = builder.min_space_savings(savings)->build();
ASSERT_EQ(std::optional<double>(savings), properties->min_space_savings());

auto copied = WriterProperties::Builder(*properties).build();
ASSERT_EQ(properties->min_space_savings(), copied->min_space_savings());
}

auto properties = builder.min_space_savings(std::nullopt)->build();
ASSERT_FALSE(properties->min_space_savings().has_value());
}

// Tests for invalid settings
TEST(TestWriterProperties, InvalidMinSpaceSavings) {
for (double savings : {-0.1, 1.1, std::numeric_limits<double>::infinity(),
-std::numeric_limits<double>::infinity(),
std::numeric_limits<double>::quiet_NaN()}) {
EXPECT_THROW(WriterProperties::Builder().min_space_savings(savings)->build(),
ParquetException);
}
}

TEST(TestWriterProperties, AdvancedHandling) {
WriterProperties::Builder builder;
builder.compression("gzip", Compression::GZIP);
Expand Down
2 changes: 2 additions & 0 deletions python/pyarrow/_dataset_parquet.pyx
Original file line number Diff line number Diff line change
Expand Up @@ -664,6 +664,7 @@ cdef class ParquetFileWriteOptions(FileWriteOptions):
store_decimal_as_integer=self._properties["store_decimal_as_integer"],
use_content_defined_chunking=self._properties["use_content_defined_chunking"],
bloom_filter_options=self._properties["bloom_filter_options"],
min_space_savings=self._properties["min_space_savings"],
)

def _set_arrow_properties(self):
Expand Down Expand Up @@ -707,6 +708,7 @@ cdef class ParquetFileWriteOptions(FileWriteOptions):
use_byte_stream_split=False,
column_encoding=None,
data_page_version="1.0",
min_space_savings=None,
use_deprecated_int96_timestamps=False,
coerce_timestamps=None,
allow_truncated_timestamps=False,
Expand Down
3 changes: 2 additions & 1 deletion python/pyarrow/_parquet.pxd
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,8 @@ cdef shared_ptr[WriterProperties] _create_writer_properties(
sorting_columns=*,
store_decimal_as_integer=*,
use_content_defined_chunking=*,
bloom_filter_options=*
bloom_filter_options=*,
min_space_savings=*
) except *


Expand Down
16 changes: 13 additions & 3 deletions python/pyarrow/_parquet.pyx
Original file line number Diff line number Diff line change
Expand Up @@ -2076,13 +2076,21 @@ cdef shared_ptr[WriterProperties] _create_writer_properties(
sorting_columns=None,
store_decimal_as_integer=False,
use_content_defined_chunking=False,
bloom_filter_options=None) except *:
bloom_filter_options=None,
min_space_savings=None) except *:

"""General writer properties"""
cdef:
shared_ptr[WriterProperties] properties
WriterProperties.Builder props
CdcOptions cdc_options
double savings

if min_space_savings is not None:
savings = min_space_savings
if not (0.0 <= savings <= 1.0):
raise ValueError("min_space_savings must be in range [0, 1]")
props.min_space_savings(optional[double](savings))

# data_page_version

Expand Down Expand Up @@ -2405,7 +2413,8 @@ cdef class ParquetWriter(_Weakrefable):
store_decimal_as_integer=False,
use_content_defined_chunking=False,
write_time_adjusted_to_utc=False,
bloom_filter_options=None):
bloom_filter_options=None,
min_space_savings=None):
cdef:
shared_ptr[WriterProperties] properties
shared_ptr[ArrowWriterProperties] arrow_properties
Expand Down Expand Up @@ -2442,7 +2451,8 @@ cdef class ParquetWriter(_Weakrefable):
sorting_columns=sorting_columns,
store_decimal_as_integer=store_decimal_as_integer,
use_content_defined_chunking=use_content_defined_chunking,
bloom_filter_options=bloom_filter_options
bloom_filter_options=bloom_filter_options,
min_space_savings=min_space_savings,
)
arrow_properties = _create_arrow_writer_properties(
use_deprecated_int96_timestamps=use_deprecated_int96_timestamps,
Expand Down
1 change: 1 addition & 0 deletions python/pyarrow/includes/libparquet.pxd
Original file line number Diff line number Diff line change
Expand Up @@ -485,6 +485,7 @@ cdef extern from "parquet/api/writer.h" namespace "parquet" nogil:
cdef cppclass WriterProperties:
cppclass Builder:
Builder* data_page_version(ParquetDataPageVersion version)
Builder* min_space_savings(optional[double] min_space_savings)
Builder* version(ParquetVersion version)
Builder* compression(ParquetCompression codec)
Builder* compression(const c_string& path,
Expand Down
4 changes: 4 additions & 0 deletions python/pyarrow/parquet/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -1090,6 +1090,7 @@ def __init__(self, where, schema, filesystem=None,
max_rows_per_page=None,
bloom_filter_options=None,
use_content_defined_chunking=False,
min_space_savings=None,
**options):
if use_deprecated_int96_timestamps is None:
# Use int96 timestamps for Spark
Expand Down Expand Up @@ -1147,6 +1148,7 @@ def __init__(self, where, schema, filesystem=None,
max_rows_per_page=max_rows_per_page,
bloom_filter_options=bloom_filter_options,
use_content_defined_chunking=use_content_defined_chunking,
min_space_savings=min_space_savings,
**options)
self.is_open = True

Expand Down Expand Up @@ -2039,6 +2041,7 @@ def write_table(table, where, row_group_size=None, version='2.6',
max_rows_per_page=None,
bloom_filter_options=None,
use_content_defined_chunking=False,
min_space_savings=None,
**kwargs):
# Implementor's note: when adding keywords here / updating defaults, also
# update it in write_to_dataset and _dataset_parquet.pyx ParquetFileWriteOptions
Expand Down Expand Up @@ -2074,6 +2077,7 @@ def write_table(table, where, row_group_size=None, version='2.6',
max_rows_per_page=max_rows_per_page,
bloom_filter_options=bloom_filter_options,
use_content_defined_chunking=use_content_defined_chunking,
min_space_savings=min_space_savings,
**kwargs) as writer:
writer.write_table(table, row_group_size=row_group_size)
except Exception:
Expand Down
44 changes: 44 additions & 0 deletions python/pyarrow/tests/parquet/test_parquet_writer.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,50 @@
pytestmark = pytest.mark.parquet


@pytest.mark.parametrize("minimum", [None, 0.0, 0.1, 1.0])
@pytest.mark.parametrize("page_version", ["1.0", "2.0"])
@pytest.mark.parametrize("compression", ["none", "snappy"])
@pytest.mark.parametrize("use_writer", [False, True])
def test_min_space_savings(
minimum, page_version, compression, use_writer):
table = pa.table({
"value": pa.array([0.0] * 1000, type=pa.float32())
})
sink = pa.BufferOutputStream()
options = dict(
compression=compression,
use_dictionary=False,
data_page_version=page_version,
min_space_savings=minimum,
)

if use_writer:
with pq.ParquetWriter(sink, table.schema, **options) as writer:
writer.write_table(table)
else:
pq.write_table(table, sink, **options)

assert pq.read_table(sink.getvalue()).equals(table)

# Tests invalid settings
@pytest.mark.parametrize(
"minimum",
[-0.1, 1.1, float("inf"), -float("inf"), float("nan")]
)
@pytest.mark.parametrize("use_writer", [False, True])
def test_invalid_min_space_savings(minimum, use_writer):
table = pa.table({"value": [1.0]})
sink = pa.BufferOutputStream()

with pytest.raises(ValueError, match="min_space_savings"):
if use_writer:
pq.ParquetWriter(
sink, table.schema, min_space_savings=minimum)
else:
pq.write_table(
table, sink, min_space_savings=minimum)


@pytest.mark.pandas
def test_parquet_incremental_file_build(tempdir):
df = _test_dataframe(100)
Expand Down
Loading
Loading