From 464520a4c1b51663ba8cdccbd7dde5a1937a18d5 Mon Sep 17 00:00:00 2001 From: AnuragRaut08 Date: Wed, 16 Sep 2026 18:11:08 +0530 Subject: [PATCH 1/3] GH-51097: Fix Parquet null counts for fixed-width leaves Signed-off-by: AnuragRaut08 --- .../parquet/arrow/arrow_statistics_test.cc | 45 +++++++++++++++++++ cpp/src/parquet/column_writer.cc | 14 +++--- 2 files changed, 54 insertions(+), 5 deletions(-) diff --git a/cpp/src/parquet/arrow/arrow_statistics_test.cc b/cpp/src/parquet/arrow/arrow_statistics_test.cc index 27a76fd72bea..f755fd340889 100644 --- a/cpp/src/parquet/arrow/arrow_statistics_test.cc +++ b/cpp/src/parquet/arrow/arrow_statistics_test.cc @@ -160,6 +160,51 @@ INSTANTIATE_TEST_SUITE_P( /*expected_min=*/"z", /*expected_max=*/"z"})); +TEST(StatisticsTest, FixedWidthLeafUnderListStructNullCount) { + // Null counts for fixed-width leaves under list> + // must include null and empty list entries from the repeated ancestor. + auto schema = ::arrow::schema({::arrow::field( + "col", ::arrow::list(::arrow::struct_( + {::arrow::field("s", ::arrow::utf8()), + ::arrow::field("i32", ::arrow::int32())})))}); + + auto table = ::arrow::Table::Make( + schema, + {::arrow::ArrayFromJSON( + ::arrow::list(::arrow::struct_( + {::arrow::field("s", ::arrow::utf8()), + ::arrow::field("i32", ::arrow::int32())})), + R"([[{"s":"a","i32":1}],null,[],[{"s":null,"i32":null},{"s":"b","i32":2}]])")}); + + std::shared_ptr<::arrow::ResizableBuffer> serialized_data = AllocateBuffer(); + auto out_stream = + std::make_shared<::arrow::io::BufferOutputStream>(serialized_data); + + ASSERT_OK_AND_ASSIGN( + std::unique_ptr writer, + FileWriter::Open(*schema, default_memory_pool(), out_stream, + default_writer_properties(), + default_arrow_writer_properties())); + ASSERT_OK(writer->WriteTable(*table)); + ASSERT_OK(writer->Close()); + ASSERT_OK(out_stream->Close()); + + auto buffer_reader = std::make_shared<::arrow::io::BufferReader>(serialized_data); + auto parquet_reader = ParquetFileReader::Open(std::move(buffer_reader)); + auto metadata = parquet_reader->metadata(); + auto row_group = metadata->RowGroup(0); + + ASSERT_EQ(row_group->num_columns(), 2); + + auto int32_stats = row_group->ColumnChunk(1)->statistics(); + ASSERT_NE(int32_stats, nullptr); + + // Fixed-width leaves must include nulls from repeated ancestors + // (e.g. null or empty lists) in the column statistics. + EXPECT_EQ(int32_stats->null_count(), 3); + EXPECT_EQ(int32_stats->num_values(), 2); +} + TEST(StatisticsTest, TruncateOnlyHalfMinMax) { // GH-43382: Tests when we only have min or max, the `HasMinMax` should be false. std::shared_ptr<::arrow::ResizableBuffer> serialized_data = AllocateBuffer(); diff --git a/cpp/src/parquet/column_writer.cc b/cpp/src/parquet/column_writer.cc index 653f28f64bde..0f3bb7c42067 100644 --- a/cpp/src/parquet/column_writer.cc +++ b/cpp/src/parquet/column_writer.cc @@ -1394,21 +1394,22 @@ class TypedColumnWriterImpl : public ColumnWriterImpl, MaybeCalculateValidityBits(AddIfNotNull(def_levels, offset), batch_size, &batch_num_values, &batch_num_spaced_values, &null_count); + const int64_t total_null_count = batch_size - batch_num_values; WriteLevelsSpaced(batch_size, AddIfNotNull(def_levels, offset), AddIfNotNull(rep_levels, offset)); if (bits_buffer_ != nullptr) { WriteValuesSpaced(AddIfNotNull(values, value_offset), batch_num_values, batch_num_spaced_values, bits_buffer_->data(), /*offset=*/0, - /*num_levels=*/batch_size, null_count); + /*num_levels=*/batch_size, total_null_count); } else { WriteValuesSpaced(AddIfNotNull(values, value_offset), batch_num_values, batch_num_spaced_values, valid_bits, valid_bits_offset + value_offset, /*num_levels=*/batch_size, - null_count); + total_null_count); } - CommitWriteAndCheckPageLimit(batch_size, batch_num_spaced_values, null_count, - check_page); + CommitWriteAndCheckPageLimit(batch_size, batch_num_spaced_values, + total_null_count, check_page); value_offset += batch_num_spaced_values; // Dictionary size checked separately from data page size since we @@ -1750,6 +1751,7 @@ class TypedColumnWriterImpl : public ColumnWriterImpl, internal::DefLevelsToBitmap(def_levels, batch_size, level_info_, &io); *out_values_to_write = io.values_read - io.null_count; *out_spaced_values_to_write = io.values_read; + // io.null_count excludes nulls from repeated ancestors. *null_count = io.null_count; } @@ -2038,6 +2040,7 @@ Status TypedColumnWriterImpl::WriteArrowDictionary( // had so we need to recompute it from def levels. MaybeCalculateValidityBits(AddIfNotNull(def_levels, offset), batch_size, &batch_num_values, &batch_num_spaced_values, &null_count); + const int64_t total_null_count = batch_size - batch_num_values; WriteLevelsSpaced(batch_size, AddIfNotNull(def_levels, offset), AddIfNotNull(rep_levels, offset)); std::shared_ptr writeable_indices = @@ -2051,7 +2054,8 @@ Status TypedColumnWriterImpl::WriteArrowDictionary( dict_encoder->PutIndices(*writeable_indices); // Update unencoded byte array data size to size statistics UpdateUnencodedDataBytes(); - CommitWriteAndCheckPageLimit(batch_size, batch_num_values, null_count, check_page); + CommitWriteAndCheckPageLimit(batch_size, batch_num_values, total_null_count, + check_page); value_offset += batch_num_spaced_values; }; From b80a51e7b3ef15fb0f4e7cf905b8d6e414c77ed4 Mon Sep 17 00:00:00 2001 From: AnuragRaut08 Date: Tue, 22 Sep 2026 21:58:38 +0530 Subject: [PATCH 2/3] GH-51097: Fix Parquet null counts for fixed-width leaves --- .../parquet/arrow/arrow_statistics_test.cc | 127 ++++++++++++------ cpp/src/parquet/column_writer.cc | 12 +- 2 files changed, 91 insertions(+), 48 deletions(-) diff --git a/cpp/src/parquet/arrow/arrow_statistics_test.cc b/cpp/src/parquet/arrow/arrow_statistics_test.cc index f755fd340889..33787b4726d9 100644 --- a/cpp/src/parquet/arrow/arrow_statistics_test.cc +++ b/cpp/src/parquet/arrow/arrow_statistics_test.cc @@ -18,6 +18,7 @@ #include "gtest/gtest.h" #include "arrow/array.h" +#include "arrow/compute/api.h" #include "arrow/array/builder_primitive.h" #include "arrow/array/builder_time.h" #include "arrow/table.h" @@ -161,48 +162,90 @@ INSTANTIATE_TEST_SUITE_P( /*expected_max=*/"z"})); TEST(StatisticsTest, FixedWidthLeafUnderListStructNullCount) { - // Null counts for fixed-width leaves under list> - // must include null and empty list entries from the repeated ancestor. - auto schema = ::arrow::schema({::arrow::field( - "col", ::arrow::list(::arrow::struct_( - {::arrow::field("s", ::arrow::utf8()), - ::arrow::field("i32", ::arrow::int32())})))}); - - auto table = ::arrow::Table::Make( - schema, - {::arrow::ArrayFromJSON( - ::arrow::list(::arrow::struct_( - {::arrow::field("s", ::arrow::utf8()), - ::arrow::field("i32", ::arrow::int32())})), - R"([[{"s":"a","i32":1}],null,[],[{"s":null,"i32":null},{"s":"b","i32":2}]])")}); - - std::shared_ptr<::arrow::ResizableBuffer> serialized_data = AllocateBuffer(); - auto out_stream = - std::make_shared<::arrow::io::BufferOutputStream>(serialized_data); - - ASSERT_OK_AND_ASSIGN( - std::unique_ptr writer, - FileWriter::Open(*schema, default_memory_pool(), out_stream, - default_writer_properties(), - default_arrow_writer_properties())); - ASSERT_OK(writer->WriteTable(*table)); - ASSERT_OK(writer->Close()); - ASSERT_OK(out_stream->Close()); - - auto buffer_reader = std::make_shared<::arrow::io::BufferReader>(serialized_data); - auto parquet_reader = ParquetFileReader::Open(std::move(buffer_reader)); - auto metadata = parquet_reader->metadata(); - auto row_group = metadata->RowGroup(0); - - ASSERT_EQ(row_group->num_columns(), 2); - - auto int32_stats = row_group->ColumnChunk(1)->statistics(); - ASSERT_NE(int32_stats, nullptr); - - // Fixed-width leaves must include nulls from repeated ancestors - // (e.g. null or empty lists) in the column statistics. - EXPECT_EQ(int32_stats->null_count(), 3); - EXPECT_EQ(int32_stats->num_values(), 2); + // Null counts for leaves under list> must include null and empty + // list entries from the repeated ancestor. + for (const auto data_page_version : + {ParquetDataPageVersion::V1, ParquetDataPageVersion::V2}) { + for (const bool use_dictionary : {false, true}) { + SCOPED_TRACE(::testing::Message() + << "data_page_version=" << static_cast(data_page_version) + << ", use_dictionary=" << use_dictionary); + + auto string_type = + use_dictionary ? ::arrow::dictionary(::arrow::int32(), ::arrow::utf8()) + : ::arrow::utf8(); + auto list_type = ::arrow::list(::arrow::struct_( + {::arrow::field("s", string_type), + ::arrow::field("i32", ::arrow::int32())})); + auto schema = ::arrow::schema({::arrow::field("col", list_type)}); + + auto table = Table::Make( + schema, {ArrayFromJSON( + list_type, + R"([[{"s":"a","i32":1}],null,[],[{"s":null,"i32":null},{"s":"b","i32":2}]])")}); + + WriterProperties::Builder properties_builder; + properties_builder.data_page_version(data_page_version); + if (use_dictionary) { + properties_builder.enable_dictionary(); + } + + std::shared_ptr<::arrow::ResizableBuffer> serialized_data = AllocateBuffer(); + auto out_stream = + std::make_shared<::arrow::io::BufferOutputStream>(serialized_data); + + ASSERT_OK_AND_ASSIGN( + std::unique_ptr writer, + FileWriter::Open(*schema, default_memory_pool(), out_stream, + properties_builder.build(), + default_arrow_writer_properties())); + ASSERT_OK(writer->WriteTable(*table)); + ASSERT_OK(writer->Close()); + ASSERT_OK(out_stream->Close()); + + auto buffer_reader = std::make_shared<::arrow::io::BufferReader>(serialized_data); + auto parquet_reader = ParquetFileReader::Open(std::move(buffer_reader)); + auto metadata = parquet_reader->metadata(); + auto row_group = metadata->RowGroup(0); + + ASSERT_EQ(row_group->num_columns(), 2); + + for (int i = 0; i < 2; ++i) { + auto stats = row_group->ColumnChunk(i)->statistics(); + ASSERT_NE(stats, nullptr); + EXPECT_EQ(stats->null_count(), 3); + EXPECT_EQ(stats->num_values(), 2); + } + + ASSERT_OK_AND_ASSIGN( + auto file_reader, + FileReader::Make(default_memory_pool(), std::move(parquet_reader), + default_arrow_reader_properties())); + + ASSERT_OK_AND_ASSIGN(auto read_table, file_reader->ReadTable()); + + if (use_dictionary) { + auto plain_list_type = ::arrow::list(::arrow::struct_( + {::arrow::field("s", ::arrow::utf8()), + ::arrow::field("i32", ::arrow::int32())})); + + ASSERT_OK_AND_ASSIGN( + auto read_array, + ::arrow::compute::Cast(read_table->column(0)->chunk(0), + plain_list_type)); + + auto expected_array = table->column(0)->chunk(0); + + ASSERT_OK_AND_ASSIGN( + auto expected_plain, + ::arrow::compute::Cast(expected_array, plain_list_type)); + + ASSERT_TRUE(read_array.Equals(expected_plain)); + } else { + ASSERT_TRUE(read_table->Equals(*table)); + } + } + } } TEST(StatisticsTest, TruncateOnlyHalfMinMax) { diff --git a/cpp/src/parquet/column_writer.cc b/cpp/src/parquet/column_writer.cc index 0f3bb7c42067..da6d4bd9130d 100644 --- a/cpp/src/parquet/column_writer.cc +++ b/cpp/src/parquet/column_writer.cc @@ -1394,22 +1394,22 @@ class TypedColumnWriterImpl : public ColumnWriterImpl, MaybeCalculateValidityBits(AddIfNotNull(def_levels, offset), batch_size, &batch_num_values, &batch_num_spaced_values, &null_count); - const int64_t total_null_count = batch_size - batch_num_values; + const int64_t parquet_null_count = batch_size - batch_num_values; WriteLevelsSpaced(batch_size, AddIfNotNull(def_levels, offset), AddIfNotNull(rep_levels, offset)); if (bits_buffer_ != nullptr) { WriteValuesSpaced(AddIfNotNull(values, value_offset), batch_num_values, batch_num_spaced_values, bits_buffer_->data(), /*offset=*/0, - /*num_levels=*/batch_size, total_null_count); + /*num_levels=*/batch_size, parquet_null_count); } else { WriteValuesSpaced(AddIfNotNull(values, value_offset), batch_num_values, batch_num_spaced_values, valid_bits, valid_bits_offset + value_offset, /*num_levels=*/batch_size, - total_null_count); + parquet_null_count); } CommitWriteAndCheckPageLimit(batch_size, batch_num_spaced_values, - total_null_count, check_page); + parquet_null_count, check_page); value_offset += batch_num_spaced_values; // Dictionary size checked separately from data page size since we @@ -2040,7 +2040,7 @@ Status TypedColumnWriterImpl::WriteArrowDictionary( // had so we need to recompute it from def levels. MaybeCalculateValidityBits(AddIfNotNull(def_levels, offset), batch_size, &batch_num_values, &batch_num_spaced_values, &null_count); - const int64_t total_null_count = batch_size - batch_num_values; + const int64_t parquet_null_count = batch_size - batch_num_values; WriteLevelsSpaced(batch_size, AddIfNotNull(def_levels, offset), AddIfNotNull(rep_levels, offset)); std::shared_ptr writeable_indices = @@ -2054,7 +2054,7 @@ Status TypedColumnWriterImpl::WriteArrowDictionary( dict_encoder->PutIndices(*writeable_indices); // Update unencoded byte array data size to size statistics UpdateUnencodedDataBytes(); - CommitWriteAndCheckPageLimit(batch_size, batch_num_values, total_null_count, + CommitWriteAndCheckPageLimit(batch_size, batch_num_values, parquet_null_count, check_page); value_offset += batch_num_spaced_values; }; From 6c6b62e96ae99eefd949bf9401e6cacd01f195ff Mon Sep 17 00:00:00 2001 From: Gang Wu Date: Thu, 24 Sep 2026 23:19:01 +0800 Subject: [PATCH 3/3] GH-51097: [C++][Parquet] Cover nested null counts across page versions Signed-off-by: Gang Wu --- .../parquet/arrow/arrow_statistics_test.cc | 63 ++++++++++++------- cpp/src/parquet/column_writer.cc | 5 +- 2 files changed, 41 insertions(+), 27 deletions(-) diff --git a/cpp/src/parquet/arrow/arrow_statistics_test.cc b/cpp/src/parquet/arrow/arrow_statistics_test.cc index 33787b4726d9..a6817f1cbc2d 100644 --- a/cpp/src/parquet/arrow/arrow_statistics_test.cc +++ b/cpp/src/parquet/arrow/arrow_statistics_test.cc @@ -18,9 +18,9 @@ #include "gtest/gtest.h" #include "arrow/array.h" -#include "arrow/compute/api.h" #include "arrow/array/builder_primitive.h" #include "arrow/array/builder_time.h" +#include "arrow/compute/api.h" #include "arrow/table.h" #include "arrow/testing/gtest_util.h" @@ -171,34 +171,35 @@ TEST(StatisticsTest, FixedWidthLeafUnderListStructNullCount) { << "data_page_version=" << static_cast(data_page_version) << ", use_dictionary=" << use_dictionary); - auto string_type = - use_dictionary ? ::arrow::dictionary(::arrow::int32(), ::arrow::utf8()) - : ::arrow::utf8(); + auto string_type = use_dictionary + ? ::arrow::dictionary(::arrow::int32(), ::arrow::utf8()) + : ::arrow::utf8(); auto list_type = ::arrow::list(::arrow::struct_( - {::arrow::field("s", string_type), - ::arrow::field("i32", ::arrow::int32())})); + {::arrow::field("s", string_type), ::arrow::field("i32", ::arrow::int32())})); auto schema = ::arrow::schema({::arrow::field("col", list_type)}); - - auto table = Table::Make( - schema, {ArrayFromJSON( - list_type, - R"([[{"s":"a","i32":1}],null,[],[{"s":null,"i32":null},{"s":"b","i32":2}]])")}); + auto table = ::arrow::TableFromJSON(schema, {R"([ + [[{"s":"a","i32":1}]], + [null], + [[]], + [[{"s":null,"i32":null},{"s":"b","i32":2}]] + ])"}); WriterProperties::Builder properties_builder; properties_builder.data_page_version(data_page_version); if (use_dictionary) { properties_builder.enable_dictionary(); + } else { + properties_builder.disable_dictionary(); } std::shared_ptr<::arrow::ResizableBuffer> serialized_data = AllocateBuffer(); auto out_stream = std::make_shared<::arrow::io::BufferOutputStream>(serialized_data); - ASSERT_OK_AND_ASSIGN( - std::unique_ptr writer, - FileWriter::Open(*schema, default_memory_pool(), out_stream, - properties_builder.build(), - default_arrow_writer_properties())); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, + FileWriter::Open(*schema, default_memory_pool(), out_stream, + properties_builder.build(), + default_arrow_writer_properties())); ASSERT_OK(writer->WriteTable(*table)); ASSERT_OK(writer->Close()); ASSERT_OK(out_stream->Close()); @@ -215,6 +216,22 @@ TEST(StatisticsTest, FixedWidthLeafUnderListStructNullCount) { ASSERT_NE(stats, nullptr); EXPECT_EQ(stats->null_count(), 3); EXPECT_EQ(stats->num_values(), 2); + + if (data_page_version == ParquetDataPageVersion::V2) { + auto page_reader = parquet_reader->RowGroup(0)->GetColumnPageReader(i); + if (use_dictionary) { + auto dictionary_page = page_reader->NextPage(); + ASSERT_NE(dictionary_page, nullptr); + ASSERT_EQ(dictionary_page->type(), PageType::DICTIONARY_PAGE); + } + auto page = page_reader->NextPage(); + ASSERT_NE(page, nullptr); + ASSERT_EQ(page->type(), PageType::DATA_PAGE_V2); + auto data_page = std::static_pointer_cast(page); + EXPECT_EQ(data_page->num_values(), 5); + EXPECT_EQ(data_page->num_nulls(), 3); + EXPECT_EQ(page_reader->NextPage(), nullptr); + } } ASSERT_OK_AND_ASSIGN( @@ -225,20 +242,18 @@ TEST(StatisticsTest, FixedWidthLeafUnderListStructNullCount) { ASSERT_OK_AND_ASSIGN(auto read_table, file_reader->ReadTable()); if (use_dictionary) { - auto plain_list_type = ::arrow::list(::arrow::struct_( - {::arrow::field("s", ::arrow::utf8()), - ::arrow::field("i32", ::arrow::int32())})); + auto plain_list_type = + ::arrow::list(::arrow::struct_({::arrow::field("s", ::arrow::utf8()), + ::arrow::field("i32", ::arrow::int32())})); ASSERT_OK_AND_ASSIGN( auto read_array, - ::arrow::compute::Cast(read_table->column(0)->chunk(0), - plain_list_type)); + ::arrow::compute::Cast(read_table->column(0)->chunk(0), plain_list_type)); auto expected_array = table->column(0)->chunk(0); - ASSERT_OK_AND_ASSIGN( - auto expected_plain, - ::arrow::compute::Cast(expected_array, plain_list_type)); + ASSERT_OK_AND_ASSIGN(auto expected_plain, + ::arrow::compute::Cast(expected_array, plain_list_type)); ASSERT_TRUE(read_array.Equals(expected_plain)); } else { diff --git a/cpp/src/parquet/column_writer.cc b/cpp/src/parquet/column_writer.cc index da6d4bd9130d..3296af62f0c4 100644 --- a/cpp/src/parquet/column_writer.cc +++ b/cpp/src/parquet/column_writer.cc @@ -1394,10 +1394,10 @@ class TypedColumnWriterImpl : public ColumnWriterImpl, MaybeCalculateValidityBits(AddIfNotNull(def_levels, offset), batch_size, &batch_num_values, &batch_num_spaced_values, &null_count); - const int64_t parquet_null_count = batch_size - batch_num_values; WriteLevelsSpaced(batch_size, AddIfNotNull(def_levels, offset), AddIfNotNull(rep_levels, offset)); + const int64_t parquet_null_count = batch_size - batch_num_values; if (bits_buffer_ != nullptr) { WriteValuesSpaced(AddIfNotNull(values, value_offset), batch_num_values, batch_num_spaced_values, bits_buffer_->data(), /*offset=*/0, @@ -1751,7 +1751,6 @@ class TypedColumnWriterImpl : public ColumnWriterImpl, internal::DefLevelsToBitmap(def_levels, batch_size, level_info_, &io); *out_values_to_write = io.values_read - io.null_count; *out_spaced_values_to_write = io.values_read; - // io.null_count excludes nulls from repeated ancestors. *null_count = io.null_count; } @@ -2040,7 +2039,6 @@ Status TypedColumnWriterImpl::WriteArrowDictionary( // had so we need to recompute it from def levels. MaybeCalculateValidityBits(AddIfNotNull(def_levels, offset), batch_size, &batch_num_values, &batch_num_spaced_values, &null_count); - const int64_t parquet_null_count = batch_size - batch_num_values; WriteLevelsSpaced(batch_size, AddIfNotNull(def_levels, offset), AddIfNotNull(rep_levels, offset)); std::shared_ptr writeable_indices = @@ -2054,6 +2052,7 @@ Status TypedColumnWriterImpl::WriteArrowDictionary( dict_encoder->PutIndices(*writeable_indices); // Update unencoded byte array data size to size statistics UpdateUnencodedDataBytes(); + const int64_t parquet_null_count = batch_size - batch_num_values; CommitWriteAndCheckPageLimit(batch_size, batch_num_values, parquet_null_count, check_page); value_offset += batch_num_spaced_values;