Skip to content

Commit 8c1a8aa

Browse files
PingLiuPingmeta-codesync[bot]
authored andcommitted
feat(parquet): Add NaN statistics to Parquet writer (#14725)
Summary: Add NaN statistic to Parquet writer This change introduces a NaN statistic in the Parquet writer. Unlike other statistics (e.g., null_count, distinct_count), the NaN statistic is not written to the Parquet file footer. It is only reported to the Parquet writer caller when needed, such as when writing Iceberg Parquet data files. Pull Request resolved: #14725 Reviewed By: tanjialiang Differential Revision: D92527330 Pulled By: xiaoxmeng fbshipit-source-id: 96e1e52efe21af1c707edf3c6867ec40103c6dfe
1 parent b9e6b55 commit 8c1a8aa

5 files changed

Lines changed: 234 additions & 8 deletions

File tree

velox/dwio/parquet/writer/arrow/Metadata.cpp

Lines changed: 95 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -101,10 +101,12 @@ static std::shared_ptr<Statistics> MakeTypedColumnStats(
101101
metadata.num_values - metadata.statistics.null_count,
102102
metadata.statistics.null_count,
103103
metadata.statistics.distinct_count,
104+
/*nan_count=*/0,
104105
metadata.statistics.__isset.max_value ||
105106
metadata.statistics.__isset.min_value,
106107
metadata.statistics.__isset.null_count,
107-
metadata.statistics.__isset.distinct_count);
108+
metadata.statistics.__isset.distinct_count,
109+
/*has_nan_count=*/false);
108110
}
109111
// Default behavior
110112
return MakeStatistics<DType>(
@@ -114,9 +116,11 @@ static std::shared_ptr<Statistics> MakeTypedColumnStats(
114116
metadata.num_values - metadata.statistics.null_count,
115117
metadata.statistics.null_count,
116118
metadata.statistics.distinct_count,
119+
/*nan_count=*/0,
117120
metadata.statistics.__isset.max || metadata.statistics.__isset.min,
118121
metadata.statistics.__isset.null_count,
119-
metadata.statistics.__isset.distinct_count);
122+
metadata.statistics.__isset.distinct_count,
123+
/*has_nan_count=*/false);
120124
}
121125

