From 34be8e8ea0868613a3c44b71c252c6022556dc42 Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Thu, 11 Jun 2026 15:47:22 +0800 Subject: [PATCH 1/9] GH-50156: [C++][Parquet] Ignore min/max for unknown column order --- cpp/src/parquet/metadata.cc | 134 ++++++++++++++++++++++------- cpp/src/parquet/metadata_test.cc | 98 +++++++++++++++++++++ cpp/src/parquet/page_index.cc | 9 ++ cpp/src/parquet/page_index_test.cc | 52 +++++++++++ cpp/src/parquet/statistics_test.cc | 7 ++ cpp/src/parquet/types.cc | 1 + cpp/src/parquet/types.h | 10 ++- 7 files changed, 279 insertions(+), 32 deletions(-) diff --git a/cpp/src/parquet/metadata.cc b/cpp/src/parquet/metadata.cc index 46cd1e442461..718c18706334 100644 --- a/cpp/src/parquet/metadata.cc +++ b/cpp/src/parquet/metadata.cc @@ -90,36 +90,111 @@ std::string ParquetVersionToString(ParquetVersion::type ver) { return "UNKNOWN"; } +namespace { + +enum class StatsMinMaxMode { + // Ignore min/max fields because their ordering is unknown or unsupported. + kDiscard, + // Use legacy min/max fields for files without column orders. + kLegacy, + // Use min_value/max_value fields with the column's well-defined order. + kNormal, +}; + +inline StatsMinMaxMode GetStatsMinMaxMode(const ColumnDescriptor& descr) { + switch (descr.column_order().get_order()) { + case ColumnOrder::TYPE_DEFINED_ORDER: + return descr.sort_order() != SortOrder::UNKNOWN ? StatsMinMaxMode::kNormal + : StatsMinMaxMode::kDiscard; + case ColumnOrder::UNDEFINED: + return descr.sort_order() != SortOrder::UNKNOWN ? StatsMinMaxMode::kLegacy + : StatsMinMaxMode::kDiscard; + case ColumnOrder::UNKNOWN: + return StatsMinMaxMode::kDiscard; + } + return StatsMinMaxMode::kDiscard; +} + +} // namespace + +static EncodedStatistics EncodedStatisticsFromThrift(const format::Statistics& statistics, + StatsMinMaxMode min_max) { + EncodedStatistics out; + + switch (min_max) { + case StatsMinMaxMode::kNormal: + if (statistics.__isset.max_value) { + out.set_max(statistics.max_value); + if (statistics.__isset.is_max_value_exact) { + out.is_max_value_exact = statistics.is_max_value_exact; + } + } + if (statistics.__isset.min_value) { + out.set_min(statistics.min_value); + if (statistics.__isset.is_min_value_exact) { + out.is_min_value_exact = statistics.is_min_value_exact; + } + } + break; + case StatsMinMaxMode::kLegacy: + if (statistics.__isset.max) { + out.set_max(statistics.max); + } + if (statistics.__isset.min) { + out.set_min(statistics.min); + } + break; + case StatsMinMaxMode::kDiscard: + break; + } + if (statistics.__isset.null_count) { + out.set_null_count(statistics.null_count); + } + if (statistics.__isset.distinct_count) { + out.set_distinct_count(statistics.distinct_count); + } + + return out; +} + template static std::shared_ptr MakeTypedColumnStats( const format::ColumnMetaData& metadata, const ColumnDescriptor* descr, ::arrow::MemoryPool* pool) { - std::optional min_exact = - metadata.statistics.__isset.is_min_value_exact - ? std::optional(metadata.statistics.is_min_value_exact) - : std::nullopt; - std::optional max_exact = - metadata.statistics.__isset.is_max_value_exact - ? std::optional(metadata.statistics.is_max_value_exact) - : std::nullopt; - // If ColumnOrder is defined, return max_value and min_value - if (descr->column_order().get_order() == ColumnOrder::TYPE_DEFINED_ORDER) { - return MakeStatistics( - descr, metadata.statistics.min_value, metadata.statistics.max_value, - metadata.num_values - metadata.statistics.null_count, - metadata.statistics.null_count, metadata.statistics.distinct_count, - metadata.statistics.__isset.max_value && metadata.statistics.__isset.min_value, - metadata.statistics.__isset.null_count, - metadata.statistics.__isset.distinct_count, min_exact, max_exact, pool); - } - // Default behavior + const auto& statistics = metadata.statistics; + const std::string kEmpty = ""; + const std::string* encoded_min = &kEmpty; + const std::string* encoded_max = &kEmpty; + bool has_min_max = false; + std::optional min_exact = std::nullopt; + std::optional max_exact = std::nullopt; + + switch (GetStatsMinMaxMode(*descr)) { + case StatsMinMaxMode::kNormal: + encoded_min = &statistics.min_value; + encoded_max = &statistics.max_value; + has_min_max = statistics.__isset.max_value && statistics.__isset.min_value; + min_exact = statistics.__isset.is_min_value_exact + ? std::optional(statistics.is_min_value_exact) + : std::nullopt; + max_exact = statistics.__isset.is_max_value_exact + ? std::optional(statistics.is_max_value_exact) + : std::nullopt; + break; + case StatsMinMaxMode::kLegacy: + encoded_min = &statistics.min; + encoded_max = &statistics.max; + has_min_max = statistics.__isset.max && statistics.__isset.min; + break; + case StatsMinMaxMode::kDiscard: + break; + } + return MakeStatistics( - descr, metadata.statistics.min, metadata.statistics.max, - metadata.num_values - metadata.statistics.null_count, - metadata.statistics.null_count, metadata.statistics.distinct_count, - metadata.statistics.__isset.max && metadata.statistics.__isset.min, - metadata.statistics.__isset.null_count, metadata.statistics.__isset.distinct_count, - min_exact, max_exact, pool); + descr, *encoded_min, *encoded_max, metadata.num_values - statistics.null_count, + statistics.null_count, statistics.distinct_count, has_min_max, + statistics.__isset.null_count, statistics.__isset.distinct_count, min_exact, + max_exact, pool); } namespace { @@ -337,11 +412,8 @@ class ColumnChunkMetaData::ColumnChunkMetaDataImpl { const std::lock_guard guard(stats_mutex_); if (possible_encoded_stats_ == nullptr) { possible_encoded_stats_ = - std::make_shared(FromThrift(column_metadata_->statistics)); - if (descr_->sort_order() == SortOrder::UNKNOWN) { - // If the column SortOrder is Unknown we can't trust max/min. - possible_encoded_stats_->ClearMinMax(); - } + std::make_shared(EncodedStatisticsFromThrift( + column_metadata_->statistics, GetStatsMinMaxMode(*descr_))); } } return writer_version_->HasCorrectStatistics(type(), *possible_encoded_stats_, @@ -1037,7 +1109,7 @@ class FileMetaData::FileMetaDataImpl { if (column_order.__isset.TYPE_ORDER) { column_orders.push_back(ColumnOrder::type_defined_); } else { - column_orders.push_back(ColumnOrder::undefined_); + column_orders.push_back(ColumnOrder::unknown_); } } } else { diff --git a/cpp/src/parquet/metadata_test.cc b/cpp/src/parquet/metadata_test.cc index 572f053179cd..b2a09be1eefb 100644 --- a/cpp/src/parquet/metadata_test.cc +++ b/cpp/src/parquet/metadata_test.cc @@ -275,6 +275,104 @@ TEST(Metadata, TestBuildAccess) { ASSERT_TRUE(f_accessor_1->Equals(*f_accessor->Subset({2, 0}))); } +namespace { + +std::string EncodeInt32(int32_t value) { + return std::string(reinterpret_cast(&value), sizeof(value)); +} + +constexpr int32_t kLegacyMin = 100, kLegacyMax = 200; + +std::string SerializeMetadata(const format::FileMetaData& thrift_metadata) { + std::string out; + ThriftSerializer{}.SerializeToString(&thrift_metadata, &out); + return out; +} + +std::shared_ptr ParseMetadata(std::string serialized_metadata) { + uint32_t decoded_len = static_cast(serialized_metadata.size()); + return FileMetaData::Make(serialized_metadata.data(), &decoded_len); +} + +format::FileMetaData SingleInt32MetadataWithStats() { + format::FileMetaData metadata; + format::SchemaElement root, leaf; + schema::NodeVector fields = {schema::Int32("int_col", Repetition::REQUIRED)}; + schema::GroupNode::Make("schema", Repetition::REPEATED, fields)->ToParquet(&root); + fields.back()->ToParquet(&leaf); + metadata.schema = {std::move(root), std::move(leaf)}; + + auto& column = metadata.row_groups.emplace_back().columns.emplace_back(); + column.__isset.meta_data = true; + auto& column_metadata = column.meta_data; + column_metadata.__set_type(format::Type::INT32); + column_metadata.__isset.statistics = true; + auto& statistics = column_metadata.statistics; + statistics.__set_min(EncodeInt32(kLegacyMin)); + statistics.__set_max(EncodeInt32(kLegacyMax)); + statistics.__set_min_value(EncodeInt32(kLegacyMin)); + statistics.__set_max_value(EncodeInt32(kLegacyMax)); + metadata.column_orders.emplace_back().__set_TYPE_ORDER(format::TypeDefinedOrder{}); + metadata.__isset.column_orders = true; + return metadata; +} + +std::unique_ptr GetOnlyColumnChunk( + const FileMetaData& metadata, ColumnOrder::type expected_order) { + EXPECT_EQ(expected_order, metadata.schema()->Column(0)->column_order().get_order()); + return metadata.RowGroup(0)->ColumnChunk(0); +} + +void AssertColumnChunkHasNoMinMax(const FileMetaData& metadata, + ColumnOrder::type expected_order) { + auto column = GetOnlyColumnChunk(metadata, expected_order); + ASSERT_NE(nullptr, column->encoded_statistics()); + EXPECT_FALSE(column->encoded_statistics()->has_min); + EXPECT_FALSE(column->encoded_statistics()->has_max); + ASSERT_NE(nullptr, column->statistics()); + EXPECT_FALSE(column->statistics()->HasMinMax()); +} + +void AssertColumnChunkMinMax(const FileMetaData& metadata, + ColumnOrder::type expected_order, int32_t min, int32_t max) { + auto column = GetOnlyColumnChunk(metadata, expected_order); + const std::string encoded_min = EncodeInt32(min); + const std::string encoded_max = EncodeInt32(max); + + ASSERT_NE(nullptr, column->encoded_statistics()); + EXPECT_EQ(encoded_min, column->encoded_statistics()->min()); + EXPECT_EQ(encoded_max, column->encoded_statistics()->max()); + ASSERT_NE(nullptr, column->statistics()); + EXPECT_EQ(encoded_min, column->statistics()->EncodeMin()); + EXPECT_EQ(encoded_max, column->statistics()->EncodeMax()); +} + +} // namespace + +TEST(Metadata, UnknownColumnOrderIgnoresMinMax) { + std::string serialized_metadata = SerializeMetadata(SingleInt32MetadataWithStats()); + const std::string kTypeDefinedOrder("\x1c\x00", 2); + const std::string kUnsupportedOrder("\x2c\x00", 2); + const auto pos = serialized_metadata.find(kTypeDefinedOrder); + ASSERT_NE(std::string::npos, pos); + serialized_metadata.replace(pos, kTypeDefinedOrder.size(), kUnsupportedOrder); + + auto metadata = ParseMetadata(serialized_metadata); + AssertColumnChunkHasNoMinMax(*metadata, ColumnOrder::UNKNOWN); +} + +TEST(Metadata, MissingColumnOrderUsesLegacyMinMax) { + format::FileMetaData thrift_metadata = SingleInt32MetadataWithStats(); + thrift_metadata.column_orders.clear(); + thrift_metadata.__isset.column_orders = false; + auto& statistics = thrift_metadata.row_groups.at(0).columns.at(0).meta_data.statistics; + statistics.__set_min_value(EncodeInt32(kLegacyMin - 100)); + statistics.__set_max_value(EncodeInt32(kLegacyMax + 100)); + + auto metadata = ParseMetadata(SerializeMetadata(thrift_metadata)); + AssertColumnChunkMinMax(*metadata, ColumnOrder::UNDEFINED, kLegacyMin, kLegacyMax); +} + TEST(Metadata, TestV1Version) { // PARQUET-839 parquet::schema::NodeVector fields; diff --git a/cpp/src/parquet/page_index.cc b/cpp/src/parquet/page_index.cc index 7434f2828da2..9556356be84e 100644 --- a/cpp/src/parquet/page_index.cc +++ b/cpp/src/parquet/page_index.cc @@ -38,6 +38,12 @@ namespace parquet { namespace { +inline bool CanTrustPageIndexMinMax(const ColumnDescriptor& descr) { + const auto column_order = descr.column_order().get_order(); + return column_order != ColumnOrder::UNKNOWN && column_order != ColumnOrder::UNDEFINED && + descr.sort_order() != SortOrder::UNKNOWN; +} + template void Decode(std::unique_ptr::Decoder>& decoder, const std::string& input, std::vector* output, @@ -973,6 +979,9 @@ std::unique_ptr ColumnIndex::Make(const ColumnDescriptor& descr, // Guard against UB when moving column_index throw ParquetException("Invalid ColumnIndex boundary_order"); } + if (!CanTrustPageIndexMinMax(descr)) { + return nullptr; + } switch (descr.physical_type()) { case Type::BOOLEAN: return std::make_unique>(descr, diff --git a/cpp/src/parquet/page_index_test.cc b/cpp/src/parquet/page_index_test.cc index 3a7308c1c6bc..df6ec6c6b098 100644 --- a/cpp/src/parquet/page_index_test.cc +++ b/cpp/src/parquet/page_index_test.cc @@ -561,6 +561,58 @@ void TestWriteTypedColumnIndex(schema::NodePtr node, } } +template +std::shared_ptr SerializedColumnIndex(const T& min, const T& max) { + auto encode = [](const T& value) { + return std::string(reinterpret_cast(&value), sizeof(value)); + }; + format::ColumnIndex column_index; + column_index.__set_null_pages({false}); + column_index.__set_min_values({encode(min)}); + column_index.__set_max_values({encode(max)}); + column_index.__set_boundary_order(format::BoundaryOrder::UNORDERED); + + auto sink = CreateOutputStream(); + ThriftSerializer{}.Serialize(&column_index, sink.get()); + PARQUET_ASSIGN_OR_THROW(auto buffer, sink->Finish()); + return buffer; +} + +void AssertColumnIndexIgnored(const ColumnDescriptor& read_descr, + const std::shared_ptr& buffer) { + auto column_index = + ColumnIndex::Make(read_descr, buffer->data(), static_cast(buffer->size()), + default_reader_properties()); + ASSERT_EQ(nullptr, column_index); +} + +void AssertColumnIndexIgnoredWithColumnOrder(ColumnOrder column_order) { + auto node = schema::Int32("c1"); + auto buffer = SerializedColumnIndex(/*min=*/1, /*max=*/2); + + std::static_pointer_cast(node)->SetColumnOrder(column_order); + ColumnDescriptor read_descr(node, /*max_definition_level=*/1, + /*max_repetition_level=*/0); + AssertColumnIndexIgnored(read_descr, buffer); +} + +TEST(PageIndex, ReadColumnIndexWithUnsupportedColumnOrder) { + AssertColumnIndexIgnoredWithColumnOrder(ColumnOrder::unknown_); + AssertColumnIndexIgnoredWithColumnOrder(ColumnOrder::undefined_); +} + +TEST(PageIndex, ReadColumnIndexWithUnknownSortOrder) { + auto node = schema::PrimitiveNode::Make("c1", Repetition::REQUIRED, Type::INT96); + ColumnDescriptor descr(node, /*max_definition_level=*/0, /*max_repetition_level=*/0); + ASSERT_EQ(SortOrder::UNKNOWN, descr.sort_order()); + + Int96 min{{1, 2, 3}}; + Int96 max{{4, 5, 6}}; + auto buffer = SerializedColumnIndex(min, max); + + AssertColumnIndexIgnored(descr, buffer); +} + TEST(PageIndex, WriteInt32ColumnIndex) { auto encode = [=](int32_t value) { return std::string(reinterpret_cast(&value), sizeof(int32_t)); diff --git a/cpp/src/parquet/statistics_test.cc b/cpp/src/parquet/statistics_test.cc index 905502cb0a57..0d300b856c2f 100644 --- a/cpp/src/parquet/statistics_test.cc +++ b/cpp/src/parquet/statistics_test.cc @@ -1660,6 +1660,13 @@ TEST(TestStatisticsSortOrder, UNKNOWN) { ASSERT_EQ(1, enc_stats->null_count); ASSERT_FALSE(enc_stats->is_max_value_exact.has_value()); ASSERT_FALSE(enc_stats->is_min_value_exact.has_value()); + + // Unknown sort order should not cause min/max to be set + std::shared_ptr stats = column_chunk->statistics(); + ASSERT_NE(nullptr, stats); + ASSERT_FALSE(stats->HasMinMax()); + ASSERT_TRUE(stats->HasNullCount()); + ASSERT_EQ(1, stats->null_count()); } // Test statistics for binary column with UNSIGNED sort order diff --git a/cpp/src/parquet/types.cc b/cpp/src/parquet/types.cc index fb4eb92a7544..ff5bd4eafb90 100644 --- a/cpp/src/parquet/types.cc +++ b/cpp/src/parquet/types.cc @@ -444,6 +444,7 @@ SortOrder::type GetSortOrder(const std::shared_ptr& logical_t ColumnOrder ColumnOrder::undefined_ = ColumnOrder(ColumnOrder::UNDEFINED); ColumnOrder ColumnOrder::type_defined_ = ColumnOrder(ColumnOrder::TYPE_DEFINED_ORDER); +ColumnOrder ColumnOrder::unknown_ = ColumnOrder(ColumnOrder::UNKNOWN); // Static methods for LogicalType class diff --git a/cpp/src/parquet/types.h b/cpp/src/parquet/types.h index ad4df5119e75..846a6cfaf9ce 100644 --- a/cpp/src/parquet/types.h +++ b/cpp/src/parquet/types.h @@ -599,7 +599,14 @@ bool PageCanUseChecksum(PageType::type pageType); class ColumnOrder { public: - enum type { UNDEFINED, TYPE_DEFINED_ORDER }; + enum type { + // File metadata has no column order, only legacy min/max in stats are defined. + UNDEFINED, + // File metadata uses TypeDefinedOrder from the Parquet format. + TYPE_DEFINED_ORDER, + // Column order value unsupported by this reader. + UNKNOWN + }; explicit ColumnOrder(ColumnOrder::type column_order) : column_order_(column_order) {} // Default to Type Defined Order ColumnOrder() : column_order_(type::TYPE_DEFINED_ORDER) {} @@ -607,6 +614,7 @@ class ColumnOrder { static ColumnOrder undefined_; static ColumnOrder type_defined_; + static ColumnOrder unknown_; private: ColumnOrder::type column_order_; From fed1776724c620ea3422df295e29b084830c6ac5 Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Thu, 11 Jun 2026 16:01:36 +0800 Subject: [PATCH 2/9] Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- cpp/src/parquet/metadata_test.cc | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/cpp/src/parquet/metadata_test.cc b/cpp/src/parquet/metadata_test.cc index b2a09be1eefb..92cedd1b8501 100644 --- a/cpp/src/parquet/metadata_test.cc +++ b/cpp/src/parquet/metadata_test.cc @@ -350,14 +350,14 @@ void AssertColumnChunkMinMax(const FileMetaData& metadata, } // namespace TEST(Metadata, UnknownColumnOrderIgnoresMinMax) { - std::string serialized_metadata = SerializeMetadata(SingleInt32MetadataWithStats()); - const std::string kTypeDefinedOrder("\x1c\x00", 2); - const std::string kUnsupportedOrder("\x2c\x00", 2); - const auto pos = serialized_metadata.find(kTypeDefinedOrder); - ASSERT_NE(std::string::npos, pos); - serialized_metadata.replace(pos, kTypeDefinedOrder.size(), kUnsupportedOrder); - - auto metadata = ParseMetadata(serialized_metadata); + format::FileMetaData thrift_metadata = SingleInt32MetadataWithStats(); + // Simulate an unsupported ColumnOrder value: unknown union fields are skipped by Thrift, + // leaving an entry with no known field set. + thrift_metadata.column_orders.clear(); + thrift_metadata.column_orders.emplace_back(); + thrift_metadata.__isset.column_orders = true; + + auto metadata = ParseMetadata(SerializeMetadata(thrift_metadata)); AssertColumnChunkHasNoMinMax(*metadata, ColumnOrder::UNKNOWN); } From a17f23262572a9ea773f35b87ce1a3cfee636fec Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Thu, 11 Jun 2026 16:16:55 +0800 Subject: [PATCH 3/9] Apply suggestion from @wgtmac --- cpp/src/parquet/metadata_test.cc | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cpp/src/parquet/metadata_test.cc b/cpp/src/parquet/metadata_test.cc index 92cedd1b8501..ac45be1fac38 100644 --- a/cpp/src/parquet/metadata_test.cc +++ b/cpp/src/parquet/metadata_test.cc @@ -351,8 +351,8 @@ void AssertColumnChunkMinMax(const FileMetaData& metadata, TEST(Metadata, UnknownColumnOrderIgnoresMinMax) { format::FileMetaData thrift_metadata = SingleInt32MetadataWithStats(); - // Simulate an unsupported ColumnOrder value: unknown union fields are skipped by Thrift, - // leaving an entry with no known field set. + // Simulate an unsupported ColumnOrder value: unknown union fields are skipped by + // Thrift, leaving an entry with no known field set. thrift_metadata.column_orders.clear(); thrift_metadata.column_orders.emplace_back(); thrift_metadata.__isset.column_orders = true; From a92a70d57d9b1202c6ee8f4bbcde182ddc7c0436 Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Thu, 11 Jun 2026 16:17:21 +0800 Subject: [PATCH 4/9] Apply suggestion from @wgtmac --- cpp/src/parquet/metadata.cc | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cpp/src/parquet/metadata.cc b/cpp/src/parquet/metadata.cc index 718c18706334..c29a8f2dbf8e 100644 --- a/cpp/src/parquet/metadata.cc +++ b/cpp/src/parquet/metadata.cc @@ -101,7 +101,7 @@ enum class StatsMinMaxMode { kNormal, }; -inline StatsMinMaxMode GetStatsMinMaxMode(const ColumnDescriptor& descr) { +StatsMinMaxMode GetStatsMinMaxMode(const ColumnDescriptor& descr) { switch (descr.column_order().get_order()) { case ColumnOrder::TYPE_DEFINED_ORDER: return descr.sort_order() != SortOrder::UNKNOWN ? StatsMinMaxMode::kNormal From d42ce32969af14cb84e21f84fdaf66ae92128617 Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Thu, 11 Jun 2026 16:17:37 +0800 Subject: [PATCH 5/9] Apply suggestion from @wgtmac --- cpp/src/parquet/page_index.cc | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cpp/src/parquet/page_index.cc b/cpp/src/parquet/page_index.cc index 9556356be84e..a1ec4f75d39c 100644 --- a/cpp/src/parquet/page_index.cc +++ b/cpp/src/parquet/page_index.cc @@ -38,7 +38,7 @@ namespace parquet { namespace { -inline bool CanTrustPageIndexMinMax(const ColumnDescriptor& descr) { +bool CanTrustPageIndexMinMax(const ColumnDescriptor& descr) { const auto column_order = descr.column_order().get_order(); return column_order != ColumnOrder::UNKNOWN && column_order != ColumnOrder::UNDEFINED && descr.sort_order() != SortOrder::UNKNOWN; From 0ada0c0d68bb924fdedb17794f54e531b395c0c5 Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Fri, 12 Jun 2026 14:04:53 +0800 Subject: [PATCH 6/9] add can_use_stats() and revert nullable column index --- cpp/src/parquet/metadata.cc | 88 ++++++------------------------ cpp/src/parquet/page_index.cc | 9 --- cpp/src/parquet/page_index_test.cc | 52 ------------------ cpp/src/parquet/schema.h | 13 +++++ cpp/src/parquet/schema_test.cc | 26 +++++++++ cpp/src/parquet/thrift_internal.h | 29 +++++++--- 6 files changed, 79 insertions(+), 138 deletions(-) diff --git a/cpp/src/parquet/metadata.cc b/cpp/src/parquet/metadata.cc index c29a8f2dbf8e..2b901b39de58 100644 --- a/cpp/src/parquet/metadata.cc +++ b/cpp/src/parquet/metadata.cc @@ -92,75 +92,26 @@ std::string ParquetVersionToString(ParquetVersion::type ver) { namespace { -enum class StatsMinMaxMode { - // Ignore min/max fields because their ordering is unknown or unsupported. - kDiscard, - // Use legacy min/max fields for files without column orders. - kLegacy, - // Use min_value/max_value fields with the column's well-defined order. - kNormal, -}; - -StatsMinMaxMode GetStatsMinMaxMode(const ColumnDescriptor& descr) { +StatisticsMinMaxField GetStatisticsMinMaxField(const ColumnDescriptor& descr) { switch (descr.column_order().get_order()) { case ColumnOrder::TYPE_DEFINED_ORDER: - return descr.sort_order() != SortOrder::UNKNOWN ? StatsMinMaxMode::kNormal - : StatsMinMaxMode::kDiscard; + return descr.sort_order() != SortOrder::UNKNOWN + ? StatisticsMinMaxField::kMinValueMaxValue + : StatisticsMinMaxField::kInvalid; case ColumnOrder::UNDEFINED: - return descr.sort_order() != SortOrder::UNKNOWN ? StatsMinMaxMode::kLegacy - : StatsMinMaxMode::kDiscard; + return descr.sort_order() == SortOrder::SIGNED + ? StatisticsMinMaxField::kLegacyMinMax + : StatisticsMinMaxField::kInvalid; case ColumnOrder::UNKNOWN: - return StatsMinMaxMode::kDiscard; - } - return StatsMinMaxMode::kDiscard; -} - -} // namespace - -static EncodedStatistics EncodedStatisticsFromThrift(const format::Statistics& statistics, - StatsMinMaxMode min_max) { - EncodedStatistics out; - - switch (min_max) { - case StatsMinMaxMode::kNormal: - if (statistics.__isset.max_value) { - out.set_max(statistics.max_value); - if (statistics.__isset.is_max_value_exact) { - out.is_max_value_exact = statistics.is_max_value_exact; - } - } - if (statistics.__isset.min_value) { - out.set_min(statistics.min_value); - if (statistics.__isset.is_min_value_exact) { - out.is_min_value_exact = statistics.is_min_value_exact; - } - } - break; - case StatsMinMaxMode::kLegacy: - if (statistics.__isset.max) { - out.set_max(statistics.max); - } - if (statistics.__isset.min) { - out.set_min(statistics.min); - } - break; - case StatsMinMaxMode::kDiscard: - break; + return StatisticsMinMaxField::kInvalid; } - if (statistics.__isset.null_count) { - out.set_null_count(statistics.null_count); - } - if (statistics.__isset.distinct_count) { - out.set_distinct_count(statistics.distinct_count); - } - - return out; + return StatisticsMinMaxField::kInvalid; } template -static std::shared_ptr MakeTypedColumnStats( - const format::ColumnMetaData& metadata, const ColumnDescriptor* descr, - ::arrow::MemoryPool* pool) { +std::shared_ptr MakeTypedColumnStats(const format::ColumnMetaData& metadata, + const ColumnDescriptor* descr, + ::arrow::MemoryPool* pool) { const auto& statistics = metadata.statistics; const std::string kEmpty = ""; const std::string* encoded_min = &kEmpty; @@ -169,8 +120,8 @@ static std::shared_ptr MakeTypedColumnStats( std::optional min_exact = std::nullopt; std::optional max_exact = std::nullopt; - switch (GetStatsMinMaxMode(*descr)) { - case StatsMinMaxMode::kNormal: + switch (GetStatisticsMinMaxField(*descr)) { + case StatisticsMinMaxField::kMinValueMaxValue: encoded_min = &statistics.min_value; encoded_max = &statistics.max_value; has_min_max = statistics.__isset.max_value && statistics.__isset.min_value; @@ -181,12 +132,12 @@ static std::shared_ptr MakeTypedColumnStats( ? std::optional(statistics.is_max_value_exact) : std::nullopt; break; - case StatsMinMaxMode::kLegacy: + case StatisticsMinMaxField::kLegacyMinMax: encoded_min = &statistics.min; encoded_max = &statistics.max; has_min_max = statistics.__isset.max && statistics.__isset.min; break; - case StatsMinMaxMode::kDiscard: + case StatisticsMinMaxField::kInvalid: break; } @@ -197,8 +148,6 @@ static std::shared_ptr MakeTypedColumnStats( max_exact, pool); } -namespace { - std::shared_ptr MakeColumnGeometryStats( const format::ColumnMetaData& metadata, const ColumnDescriptor* descr) { if (metadata.__isset.geospatial_statistics) { @@ -411,9 +360,8 @@ class ColumnChunkMetaData::ColumnChunkMetaDataImpl { { const std::lock_guard guard(stats_mutex_); if (possible_encoded_stats_ == nullptr) { - possible_encoded_stats_ = - std::make_shared(EncodedStatisticsFromThrift( - column_metadata_->statistics, GetStatsMinMaxMode(*descr_))); + possible_encoded_stats_ = std::make_shared( + FromThrift(column_metadata_->statistics, GetStatisticsMinMaxField(*descr_))); } } return writer_version_->HasCorrectStatistics(type(), *possible_encoded_stats_, diff --git a/cpp/src/parquet/page_index.cc b/cpp/src/parquet/page_index.cc index a1ec4f75d39c..7434f2828da2 100644 --- a/cpp/src/parquet/page_index.cc +++ b/cpp/src/parquet/page_index.cc @@ -38,12 +38,6 @@ namespace parquet { namespace { -bool CanTrustPageIndexMinMax(const ColumnDescriptor& descr) { - const auto column_order = descr.column_order().get_order(); - return column_order != ColumnOrder::UNKNOWN && column_order != ColumnOrder::UNDEFINED && - descr.sort_order() != SortOrder::UNKNOWN; -} - template void Decode(std::unique_ptr::Decoder>& decoder, const std::string& input, std::vector* output, @@ -979,9 +973,6 @@ std::unique_ptr ColumnIndex::Make(const ColumnDescriptor& descr, // Guard against UB when moving column_index throw ParquetException("Invalid ColumnIndex boundary_order"); } - if (!CanTrustPageIndexMinMax(descr)) { - return nullptr; - } switch (descr.physical_type()) { case Type::BOOLEAN: return std::make_unique>(descr, diff --git a/cpp/src/parquet/page_index_test.cc b/cpp/src/parquet/page_index_test.cc index df6ec6c6b098..3a7308c1c6bc 100644 --- a/cpp/src/parquet/page_index_test.cc +++ b/cpp/src/parquet/page_index_test.cc @@ -561,58 +561,6 @@ void TestWriteTypedColumnIndex(schema::NodePtr node, } } -template -std::shared_ptr SerializedColumnIndex(const T& min, const T& max) { - auto encode = [](const T& value) { - return std::string(reinterpret_cast(&value), sizeof(value)); - }; - format::ColumnIndex column_index; - column_index.__set_null_pages({false}); - column_index.__set_min_values({encode(min)}); - column_index.__set_max_values({encode(max)}); - column_index.__set_boundary_order(format::BoundaryOrder::UNORDERED); - - auto sink = CreateOutputStream(); - ThriftSerializer{}.Serialize(&column_index, sink.get()); - PARQUET_ASSIGN_OR_THROW(auto buffer, sink->Finish()); - return buffer; -} - -void AssertColumnIndexIgnored(const ColumnDescriptor& read_descr, - const std::shared_ptr& buffer) { - auto column_index = - ColumnIndex::Make(read_descr, buffer->data(), static_cast(buffer->size()), - default_reader_properties()); - ASSERT_EQ(nullptr, column_index); -} - -void AssertColumnIndexIgnoredWithColumnOrder(ColumnOrder column_order) { - auto node = schema::Int32("c1"); - auto buffer = SerializedColumnIndex(/*min=*/1, /*max=*/2); - - std::static_pointer_cast(node)->SetColumnOrder(column_order); - ColumnDescriptor read_descr(node, /*max_definition_level=*/1, - /*max_repetition_level=*/0); - AssertColumnIndexIgnored(read_descr, buffer); -} - -TEST(PageIndex, ReadColumnIndexWithUnsupportedColumnOrder) { - AssertColumnIndexIgnoredWithColumnOrder(ColumnOrder::unknown_); - AssertColumnIndexIgnoredWithColumnOrder(ColumnOrder::undefined_); -} - -TEST(PageIndex, ReadColumnIndexWithUnknownSortOrder) { - auto node = schema::PrimitiveNode::Make("c1", Repetition::REQUIRED, Type::INT96); - ColumnDescriptor descr(node, /*max_definition_level=*/0, /*max_repetition_level=*/0); - ASSERT_EQ(SortOrder::UNKNOWN, descr.sort_order()); - - Int96 min{{1, 2, 3}}; - Int96 max{{4, 5, 6}}; - auto buffer = SerializedColumnIndex(min, max); - - AssertColumnIndexIgnored(descr, buffer); -} - TEST(PageIndex, WriteInt32ColumnIndex) { auto encode = [=](int32_t value) { return std::string(reinterpret_cast(&value), sizeof(int32_t)); diff --git a/cpp/src/parquet/schema.h b/cpp/src/parquet/schema.h index 1addc73bd367..52c6b5544e32 100644 --- a/cpp/src/parquet/schema.h +++ b/cpp/src/parquet/schema.h @@ -382,6 +382,19 @@ class PARQUET_EXPORT ColumnDescriptor { return la ? GetSortOrder(la, pt) : GetSortOrder(converted_type(), pt); } + // Whether ColumnOrder-governed min/max values have a supported ordering. + bool can_use_stats() const { + switch (column_order().get_order()) { + case ColumnOrder::TYPE_DEFINED_ORDER: + return sort_order() != SortOrder::UNKNOWN; + case ColumnOrder::UNDEFINED: + return sort_order() == SortOrder::SIGNED; + case ColumnOrder::UNKNOWN: + return false; + } + return false; + } + const std::string& name() const { return primitive_node_->name(); } const std::shared_ptr path() const; diff --git a/cpp/src/parquet/schema_test.cc b/cpp/src/parquet/schema_test.cc index 2950a7df70f8..425f34ff61db 100644 --- a/cpp/src/parquet/schema_test.cc +++ b/cpp/src/parquet/schema_test.cc @@ -660,6 +660,32 @@ TEST(TestColumnDescriptor, TestAttrs) { ASSERT_EQ(expected_descr, descr2.ToString()); } +TEST(TestColumnDescriptor, CanUseStats) { + NodePtr node = Int32("name"); + ColumnDescriptor descr(node, 0, 0); + // Type-defined column order is usable when the type has a known sort order. + EXPECT_TRUE(descr.can_use_stats()); + + auto primitive_node = std::static_pointer_cast(node); + primitive_node->SetColumnOrder(ColumnOrder::undefined_); + // Missing column order falls back to legacy min/max, which are signed. + EXPECT_TRUE(ColumnDescriptor(node, 0, 0).can_use_stats()); + + primitive_node->SetColumnOrder(ColumnOrder::unknown_); + // Unsupported column order means min/max ordering is unknown to this reader. + EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_stats()); + + node = PrimitiveNode::Make("name", Repetition::REQUIRED, Type::BYTE_ARRAY); + primitive_node = std::static_pointer_cast(node); + primitive_node->SetColumnOrder(ColumnOrder::undefined_); + // Legacy min/max are signed, so they cannot represent unsigned byte ordering. + EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_stats()); + + node = PrimitiveNode::Make("name", Repetition::REQUIRED, Type::INT96); + // INT96 has no defined sort order in the Parquet type-defined ordering. + EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_stats()); +} + class TestSchemaDescriptor : public ::testing::Test { public: void setUp() {} diff --git a/cpp/src/parquet/thrift_internal.h b/cpp/src/parquet/thrift_internal.h index 5d14ae4e289c..e6d69308edb7 100644 --- a/cpp/src/parquet/thrift_internal.h +++ b/cpp/src/parquet/thrift_internal.h @@ -252,12 +252,21 @@ static inline AadMetadata FromThrift(format::AesGcmCtrV1 aesGcmCtrV1) { aesGcmCtrV1.supply_aad_prefix}; } -static inline EncodedStatistics FromThrift(const format::Statistics& stats) { +// Selects which thrift Statistics min/max fields should populate EncodedStatistics. +enum class StatisticsMinMaxField { + // Do not populate min/max, because the ordering is undefined or unsupported. + kInvalid, + // Populate min/max from the min_value/max_value fields. + kMinValueMaxValue, + // Populate min/max from the legacy min/max fields. + kLegacyMinMax, +}; + +static inline EncodedStatistics FromThrift(const format::Statistics& stats, + StatisticsMinMaxField min_max) { EncodedStatistics out; - // Use the new V2 min-max statistics over the former one if it is filled - if (stats.__isset.max_value || stats.__isset.min_value) { - // TODO: check if the column_order is TYPE_DEFINED_ORDER. + if (min_max == StatisticsMinMaxField::kMinValueMaxValue) { if (stats.__isset.max_value) { out.set_max(stats.max_value); if (stats.__isset.is_max_value_exact) { @@ -270,9 +279,7 @@ static inline EncodedStatistics FromThrift(const format::Statistics& stats) { out.is_min_value_exact = stats.is_min_value_exact; } } - } else if (stats.__isset.max || stats.__isset.min) { - // TODO: check created_by to see if it is corrupted for some types. - // TODO: check if the sort_order is SIGNED. + } else if (min_max == StatisticsMinMaxField::kLegacyMinMax) { if (stats.__isset.max) { out.set_max(stats.max); } @@ -290,6 +297,14 @@ static inline EncodedStatistics FromThrift(const format::Statistics& stats) { return out; } +static inline EncodedStatistics FromThrift(const format::Statistics& stats) { + // Use the new V2 min-max statistics over the former one if it is filled. + if (stats.__isset.max_value || stats.__isset.min_value) { + return FromThrift(stats, StatisticsMinMaxField::kMinValueMaxValue); + } + return FromThrift(stats, StatisticsMinMaxField::kLegacyMinMax); +} + static inline geospatial::EncodedGeoStatistics FromThrift( const format::GeospatialStatistics& geo_stats) { geospatial::EncodedGeoStatistics out; From b3360da1ee28da57aa76c0eac9ced1a3d237db60 Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Sat, 13 Jun 2026 20:15:42 +0800 Subject: [PATCH 7/9] rename to can_use_min_max --- cpp/src/parquet/schema.h | 2 +- cpp/src/parquet/schema_test.cc | 10 +++++----- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/cpp/src/parquet/schema.h b/cpp/src/parquet/schema.h index 52c6b5544e32..9725d2152597 100644 --- a/cpp/src/parquet/schema.h +++ b/cpp/src/parquet/schema.h @@ -383,7 +383,7 @@ class PARQUET_EXPORT ColumnDescriptor { } // Whether ColumnOrder-governed min/max values have a supported ordering. - bool can_use_stats() const { + bool can_use_min_max() const { switch (column_order().get_order()) { case ColumnOrder::TYPE_DEFINED_ORDER: return sort_order() != SortOrder::UNKNOWN; diff --git a/cpp/src/parquet/schema_test.cc b/cpp/src/parquet/schema_test.cc index 425f34ff61db..859f14a34d91 100644 --- a/cpp/src/parquet/schema_test.cc +++ b/cpp/src/parquet/schema_test.cc @@ -664,26 +664,26 @@ TEST(TestColumnDescriptor, CanUseStats) { NodePtr node = Int32("name"); ColumnDescriptor descr(node, 0, 0); // Type-defined column order is usable when the type has a known sort order. - EXPECT_TRUE(descr.can_use_stats()); + EXPECT_TRUE(descr.can_use_min_max()); auto primitive_node = std::static_pointer_cast(node); primitive_node->SetColumnOrder(ColumnOrder::undefined_); // Missing column order falls back to legacy min/max, which are signed. - EXPECT_TRUE(ColumnDescriptor(node, 0, 0).can_use_stats()); + EXPECT_TRUE(ColumnDescriptor(node, 0, 0).can_use_min_max()); primitive_node->SetColumnOrder(ColumnOrder::unknown_); // Unsupported column order means min/max ordering is unknown to this reader. - EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_stats()); + EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_min_max()); node = PrimitiveNode::Make("name", Repetition::REQUIRED, Type::BYTE_ARRAY); primitive_node = std::static_pointer_cast(node); primitive_node->SetColumnOrder(ColumnOrder::undefined_); // Legacy min/max are signed, so they cannot represent unsigned byte ordering. - EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_stats()); + EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_min_max()); node = PrimitiveNode::Make("name", Repetition::REQUIRED, Type::INT96); // INT96 has no defined sort order in the Parquet type-defined ordering. - EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_stats()); + EXPECT_FALSE(ColumnDescriptor(node, 0, 0).can_use_min_max()); } class TestSchemaDescriptor : public ::testing::Test { From 4a691c7c0aebd3567ba6d8ff28c0650010ef689b Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Mon, 15 Jun 2026 20:33:51 +0800 Subject: [PATCH 8/9] address comment --- cpp/src/parquet/metadata.cc | 1 + cpp/src/parquet/schema.h | 3 +++ cpp/src/parquet/types.h | 2 +- 3 files changed, 5 insertions(+), 1 deletion(-) diff --git a/cpp/src/parquet/metadata.cc b/cpp/src/parquet/metadata.cc index 2b901b39de58..396ce2e57b6c 100644 --- a/cpp/src/parquet/metadata.cc +++ b/cpp/src/parquet/metadata.cc @@ -92,6 +92,7 @@ std::string ParquetVersionToString(ParquetVersion::type ver) { namespace { +// Keep this logic consistent with ColumnDescriptor::can_use_min_max(). StatisticsMinMaxField GetStatisticsMinMaxField(const ColumnDescriptor& descr) { switch (descr.column_order().get_order()) { case ColumnOrder::TYPE_DEFINED_ORDER: diff --git a/cpp/src/parquet/schema.h b/cpp/src/parquet/schema.h index 9725d2152597..65732603ea1d 100644 --- a/cpp/src/parquet/schema.h +++ b/cpp/src/parquet/schema.h @@ -388,6 +388,9 @@ class PARQUET_EXPORT ColumnDescriptor { case ColumnOrder::TYPE_DEFINED_ORDER: return sort_order() != SortOrder::UNKNOWN; case ColumnOrder::UNDEFINED: + // If there is no defined column order, the obsolete min and max fields + // in the Statistics object are to be used, and they are always sorted + // by signed comparison, so this has to be the column type's sort order. return sort_order() == SortOrder::SIGNED; case ColumnOrder::UNKNOWN: return false; diff --git a/cpp/src/parquet/types.h b/cpp/src/parquet/types.h index 846a6cfaf9ce..687353aa9bcb 100644 --- a/cpp/src/parquet/types.h +++ b/cpp/src/parquet/types.h @@ -597,7 +597,7 @@ struct PageType { bool PageCanUseChecksum(PageType::type pageType); -class ColumnOrder { +class PARQUET_EXPORT ColumnOrder { public: enum type { // File metadata has no column order, only legacy min/max in stats are defined. From 7edb1625fc8aefbc6a996cd478f1eb0ec5e66eb9 Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Tue, 16 Jun 2026 10:10:10 +0800 Subject: [PATCH 9/9] add new PageReader::Open with descr --- cpp/src/parquet/column_io_benchmark.cc | 3 +- cpp/src/parquet/column_reader.cc | 35 +++++++++++++++++------- cpp/src/parquet/column_reader.h | 10 +++++++ cpp/src/parquet/column_writer_test.cc | 9 +++--- cpp/src/parquet/file_deserialize_test.cc | 34 +++++++++++++++++------ cpp/src/parquet/file_reader.cc | 5 ++-- cpp/src/parquet/metadata.cc | 17 ------------ cpp/src/parquet/thrift_internal.h | 29 ++++++++++++++------ 8 files changed, 91 insertions(+), 51 deletions(-) diff --git a/cpp/src/parquet/column_io_benchmark.cc b/cpp/src/parquet/column_io_benchmark.cc index 4b29a1284d1a..0dc1b0f8c965 100644 --- a/cpp/src/parquet/column_io_benchmark.cc +++ b/cpp/src/parquet/column_io_benchmark.cc @@ -161,7 +161,8 @@ std::shared_ptr BuildReader(std::shared_ptr& buffer, int64_t num_values, Compression::type codec, ColumnDescriptor* schema) { auto source = std::make_shared<::arrow::io::BufferReader>(buffer); - std::unique_ptr page_reader = PageReader::Open(source, num_values, codec); + std::unique_ptr page_reader = + PageReader::Open(source, num_values, codec, ReaderProperties(), *schema); return std::static_pointer_cast( ColumnReader::Make(schema, std::move(page_reader))); } diff --git a/cpp/src/parquet/column_reader.cc b/cpp/src/parquet/column_reader.cc index 79b837f755c3..7241632e4b62 100644 --- a/cpp/src/parquet/column_reader.cc +++ b/cpp/src/parquet/column_reader.cc @@ -200,10 +200,10 @@ namespace { // Extracts encoded statistics from V1 and V2 data page headers template -EncodedStatistics ExtractStatsFromHeader(const H& header) { +EncodedStatistics ExtractStatsFromHeader(const H& header, StatisticsMinMaxField min_max) { EncodedStatistics page_statistics; if (header.__isset.statistics) { - page_statistics = FromThrift(header.statistics); + page_statistics = FromThrift(header.statistics, min_max); } return page_statistics; } @@ -225,13 +225,15 @@ class SerializedPageReader : public PageReader { public: SerializedPageReader(std::shared_ptr stream, int64_t total_num_values, Compression::type codec, const ReaderProperties& properties, - const CryptoContext* crypto_ctx, bool always_compressed) + const CryptoContext* crypto_ctx, bool always_compressed, + StatisticsMinMaxField stats_min_max_field) : properties_(properties), stream_(std::move(stream)), decompression_buffer_(AllocateBuffer(properties_.memory_pool(), 0)), page_ordinal_(0), seen_num_values_(0), - total_num_values_(total_num_values) { + total_num_values_(total_num_values), + stats_min_max_field_(stats_min_max_field) { if (crypto_ctx != nullptr) { crypto_ctx_ = *crypto_ctx; InitDecryption(); @@ -302,6 +304,8 @@ class SerializedPageReader : public PageReader { // Number of values in all the data pages int64_t total_num_values_; + StatisticsMinMaxField stats_min_max_field_; + // data_page_aad_ and data_page_header_aad_ contain the AAD for data page and data page // header in a single column respectively. // While calculating AAD for different pages in a single column the pages AAD is @@ -349,7 +353,7 @@ bool SerializedPageReader::ShouldSkipPage(EncodedStatistics* data_page_statistic if (page_type == PageType::DATA_PAGE) { const format::DataPageHeader& header = current_page_header_.data_page_header; CheckNumValuesInHeader(header.num_values); - *data_page_statistics = ExtractStatsFromHeader(header); + *data_page_statistics = ExtractStatsFromHeader(header, stats_min_max_field_); seen_num_values_ += header.num_values; if (data_page_filter_) { const EncodedStatistics* filter_statistics = @@ -370,7 +374,7 @@ bool SerializedPageReader::ShouldSkipPage(EncodedStatistics* data_page_statistic header.repetition_levels_byte_length < 0) { throw ParquetException("Invalid page header (negative levels byte length)"); } - *data_page_statistics = ExtractStatsFromHeader(header); + *data_page_statistics = ExtractStatsFromHeader(header, stats_min_max_field_); seen_num_values_ += header.num_values; if (data_page_filter_) { const EncodedStatistics* filter_statistics = @@ -591,6 +595,16 @@ std::shared_ptr SerializedPageReader::DecompressIfNeeded( } // namespace +std::unique_ptr PageReader::Open( + std::shared_ptr stream, int64_t total_num_values, + Compression::type codec, const ReaderProperties& properties, + const ColumnDescriptor& descr, bool always_compressed, const CryptoContext* ctx) { + const auto stats_min_max_field = GetStatisticsMinMaxField(descr); + return std::unique_ptr( + new SerializedPageReader(std::move(stream), total_num_values, codec, properties, + ctx, always_compressed, stats_min_max_field)); +} + std::unique_ptr PageReader::Open(std::shared_ptr stream, int64_t total_num_values, Compression::type codec, @@ -598,7 +612,8 @@ std::unique_ptr PageReader::Open(std::shared_ptr s bool always_compressed, const CryptoContext* ctx) { return std::unique_ptr(new SerializedPageReader( - std::move(stream), total_num_values, codec, properties, ctx, always_compressed)); + std::move(stream), total_num_values, codec, properties, ctx, always_compressed, + StatisticsMinMaxField::kMinValueMaxValue)); } std::unique_ptr PageReader::Open(std::shared_ptr stream, @@ -607,9 +622,9 @@ std::unique_ptr PageReader::Open(std::shared_ptr s bool always_compressed, ::arrow::MemoryPool* pool, const CryptoContext* ctx) { - return std::unique_ptr( - new SerializedPageReader(std::move(stream), total_num_values, codec, - ReaderProperties(pool), ctx, always_compressed)); + return std::unique_ptr(new SerializedPageReader( + std::move(stream), total_num_values, codec, ReaderProperties(pool), ctx, + always_compressed, StatisticsMinMaxField::kMinValueMaxValue)); } namespace { diff --git a/cpp/src/parquet/column_reader.h b/cpp/src/parquet/column_reader.h index ac4469b1904f..ad20d4a9e6e6 100644 --- a/cpp/src/parquet/column_reader.h +++ b/cpp/src/parquet/column_reader.h @@ -117,11 +117,21 @@ class PARQUET_EXPORT PageReader { public: virtual ~PageReader() = default; + static std::unique_ptr Open(std::shared_ptr stream, + int64_t total_num_values, + Compression::type codec, + const ReaderProperties& properties, + const ColumnDescriptor& descr, + bool always_compressed = false, + const CryptoContext* ctx = NULLPTR); + + PARQUET_DEPRECATED("Deprecated in 25.0.0. Use the ColumnDescriptor overload instead.") static std::unique_ptr Open( std::shared_ptr stream, int64_t total_num_values, Compression::type codec, bool always_compressed = false, ::arrow::MemoryPool* pool = ::arrow::default_memory_pool(), const CryptoContext* ctx = NULLPTR); + PARQUET_DEPRECATED("Deprecated in 25.0.0. Use the ColumnDescriptor overload instead.") static std::unique_ptr Open(std::shared_ptr stream, int64_t total_num_values, Compression::type codec, diff --git a/cpp/src/parquet/column_writer_test.cc b/cpp/src/parquet/column_writer_test.cc index a45394917203..69ec07c1b78f 100644 --- a/cpp/src/parquet/column_writer_test.cc +++ b/cpp/src/parquet/column_writer_test.cc @@ -100,8 +100,8 @@ class TestPrimitiveWriter : public PrimitiveTypedTest { auto source = std::make_shared<::arrow::io::BufferReader>(buffer); ReaderProperties readerProperties; readerProperties.set_page_checksum_verification(page_checksum_verify); - std::unique_ptr page_reader = - PageReader::Open(std::move(source), num_rows, compression, readerProperties); + std::unique_ptr page_reader = PageReader::Open( + std::move(source), num_rows, compression, readerProperties, *this->descr_); reader_ = std::static_pointer_cast>( ColumnReader::Make(this->descr_, std::move(page_reader))); } @@ -2058,8 +2058,9 @@ TEST_F(TestValuesWriterInt32Type, AvoidCompressedInDataPageV2) { ASSERT_OK_AND_ASSIGN(auto buffer, this->sink_->Finish()); auto source = std::make_shared<::arrow::io::BufferReader>(buffer); ReaderProperties readerProperties; - std::unique_ptr page_reader = PageReader::Open( - std::move(source), total_num_values, compression, readerProperties); + std::unique_ptr page_reader = + PageReader::Open(std::move(source), total_num_values, compression, + readerProperties, *this->descr_); auto data_page = std::static_pointer_cast(page_reader->NextPage()); ASSERT_TRUE(data_page != nullptr); ASSERT_FALSE(data_page->is_compressed()); diff --git a/cpp/src/parquet/file_deserialize_test.cc b/cpp/src/parquet/file_deserialize_test.cc index 7fa5e2f167e2..e81f7689a6cc 100644 --- a/cpp/src/parquet/file_deserialize_test.cc +++ b/cpp/src/parquet/file_deserialize_test.cc @@ -46,6 +46,14 @@ namespace parquet { using ::arrow::io::BufferReader; using ::parquet::DataPageStats; +// Make a column descriptor with an undefined column order. +ColumnDescriptor MakeColumnDescriptor() { + auto node = schema::Int32("int_col", Repetition::REQUIRED); + static_cast(node.get()) + ->SetColumnOrder(ColumnOrder::undefined_); + return ColumnDescriptor(node, /*max_definition_level=*/0, /*max_repetition_level=*/0); +} + // Adds page statistics occupying a certain amount of bytes (for testing very // large page headers) template @@ -110,6 +118,8 @@ static std::vector GetSupportedCodecTypes() { class TestPageSerde : public ::testing::Test { public: + TestPageSerde() : descr_(MakeColumnDescriptor()) {} + void SetUp() { data_page_header_.encoding = format::Encoding::PLAIN; data_page_header_.definition_level_encoding = format::Encoding::RLE; @@ -124,7 +134,7 @@ class TestPageSerde : public ::testing::Test { EndStream(); auto stream = std::make_shared<::arrow::io::BufferReader>(out_buffer_); - page_reader_ = PageReader::Open(stream, num_rows, codec, properties); + page_reader_ = PageReader::Open(stream, num_rows, codec, properties, descr_); } void WriteDataPageHeader(int max_serialized_len = 1024, int32_t uncompressed_size = 0, @@ -204,6 +214,7 @@ class TestPageSerde : public ::testing::Test { std::shared_ptr<::arrow::io::BufferOutputStream> out_stream_; std::shared_ptr out_buffer_; + ColumnDescriptor descr_; std::unique_ptr page_reader_; format::PageHeader page_header_; format::DataPageHeader data_page_header_; @@ -466,7 +477,8 @@ TYPED_TEST(PageFilterTest, TestPageWithoutStatistics) { auto stream = std::make_shared<::arrow::io::BufferReader>(this->out_buffer_); this->page_reader_ = - PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED); + PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED, + ReaderProperties(), this->descr_); int num_pages = 0; bool is_stats_null = false; @@ -492,7 +504,8 @@ TYPED_TEST(PageFilterTest, TestPageFilterCallback) { // are right. auto stream = std::make_shared<::arrow::io::BufferReader>(this->out_buffer_); this->page_reader_ = - PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED); + PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED, + ReaderProperties(), this->descr_); std::vector read_stats; std::vector read_num_values; @@ -526,7 +539,8 @@ TYPED_TEST(PageFilterTest, TestPageFilterCallback) { { // Skip all pages. auto stream = std::make_shared<::arrow::io::BufferReader>(this->out_buffer_); this->page_reader_ = - PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED); + PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED, + ReaderProperties(), this->descr_); auto skip_all_pages = [](const DataPageStats& stats) -> bool { return true; }; @@ -538,7 +552,8 @@ TYPED_TEST(PageFilterTest, TestPageFilterCallback) { { // Skip every other page. auto stream = std::make_shared<::arrow::io::BufferReader>(this->out_buffer_); this->page_reader_ = - PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED); + PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED, + ReaderProperties(), this->descr_); // Skip pages with even number of values. auto skip_even_pages = [](const DataPageStats& stats) -> bool { @@ -569,7 +584,8 @@ TYPED_TEST(PageFilterTest, TestChangingPageFilter) { auto stream = std::make_shared<::arrow::io::BufferReader>(this->out_buffer_); this->page_reader_ = - PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED); + PageReader::Open(stream, this->total_rows_, Compression::UNCOMPRESSED, + ReaderProperties(), this->descr_); // This callback will always return false. auto read_all_pages = [](const DataPageStats& stats) -> bool { return false; }; @@ -604,7 +620,8 @@ TEST_F(TestPageSerde, DoesNotFilterDictionaryPages) { // Try to read it back while asking for all data pages to be skipped. auto stream = std::make_shared<::arrow::io::BufferReader>(out_buffer_); - page_reader_ = PageReader::Open(stream, /*num_rows=*/100, Compression::UNCOMPRESSED); + page_reader_ = PageReader::Open(stream, /*num_rows=*/100, Compression::UNCOMPRESSED, + ReaderProperties(), descr_); auto skip_all_pages = [](const DataPageStats& stats) -> bool { return true; }; @@ -641,7 +658,8 @@ TEST_F(TestPageSerde, SkipsNonDataPages) { EndStream(); auto stream = std::make_shared<::arrow::io::BufferReader>(out_buffer_); - page_reader_ = PageReader::Open(stream, /*num_rows=*/100, Compression::UNCOMPRESSED); + page_reader_ = PageReader::Open(stream, /*num_rows=*/100, Compression::UNCOMPRESSED, + ReaderProperties(), descr_); // Only the two data pages are returned. std::shared_ptr current_page = page_reader_->NextPage(); diff --git a/cpp/src/parquet/file_reader.cc b/cpp/src/parquet/file_reader.cc index d0552adcee53..2f46a5e296f8 100644 --- a/cpp/src/parquet/file_reader.cc +++ b/cpp/src/parquet/file_reader.cc @@ -239,6 +239,7 @@ class SerializedRowGroup : public RowGroupReader::Contents { std::unique_ptr GetColumnPageReader(int i) override { // Read column chunk from the file auto col = row_group_metadata_->ColumnChunk(i); + const ColumnDescriptor* descr = row_group_metadata_->schema()->Column(i); ::arrow::io::ReadRange col_range = ComputeColumnChunkRange(file_metadata_, source_size_, row_group_ordinal_, i); @@ -263,7 +264,7 @@ class SerializedRowGroup : public RowGroupReader::Contents { // Column is encrypted only if crypto_metadata exists. if (!crypto_metadata) { return PageReader::Open(stream, col->num_values(), col->compression(), properties_, - always_compressed); + *descr, always_compressed); } // The column is encrypted @@ -286,7 +287,7 @@ class SerializedRowGroup : public RowGroupReader::Contents { std::move(meta_decryptor_factory), std::move(data_decryptor_factory)}; return PageReader::Open(stream, col->num_values(), col->compression(), properties_, - always_compressed, &ctx); + *descr, always_compressed, &ctx); } private: diff --git a/cpp/src/parquet/metadata.cc b/cpp/src/parquet/metadata.cc index 396ce2e57b6c..98f60df63dd4 100644 --- a/cpp/src/parquet/metadata.cc +++ b/cpp/src/parquet/metadata.cc @@ -92,23 +92,6 @@ std::string ParquetVersionToString(ParquetVersion::type ver) { namespace { -// Keep this logic consistent with ColumnDescriptor::can_use_min_max(). -StatisticsMinMaxField GetStatisticsMinMaxField(const ColumnDescriptor& descr) { - switch (descr.column_order().get_order()) { - case ColumnOrder::TYPE_DEFINED_ORDER: - return descr.sort_order() != SortOrder::UNKNOWN - ? StatisticsMinMaxField::kMinValueMaxValue - : StatisticsMinMaxField::kInvalid; - case ColumnOrder::UNDEFINED: - return descr.sort_order() == SortOrder::SIGNED - ? StatisticsMinMaxField::kLegacyMinMax - : StatisticsMinMaxField::kInvalid; - case ColumnOrder::UNKNOWN: - return StatisticsMinMaxField::kInvalid; - } - return StatisticsMinMaxField::kInvalid; -} - template std::shared_ptr MakeTypedColumnStats(const format::ColumnMetaData& metadata, const ColumnDescriptor* descr, diff --git a/cpp/src/parquet/thrift_internal.h b/cpp/src/parquet/thrift_internal.h index e6d69308edb7..971e6ccebc9d 100644 --- a/cpp/src/parquet/thrift_internal.h +++ b/cpp/src/parquet/thrift_internal.h @@ -44,6 +44,7 @@ #include "parquet/geospatial/statistics.h" #include "parquet/platform.h" #include "parquet/properties.h" +#include "parquet/schema.h" #include "parquet/size_statistics.h" #include "parquet/statistics.h" #include "parquet/types.h" @@ -252,7 +253,7 @@ static inline AadMetadata FromThrift(format::AesGcmCtrV1 aesGcmCtrV1) { aesGcmCtrV1.supply_aad_prefix}; } -// Selects which thrift Statistics min/max fields should populate EncodedStatistics. +// Selects how thrift Statistics min/max fields should populate EncodedStatistics. enum class StatisticsMinMaxField { // Do not populate min/max, because the ordering is undefined or unsupported. kInvalid, @@ -262,6 +263,24 @@ enum class StatisticsMinMaxField { kLegacyMinMax, }; +// Keep this field-selection logic consistent with ColumnDescriptor::can_use_min_max(). +static inline StatisticsMinMaxField GetStatisticsMinMaxField( + const ColumnDescriptor& descr) { + switch (descr.column_order().get_order()) { + case ColumnOrder::TYPE_DEFINED_ORDER: + return descr.sort_order() != SortOrder::UNKNOWN + ? StatisticsMinMaxField::kMinValueMaxValue + : StatisticsMinMaxField::kInvalid; + case ColumnOrder::UNDEFINED: + return descr.sort_order() == SortOrder::SIGNED + ? StatisticsMinMaxField::kLegacyMinMax + : StatisticsMinMaxField::kInvalid; + case ColumnOrder::UNKNOWN: + return StatisticsMinMaxField::kInvalid; + } + return StatisticsMinMaxField::kInvalid; +} + static inline EncodedStatistics FromThrift(const format::Statistics& stats, StatisticsMinMaxField min_max) { EncodedStatistics out; @@ -297,14 +316,6 @@ static inline EncodedStatistics FromThrift(const format::Statistics& stats, return out; } -static inline EncodedStatistics FromThrift(const format::Statistics& stats) { - // Use the new V2 min-max statistics over the former one if it is filled. - if (stats.__isset.max_value || stats.__isset.min_value) { - return FromThrift(stats, StatisticsMinMaxField::kMinValueMaxValue); - } - return FromThrift(stats, StatisticsMinMaxField::kLegacyMinMax); -} - static inline geospatial::EncodedGeoStatistics FromThrift( const format::GeospatialStatistics& geo_stats) { geospatial::EncodedGeoStatistics out;