diff --git a/cpp/src/parquet/arrow/arrow_statistics_test.cc b/cpp/src/parquet/arrow/arrow_statistics_test.cc index 27a76fd72bea..a6817f1cbc2d 100644 --- a/cpp/src/parquet/arrow/arrow_statistics_test.cc +++ b/cpp/src/parquet/arrow/arrow_statistics_test.cc @@ -20,6 +20,7 @@ #include "arrow/array.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" @@ -160,6 +161,108 @@ INSTANTIATE_TEST_SUITE_P( /*expected_min=*/"z", /*expected_max=*/"z"})); +TEST(StatisticsTest, FixedWidthLeafUnderListStructNullCount) { + // 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 = ::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(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); + + 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( + 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) { // 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..3296af62f0c4 100644 --- a/cpp/src/parquet/column_writer.cc +++ b/cpp/src/parquet/column_writer.cc @@ -1397,18 +1397,19 @@ class TypedColumnWriterImpl : public ColumnWriterImpl, 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, - /*num_levels=*/batch_size, 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, - null_count); + parquet_null_count); } - CommitWriteAndCheckPageLimit(batch_size, batch_num_spaced_values, null_count, - check_page); + CommitWriteAndCheckPageLimit(batch_size, batch_num_spaced_values, + parquet_null_count, check_page); value_offset += batch_num_spaced_values; // Dictionary size checked separately from data page size since we @@ -2051,7 +2052,9 @@ 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); + 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; };