122126
std::shared_ptr<Statistics> MakeColumnStats(
@@ -1015,6 +1019,22 @@ class FileMetaData::FileMetaDataImpl {
10151019
file_decryptor_ = file_decryptor;
10161020
}
10171021

1022+
// Set NaN counts from the builder (called during Finish)
1023+
// This stores total NaN counts per field ID across all row groups.
1024+
void setNaNCounts(
1025+
std::unordered_map<int32_t, std::pair<int64_t, bool>> nan_counts) {
1026+
field_nan_counts_ = std::move(nan_counts);
1027+
}
1028+
1029+
// Get total NaN count for a specific field ID across all row groups.
1030+
std::pair<int64_t, bool> getNaNCount(int32_t fieldId) const {
1031+
auto it = field_nan_counts_.find(fieldId);
1032+
if (it != field_nan_counts_.end()) {
1033+
return it->second;
1034+
}
1035+
return {0, false};
1036+
}
1037+
10181038
private:
10191039
friend FileMetaDataBuilder;
10201040
uint32_t metadata_len_ = 0;
@@ -1024,6 +1044,9 @@ class FileMetaData::FileMetaDataImpl {
10241044
std::shared_ptr<const KeyValueMetadata> key_value_metadata_;
10251045
const ReaderProperties properties_;
10261046
std::shared_ptr<InternalFileDecryptor> file_decryptor_;
1047+
// Total NaN counts per field ID across all row groups: field_id ->
1048+
// (nan_count, has_nan_count).
1049+
std::unordered_map<int32_t, std::pair<int64_t, bool>> field_nan_counts_;
10271050

10281051
void InitSchema() {
10291052
if (metadata_->schema.empty()) {
@@ -1200,6 +1223,10 @@ std::shared_ptr<FileMetaData> FileMetaData::Subset(
12001223
return impl_->Subset(row_groups);
12011224
}
12021225

1226+
std::pair<int64_t, bool> FileMetaData::getNaNCount(int32_t fieldId) const {
1227+
return impl_->getNaNCount(fieldId);
1228+
}
1229+
12031230
void FileMetaData::WriteTo(
12041231
::arrow::io::OutputStream* dst,
12051232
const std::shared_ptr<Encryptor>& encryptor) const {
@@ -1715,6 +1742,19 @@ class ColumnChunkMetaDataBuilder::ColumnChunkMetaDataBuilderImpl {
17151742
// column metadata
17161743
void SetStatistics(const EncodedStatistics& val) {
17171744
column_chunk_->meta_data.__set_statistics(ToThrift(val));
1745+
// Store NaN count separately since it's not written to the parquet file.
1746+
if (val.has_nan_count) {
1747+
nan_count_ = val.nan_count;
1748+
has_nan_count_ = true;
1749+
}
1750+
}
1751+
1752+
int64_t nan_count() const {
1753+
return nan_count_;
1754+
}
1755+
1756+
bool has_nan_count() const {
1757+
return has_nan_count_;
17181758
}
17191759

17201760
void Finish(
@@ -1883,6 +1923,9 @@ class ColumnChunkMetaDataBuilder::ColumnChunkMetaDataBuilderImpl {
18831923
owned_column_chunk_;
18841924
const std::shared_ptr<WriterProperties> properties_;
18851925
const ColumnDescriptor* column_;
1926+
// NaN count is stored separately since it's not written to the parquet file.
1927+
int64_t nan_count_ = 0;
1928+
bool has_nan_count_ = false;
18861929
};
18871930

18881931
std::unique_ptr<ColumnChunkMetaDataBuilder> ColumnChunkMetaDataBuilder::Make(
@@ -1970,6 +2013,14 @@ int64_t ColumnChunkMetaDataBuilder::total_compressed_size() const {
19702013
return impl_->total_compressed_size();
19712014
}
19722015

2016+
int64_t ColumnChunkMetaDataBuilder::nan_count() const {
2017+
return impl_->nan_count();
2018+
}
2019+
2020+
bool ColumnChunkMetaDataBuilder::has_nan_count() const {
2021+
return impl_->has_nan_count();
2022+
}
2023+
19732024
class RowGroupMetaDataBuilder::RowGroupMetaDataBuilderImpl {
19742025
public:
19752026
explicit RowGroupMetaDataBuilderImpl(
@@ -2062,6 +2113,16 @@ class RowGroupMetaDataBuilder::RowGroupMetaDataBuilderImpl {
20622113
return row_group_->num_rows;
20632114
}
20642115

2116+
// Returns a map of field_id -> (nan_count, has_nan_count).
2117+
std::unordered_map<int32_t, std::pair<int64_t, bool>> nan_counts() const {
2118+
std::unordered_map<int32_t, std::pair<int64_t, bool>> result;
2119+
for (const auto& builder : column_builders_) {
2120+
int32_t field_id = builder->descr()->schema_node()->field_id();
2121+
result[field_id] = {builder->nan_count(), builder->has_nan_count()};
2122+
}
2123+
return result;
2124+
}
2125+
20652126
private:
20662127
void InitializeColumns(int ncols) {
20672128
row_group_->columns.resize(ncols);
@@ -2119,6 +2180,11 @@ void RowGroupMetaDataBuilder::Finish(
21192180
impl_->Finish(total_bytes_written, row_group_ordinal);
21202181
}
21212182

2183+
std::unordered_map<int32_t, std::pair<int64_t, bool>>
2184+
RowGroupMetaDataBuilder::nan_counts() const {
2185+
return impl_->nan_counts();
2186+
}
2187+
21222188
// file metadata
21232189
class FileMetaDataBuilder::FileMetaDataBuilderImpl {
21242190
public:
@@ -2138,6 +2204,9 @@ class FileMetaDataBuilder::FileMetaDataBuilderImpl {
21382204
}
21392205

21402206
RowGroupMetaDataBuilder* AppendRowGroup() {
2207+
// Accumulate NaN counts from the previous row group before creating a new
2208+
// one.
2209+
accumulateNaNCountsFromCurrentRowGroup();
21412210
row_groups_.emplace_back();
21422211
current_row_group_builder_ = RowGroupMetaDataBuilder::Make(
21432212
properties_, schema_, &row_groups_.back());
@@ -2182,6 +2251,9 @@ class FileMetaDataBuilder::FileMetaDataBuilderImpl {
21822251

21832252
std::unique_ptr<FileMetaData> Finish(
21842253
const std::shared_ptr<const KeyValueMetadata>& key_value_metadata) {
2254+
// Accumulate NaN counts from the last row group.
2255+
accumulateNaNCountsFromCurrentRowGroup();
2256+
21852257
int64_t total_rows = 0;
21862258
for (auto row_group : row_groups_) {
21872259
total_rows += row_group.num_rows;
@@ -2259,6 +2331,8 @@ class FileMetaDataBuilder::FileMetaDataBuilderImpl {
22592331
file_meta_data->impl_->metadata_ = std::move(metadata_);
22602332
file_meta_data->impl_->InitSchema();
22612333
file_meta_data->impl_->InitKeyValueMetadata();
2334+
// Pass total NaN counts per field ID to FileMetaData.
2335+
file_meta_data->impl_->setNaNCounts(std::move(field_nan_counts_));
22622336
return file_meta_data;
22632337
}
22642338

@@ -2290,12 +2364,31 @@ class FileMetaDataBuilder::FileMetaDataBuilderImpl {
22902364
crypto_metadata_;
22912365

22922366
private:
2367+
// Helper to accumulate NaN counts from the current row group builder.
2368+
void accumulateNaNCountsFromCurrentRowGroup() {
2369+
if (!current_row_group_builder_) {
2370+
return;
2371+
}
2372+
auto rg_nan_counts = current_row_group_builder_->nan_counts();
2373+
// Accumulate NaN counts from this row group (keyed by field ID).
2374+
for (const auto& [fieldId, countPair] : rg_nan_counts) {
2375+
const auto& [count, has_count] = countPair;
2376+
if (has_count) {
2377+
field_nan_counts_[fieldId].first += count;
2378+
field_nan_counts_[fieldId].second = true;
2379+
}
2380+
}
2381+
}
2382+
22932383
const std::shared_ptr<WriterProperties> properties_;
22942384
std::vector<facebook::velox::parquet::thrift::RowGroup> row_groups_;
22952385

22962386
std::unique_ptr<RowGroupMetaDataBuilder> current_row_group_builder_;
22972387
const SchemaDescriptor* schema_;
22982388
std::shared_ptr<const KeyValueMetadata> key_value_metadata_;
2389+
// Total NaN counts per field ID across all row groups: field_id ->
2390+
// (nan_count, has_nan_count).
2391+
std::unordered_map<int32_t, std::pair<int64_t, bool>> field_nan_counts_;
22992392
};
23002393

23012394
std::unique_ptr<FileMetaDataBuilder> FileMetaDataBuilder::Make(

velox/dwio/parquet/writer/arrow/Metadata.h

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
#include <memory>
2424
#include <optional>
2525
#include <string>
26+
#include <unordered_map>
2627
#include <utility>
2728
#include <vector>
2829

@@ -416,6 +417,12 @@ class PARQUET_EXPORT FileMetaData {
416417
std::shared_ptr<FileMetaData> Subset(
417418
const std::vector<int>& row_groups) const;
418419

420+
/// \brief Get total NaN count for a specific field ID across all row groups.
421+
/// Returns a pair of (nan_count, has_nan_count).
422+
/// NaN counts are collected during writing but not written to the parquet
423+
/// file.
424+
std::pair<int64_t, bool> getNaNCount(int32_t fieldId) const;
425+
419426
private:
420427
friend FileMetaDataBuilder;
421428
friend class SerializedFile;
@@ -486,6 +493,13 @@ class PARQUET_EXPORT ColumnChunkMetaDataBuilder {
486493
const ColumnDescriptor* descr() const;
487494

488495
int64_t total_compressed_size() const;
496+
497+
// NaN count accessors - NaN counts are collected during writing but not
498+
// written to the parquet file.
499+
int64_t nan_count() const;
500+
501+
bool has_nan_count() const;
502+
489503
// commit the metadata
490504

491505
void Finish(
@@ -537,6 +551,10 @@ class PARQUET_EXPORT RowGroupMetaDataBuilder {
537551

538552
void set_num_rows(int64_t num_rows);
539553

554+
// Get NaN counts for all columns in current row group.
555+
// Returns a map of field_id -> (nan_count, has_nan_count).
556+
std::unordered_map<int32_t, std::pair<int64_t, bool>> nan_counts() const;
557+
540558
// commit the metadata
541559
void Finish(int64_t total_bytes_written, int16_t row_group_ordinal = -1);
542560

0 commit comments

Comments
 (0)