From 072b699dba7e9f0f19222f53bc5cc782daf4ae49 Mon Sep 17 00:00:00 2001 From: Nicholas Jiang Date: Sun, 4 Oct 2026 14:07:22 +0800 Subject: [PATCH] perf(parquet): derive data file stats from in-memory writer metadata Closing a Parquet data file used to reopen the file just written, read its footer back and decode it again to build the column stats of DataFileMeta. On remote storage this costs an open, a length lookup and ranged reads for every rolled file. ParquetFormatWriter now keeps the FileMetaData that parquet::arrow::FileWriter::metadata() returns after Close() and derives the stats from it through the conversion shared with ParquetStatsExtractor, so they stay identical to the ones read back from the footer. DataFileWriter and KeyValueDataFileWriter take the stats from a format writer implementing the new public WrittenFileStatsProvider interface in include/paimon/format/, and fall back to re-reading the file through FormatStatsExtractor for the other formats. --- docs/source/api/file_format.rst | 4 + .../format/written_file_stats_provider.h | 48 ++++++ src/paimon/core/io/data_file_writer.cpp | 2 +- src/paimon/core/io/data_file_writer_base.h | 16 ++ .../core/io/data_file_writer_base_test.cpp | 142 ++++++++++++++++++ .../core/io/key_value_data_file_writer.cpp | 2 +- src/paimon/core/io/single_file_writer.h | 6 + .../format/parquet/parquet_format_writer.cpp | 16 ++ .../format/parquet/parquet_format_writer.h | 10 +- .../parquet/parquet_stats_extractor.cpp | 20 ++- .../format/parquet/parquet_stats_extractor.h | 7 + .../parquet/parquet_stats_extractor_test.cpp | 122 +++++++++++++++ 12 files changed, 386 insertions(+), 9 deletions(-) create mode 100644 include/paimon/format/written_file_stats_provider.h diff --git a/docs/source/api/file_format.rst b/docs/source/api/file_format.rst index 23c88cc6c..f90e07f22 100644 --- a/docs/source/api/file_format.rst +++ b/docs/source/api/file_format.rst @@ -51,3 +51,7 @@ Interface .. doxygenclass:: paimon::FormatStatsExtractor :members: :undoc-members: + +.. doxygenclass:: paimon::WrittenFileStatsProvider + :members: + :undoc-members: diff --git a/include/paimon/format/written_file_stats_provider.h b/include/paimon/format/written_file_stats_provider.h new file mode 100644 index 000000000..454240021 --- /dev/null +++ b/include/paimon/format/written_file_stats_provider.h @@ -0,0 +1,48 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#pragma once + +#include + +#include "paimon/result.h" +#include "paimon/type_fwd.h" + +namespace paimon { + +/// Optionally implemented by a `FormatWriter` that can produce the column statistics of the file it +/// has finished from the metadata it still holds in memory. The data file writers then take the +/// statistics from it instead of reopening the file through `FormatStatsExtractor` to read them +/// back. +/// +/// It is implemented alongside `FormatWriter` rather than added to it, so `FormatWriter` keeps its +/// ABI and format writers that do not implement it are unaffected. +class PAIMON_EXPORT WrittenFileStatsProvider { + public: + virtual ~WrittenFileStatsProvider() = default; + + /// Extracts the statistics of each top-level column of the finished file. They must equal + /// what `FormatStatsExtractor::Extract()` reads back from that file. + /// + /// @param pool Memory pool used to build the statistics. + /// @return The statistics, or an error status if the file has not been finished successfully. + virtual Result ExtractWrittenFileStats( + const std::shared_ptr& pool) const = 0; +}; + +} // namespace paimon diff --git a/src/paimon/core/io/data_file_writer.cpp b/src/paimon/core/io/data_file_writer.cpp index 9275fcf25..5a1517613 100644 --- a/src/paimon/core/io/data_file_writer.cpp +++ b/src/paimon/core/io/data_file_writer.cpp @@ -79,7 +79,7 @@ Result>> DataFileWriter::GetFieldStats( assert(false); return Status::Invalid("simple stats extractor is null pointer."); } - return stats_extractor_->Extract(fs_, path_, pool_); + return ExtractFileStats(stats_extractor_, pool_); } } // namespace paimon diff --git a/src/paimon/core/io/data_file_writer_base.h b/src/paimon/core/io/data_file_writer_base.h index a0d018458..3583c633d 100644 --- a/src/paimon/core/io/data_file_writer_base.h +++ b/src/paimon/core/io/data_file_writer_base.h @@ -32,8 +32,11 @@ #include "paimon/core/io/data_file_index_writer.h" #include "paimon/core/io/data_file_meta.h" #include "paimon/core/io/single_file_writer.h" +#include "paimon/format/format_stats_extractor.h" +#include "paimon/format/written_file_stats_provider.h" #include "paimon/result.h" #include "paimon/status.h" +#include "paimon/type_fwd.h" namespace paimon { @@ -97,6 +100,19 @@ class DataFileWriterBase : public SingleFileWriter ExtractFileStats( + const std::shared_ptr& stats_extractor, + const std::shared_ptr& pool) const { + const auto* stats_provider = + dynamic_cast(this->GetFormatWriter()); + if (stats_provider != nullptr) { + return stats_provider->ExtractWrittenFileStats(pool); + } + return stats_extractor->Extract(this->fs_, this->path_, pool); + } + Status BeforeFinish() override { if (metadata_finalizer_) { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr updated_schema, diff --git a/src/paimon/core/io/data_file_writer_base_test.cpp b/src/paimon/core/io/data_file_writer_base_test.cpp index 9a3df7419..f2b5ac34f 100644 --- a/src/paimon/core/io/data_file_writer_base_test.cpp +++ b/src/paimon/core/io/data_file_writer_base_test.cpp @@ -18,24 +18,55 @@ #include "paimon/core/io/data_file_writer_base.h" +#include #include +#include +#include #include #include "arrow/api.h" #include "arrow/c/bridge.h" #include "arrow/ipc/json_simple.h" #include "gtest/gtest.h" +#include "paimon/common/table/special_fields.h" +#include "paimon/common/utils/arrow/status_utils.h" +#include "paimon/common/utils/long_counter.h" #include "paimon/core/core_options.h" +#include "paimon/core/io/append_data_file_writer_factory.h" #include "paimon/core/io/data_file_index_writer.h" #include "paimon/core/io/data_file_path_factory.h" #include "paimon/core/io/file_index_options.h" +#include "paimon/core/io/key_value_data_file_writer_factory.h" +#include "paimon/core/key_value.h" +#include "paimon/core/manifest/file_source.h" +#include "paimon/core/stats/simple_stats.h" +#include "paimon/core/stats/simple_stats_converter.h" #include "paimon/defs.h" #include "paimon/format/file_format.h" +#include "paimon/format/format_stats_extractor.h" +#include "paimon/fs/local/local_file_system.h" #include "paimon/memory/memory_pool.h" +#include "paimon/testing/utils/binary_row_generator.h" #include "paimon/testing/utils/testharness.h" namespace paimon::test { +namespace { + +class OpenCountingFileSystem : public LocalFileSystem { + public: + using LocalFileSystem::Open; + + Result> Open(const std::string& path) const override { + ++open_count; + return LocalFileSystem::Open(path); + } + + mutable int32_t open_count = 0; +}; + +} // namespace + class IndexedDataFileWriter : public DataFileWriterBase<::ArrowArray*> { public: IndexedDataFileWriter() : DataFileWriterBase(/*compression=*/"zstd", /*converter=*/nullptr) {} @@ -94,4 +125,115 @@ TEST(DataFileWriterBaseTest, WriteAfterCloseDoesNotConsumeIndexBatch) { ArrowArrayRelease(&second_batch); } +class DataFileWriterStatsTest : public ::testing::TestWithParam { + public: + void SetUp() override { + dir_ = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir_); + pool_ = GetDefaultPool(); + file_system_ = std::make_shared(); + ASSERT_OK_AND_ASSIGN( + options_, CoreOptions::FromMap({{Options::FILE_FORMAT, GetParam()}}, file_system_)); + path_factory_ = std::make_shared(); + ASSERT_OK(path_factory_->Init(dir_->Str(), GetParam(), "data-", nullptr)); + } + + int32_t ExpectedOpenCount() const { + return GetParam() == "parquet" ? 0 : 1; + } + + Result ReadBackStats(const std::shared_ptr& file_schema, + const std::string& path) const { + ::ArrowSchema c_schema; + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*file_schema, &c_schema)); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr stats_extractor, + options_.GetFileFormat()->CreateStatsExtractor(&c_schema)); + return stats_extractor->Extract(file_system_, path, pool_); + } + + protected: + std::unique_ptr dir_; + std::shared_ptr pool_; + std::shared_ptr file_system_; + CoreOptions options_; + std::shared_ptr path_factory_; +}; + +TEST_P(DataFileWriterStatsTest, AppendWriter) { + using AppendFileWriter = SingleFileWriter<::ArrowArray*, std::shared_ptr>; + std::shared_ptr schema = + arrow::schema({arrow::field("f0", arrow::int32()), arrow::field("f1", arrow::utf8()), + arrow::field("f2", arrow::float64())}); + AppendDataFileWriterFactory writer_factory(options_, /*schema_id=*/0, schema, + /*write_cols=*/std::nullopt, + std::make_shared(0), + FileSource::Append(), path_factory_, pool_); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, writer_factory.CreateWriter()); + + std::shared_ptr array = + arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(schema->fields()), R"([ + [1, "b", 1.5], + [null, "a", -0.5], + [3, null, 2.5] + ])") + .ValueOrDie(); + ::ArrowArray batch; + ASSERT_TRUE(arrow::ExportArray(*array, &batch).ok()); + ASSERT_OK(writer->Write(&batch)); + ASSERT_OK(writer->Close()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr file_meta, writer->GetResult()); + ASSERT_EQ(ExpectedOpenCount(), file_system_->open_count); + + ASSERT_OK_AND_ASSIGN(ColumnStatsVector read_back_stats, + ReadBackStats(schema, writer->GetPath())); + ASSERT_OK_AND_ASSIGN(SimpleStats expected_value_stats, + SimpleStatsConverter::ToBinary(read_back_stats, pool_.get())); + ASSERT_EQ(expected_value_stats, file_meta->value_stats); +} + +TEST_P(DataFileWriterStatsTest, KeyValueWriter) { + using KeyValueFileWriter = SingleFileWriter>; + std::shared_ptr schema = SpecialFields::CompleteSequenceAndValueKindField( + arrow::schema({arrow::field("k", arrow::int32()), arrow::field("v", arrow::utf8())})); + KeyValueDataFileWriterFactory writer_factory(options_, /*schema_id=*/0, schema, /*level=*/0, + FileSource::Append(), /*primary_keys=*/{"k"}, + path_factory_, /*create_stats_extractor=*/true, + /*is_changelog=*/false, pool_); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, writer_factory.CreateWriter()); + + std::shared_ptr array = + arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(schema->fields()), R"([ + [0, 0, 1, "b"], + [1, 0, 2, null], + [2, 0, 3, "a"] + ])") + .ValueOrDie(); + KeyValueBatch batch; + batch.min_sequence_number = 0; + batch.max_sequence_number = 2; + batch.min_key = BinaryRowGenerator::GenerateRowPtr({1}, pool_.get()); + batch.max_key = BinaryRowGenerator::GenerateRowPtr({3}, pool_.get()); + batch.batch = std::make_unique<::ArrowArray>(); + ASSERT_TRUE(arrow::ExportArray(*array, batch.batch.get()).ok()); + ASSERT_OK(writer->Write(std::move(batch))); + ASSERT_OK(writer->Close()); + ASSERT_OK_AND_ASSIGN(std::shared_ptr file_meta, writer->GetResult()); + ASSERT_EQ(ExpectedOpenCount(), file_system_->open_count); + + ASSERT_OK_AND_ASSIGN(ColumnStatsVector read_back_stats, + ReadBackStats(schema, writer->GetPath())); + ColumnStatsVector key_column_stats = {read_back_stats[schema->GetFieldIndex("k")]}; + ColumnStatsVector value_column_stats( + read_back_stats.begin() + SpecialFields::KEY_VALUE_SPECIAL_FIELD_COUNT, + read_back_stats.end()); + ASSERT_OK_AND_ASSIGN(SimpleStats expected_key_stats, + SimpleStatsConverter::ToBinary(key_column_stats, pool_.get())); + ASSERT_EQ(expected_key_stats, file_meta->key_stats); + ASSERT_OK_AND_ASSIGN(SimpleStats expected_value_stats, + SimpleStatsConverter::ToBinary(value_column_stats, pool_.get())); + ASSERT_EQ(expected_value_stats, file_meta->value_stats); +} + +INSTANTIATE_TEST_SUITE_P(FileFormat, DataFileWriterStatsTest, ::testing::Values("parquet", "orc")); + } // namespace paimon::test diff --git a/src/paimon/core/io/key_value_data_file_writer.cpp b/src/paimon/core/io/key_value_data_file_writer.cpp index 881712ebe..ded5422c1 100644 --- a/src/paimon/core/io/key_value_data_file_writer.cpp +++ b/src/paimon/core/io/key_value_data_file_writer.cpp @@ -197,7 +197,7 @@ Result>> KeyValueDataFileWriter::GetFie assert(false); return Status::Invalid("simple stats extractor is null pointer."); } - return stats_extractor_->Extract(fs_, path_, pool_); + return ExtractFileStats(stats_extractor_, pool_); } } // namespace paimon diff --git a/src/paimon/core/io/single_file_writer.h b/src/paimon/core/io/single_file_writer.h index 6dd1db9b8..f14421c37 100644 --- a/src/paimon/core/io/single_file_writer.h +++ b/src/paimon/core/io/single_file_writer.h @@ -155,6 +155,12 @@ class SingleFileWriter : public FileWriter { /// Serializes schema and forwards it as file metadata to FormatWriter. Status UpdateSchema(const std::shared_ptr& schema); + /// Returns the format writer of this file, which outlives Close(), or nullptr before Init() + /// has succeeded. + const FormatWriter* GetFormatWriter() const { + return writer_.get(); + } + int64_t output_bytes_ = -1; std::string compression_; std::function converter_; diff --git a/src/paimon/format/parquet/parquet_format_writer.cpp b/src/paimon/format/parquet/parquet_format_writer.cpp index 768615a97..32e0f89c4 100644 --- a/src/paimon/format/parquet/parquet_format_writer.cpp +++ b/src/paimon/format/parquet/parquet_format_writer.cpp @@ -37,8 +37,11 @@ #include "paimon/common/utils/arrow/status_utils.h" #include "paimon/common/utils/checked_cast.h" #include "paimon/core/casting/casting_utils.h" +#include "paimon/format/format_stats_extractor.h" #include "paimon/format/parquet/parquet_format_defs.h" +#include "paimon/format/parquet/parquet_stats_extractor.h" #include "parquet/arrow/writer.h" +#include "parquet/metadata.h" #include "parquet/properties.h" namespace arrow { @@ -130,9 +133,22 @@ Status ParquetFormatWriter::Flush() { Status ParquetFormatWriter::Finish() { PAIMON_RETURN_NOT_OK(Flush()); PAIMON_RETURN_NOT_OK_FROM_ARROW(writer_->Close()); + file_metadata_ = writer_->metadata(); return Status::OK(); } +Result ParquetFormatWriter::ExtractWrittenFileStats( + const std::shared_ptr& pool) const { + if (!file_metadata_) { + return Status::Invalid("Cannot extract stats before the parquet file is finished."); + } + using StatsWithFileInfo = std::pair; + ParquetStatsExtractor stats_extractor(schema_); + PAIMON_ASSIGN_OR_RAISE(StatsWithFileInfo result, + stats_extractor.ExtractFromMetadata(*file_metadata_, pool)); + return std::move(result.first); +} + Status ParquetFormatWriter::AddMetadata(const std::map& metadata) { if (metadata.empty()) { return Status::OK(); diff --git a/src/paimon/format/parquet/parquet_format_writer.h b/src/paimon/format/parquet/parquet_format_writer.h index 48f8253de..bcd2affd2 100644 --- a/src/paimon/format/parquet/parquet_format_writer.h +++ b/src/paimon/format/parquet/parquet_format_writer.h @@ -24,10 +24,12 @@ #include #include "paimon/format/format_writer.h" +#include "paimon/format/written_file_stats_provider.h" #include "paimon/fs/file_system.h" #include "paimon/metrics.h" #include "paimon/result.h" #include "paimon/status.h" +#include "paimon/type_fwd.h" #include "parquet/arrow/writer.h" namespace arrow { @@ -42,13 +44,14 @@ class OutputStream; class ArrowOutputStreamAdapter; } // namespace paimon namespace parquet { +class FileMetaData; class WriterProperties; } // namespace parquet struct ArrowArray; namespace paimon::parquet { -class ParquetFormatWriter : public FormatWriter { +class ParquetFormatWriter : public FormatWriter, public WrittenFileStatsProvider { public: static Result> Create( const std::shared_ptr& output_stream, @@ -70,6 +73,9 @@ class ParquetFormatWriter : public FormatWriter { Status AddMetadata(const std::map& metadata) override; + Result ExtractWrittenFileStats( + const std::shared_ptr& pool) const override; + private: ParquetFormatWriter(std::unique_ptr<::parquet::arrow::FileWriter> writer, const std::shared_ptr& out, @@ -103,6 +109,8 @@ class ParquetFormatWriter : public FormatWriter { // Struct view of schema_, matched against the layout of each incoming batch. std::shared_ptr logical_struct_type_; std::shared_ptr metrics_; + // Footer metadata of the file, set once Finish() has succeeded. + std::shared_ptr<::parquet::FileMetaData> file_metadata_; int64_t total_records_written_ = 0; uint64_t max_memory_use_; }; diff --git a/src/paimon/format/parquet/parquet_stats_extractor.cpp b/src/paimon/format/parquet/parquet_stats_extractor.cpp index 8dfe7faac..ddca47a69 100644 --- a/src/paimon/format/parquet/parquet_stats_extractor.cpp +++ b/src/paimon/format/parquet/parquet_stats_extractor.cpp @@ -270,17 +270,25 @@ ParquetStatsExtractor::ExtractWithFileInfo(const std::shared_ptr& fi std::shared_ptr<::parquet::FileMetaData> file_metadata = file_reader_builder.raw_reader()->metadata(); - int32_t field_count = file_metadata->schema()->group_node()->field_count(); + return ExtractFromMetadata(*file_metadata, pool); +} + +Result> +ParquetStatsExtractor::ExtractFromMetadata(const ::parquet::FileMetaData& file_metadata, + const std::shared_ptr& pool) const { + int32_t field_count = file_metadata.schema()->group_node()->field_count(); ColumnStatsVector result_stats; result_stats.reserve(field_count); std::unordered_map> merged_stats; - for (int32_t row_group_idx = 0; row_group_idx < file_metadata->num_row_groups(); + for (int32_t row_group_idx = 0; row_group_idx < file_metadata.num_row_groups(); ++row_group_idx) { - for (int32_t col_idx = 0; col_idx < file_metadata->num_columns(); ++col_idx) { - auto column_chunk = file_metadata->RowGroup(row_group_idx)->ColumnChunk(col_idx); + std::unique_ptr<::parquet::RowGroupMetaData> row_group = + file_metadata.RowGroup(row_group_idx); + for (int32_t col_idx = 0; col_idx < row_group->num_columns(); ++col_idx) { + auto column_chunk = row_group->ColumnChunk(col_idx); if (!column_chunk->is_stats_set()) { continue; } @@ -291,7 +299,7 @@ ParquetStatsExtractor::ExtractWithFileInfo(const std::shared_ptr& fi } for (int32_t field_idx = 0; field_idx < field_count; ++field_idx) { - auto node = file_metadata->schema()->group_node()->field(field_idx); + auto node = file_metadata.schema()->group_node()->field(field_idx); if (node->is_group()) { // nested type do not have parquet stats const auto& logical_type = node->logical_type(); @@ -318,7 +326,7 @@ ParquetStatsExtractor::ExtractWithFileInfo(const std::shared_ptr& fi result_stats.push_back(col_stats); } } - return std::make_pair(std::move(result_stats), FileInfo(file_metadata->num_rows())); + return std::make_pair(std::move(result_stats), FileInfo(file_metadata.num_rows())); } } // namespace paimon::parquet diff --git a/src/paimon/format/parquet/parquet_stats_extractor.h b/src/paimon/format/parquet/parquet_stats_extractor.h index 0fc852b5f..5391c27b0 100644 --- a/src/paimon/format/parquet/parquet_stats_extractor.h +++ b/src/paimon/format/parquet/parquet_stats_extractor.h @@ -67,6 +67,13 @@ class ParquetStatsExtractor : public FormatStatsExtractor { const std::shared_ptr& file_system, const std::string& path, const std::shared_ptr& pool) override; + /// Extracts the statistics of each column and the `FileInfo` from the metadata of a Parquet + /// file, whether decoded from the footer of the file or still held by the writer that has + /// just closed it. + Result> ExtractFromMetadata( + const ::parquet::FileMetaData& file_metadata, + const std::shared_ptr& pool) const; + private: void PrintConvertedType(const ::parquet::schema::Node* node); diff --git a/src/paimon/format/parquet/parquet_stats_extractor_test.cpp b/src/paimon/format/parquet/parquet_stats_extractor_test.cpp index 65998426d..11fe8f05f 100644 --- a/src/paimon/format/parquet/parquet_stats_extractor_test.cpp +++ b/src/paimon/format/parquet/parquet_stats_extractor_test.cpp @@ -49,6 +49,8 @@ #include "paimon/status.h" #include "paimon/testing/utils/testharness.h" #include "parquet/arrow/reader.h" +#include "parquet/file_reader.h" +#include "parquet/metadata.h" #include "parquet/properties.h" namespace paimon::parquet::test { @@ -91,6 +93,7 @@ class ParquetStatsExtractorTest : public ::testing::Test { ASSERT_OK_AND_ASSIGN(auto result, stats_extractor.ExtractWithFileInfo(fs, file_path, GetDefaultPool())); auto& col_stats_vec = result.first; + CheckWrittenFileStats(*format_writer, col_stats_vec); ASSERT_EQ(fields.size(), col_stats_vec.size()); ASSERT_EQ(col_stats_vec.size(), expected_stats.size()); for (size_t i = 0; i < expected_stats.size(); i++) { @@ -106,6 +109,23 @@ class ParquetStatsExtractorTest : public ::testing::Test { ASSERT_EQ(row_count, expect_row_count); } + static void CheckWrittenFileStats(const ParquetFormatWriter& format_writer, + const ColumnStatsVector& read_back_stats) { + std::shared_ptr pool = GetDefaultPool(); + ASSERT_OK_AND_ASSIGN(ColumnStatsVector written_stats, + format_writer.ExtractWrittenFileStats(pool)); + ASSERT_EQ(read_back_stats.size(), written_stats.size()); + for (size_t i = 0; i < read_back_stats.size(); ++i) { + ASSERT_EQ(read_back_stats[i]->GetFieldType(), written_stats[i]->GetFieldType()) << i; + ASSERT_EQ(read_back_stats[i]->ToString(), written_stats[i]->ToString()) << i; + } + ASSERT_OK_AND_ASSIGN(SimpleStats read_back_simple_stats, + SimpleStatsConverter::ToBinary(read_back_stats, pool.get())); + ASSERT_OK_AND_ASSIGN(SimpleStats written_simple_stats, + SimpleStatsConverter::ToBinary(written_stats, pool.get())); + ASSERT_EQ(read_back_simple_stats, written_simple_stats); + } + private: std::unique_ptr dir_; }; @@ -311,6 +331,7 @@ TEST_F(ParquetStatsExtractorTest, TestNullForAllType) { auto column_stats = ret.first; auto file_info = ret.second; + CheckWrittenFileStats(*format_writer, column_stats); ASSERT_EQ(src_array->length(), file_info.GetRowCount()); ASSERT_OK_AND_ASSIGN(auto stats, SimpleStatsConverter::ToBinary(column_stats, pool.get())); // test compatible with java @@ -362,4 +383,105 @@ TEST_F(ParquetStatsExtractorTest, TestExtractStatsTimestampType) { CheckStats(fields, data_str, expected_stats_str, /*expect_row_count=*/1); } } + +TEST_F(ParquetStatsExtractorTest, TestWrittenFileStatsAcrossRowGroups) { + arrow::FieldVector fields = { + arrow::field("f_int", arrow::int32()), + arrow::field("f_bigint", arrow::int64()), + arrow::field("f_float", arrow::float32()), + arrow::field("f_double", arrow::float64()), + arrow::field("f_string", arrow::utf8()), + arrow::field("f_binary", arrow::binary()), + arrow::field("f_decimal_int32", arrow::decimal128(9, 2)), + arrow::field("f_decimal_int64", arrow::decimal128(18, 4)), + arrow::field("f_decimal_fixed", arrow::decimal128(30, 2)), + arrow::field("f_ts_int96", arrow::timestamp(arrow::TimeUnit::NANO)), + arrow::field("f_ts_micro", arrow::timestamp(arrow::TimeUnit::MICRO)), + arrow::field("f_date", arrow::date32()), + arrow::field("f_all_null", arrow::int32()), + arrow::field("f_list", arrow::list(arrow::int32())), + arrow::field("f_struct", arrow::struct_({arrow::field("f0", arrow::int32())})), + arrow::field("f_map", + arrow::map(arrow::utf8(), arrow::map(arrow::int32(), arrow::int64()))), + }; + std::string data_str = R"([ + [3, 30, NaN, 0.5, "b", "y", "1.50", "10.0001", "12345678901234567890.12", "1970-01-01 00:00:01", "1970-01-01 00:00:00.000001", 10, null, [1, 2], [1], [["a", [[1, 10], [2, 20]]], ["b", []]]], + [1, -10, 1.5, NaN, "a", "x", "-2.25", "-3.5000", "-1.00", "1970-01-01 00:00:02", "1970-01-01 00:00:00.000002", 12, null, [3], [2], [["c", null]]], + [null, 20, NaN, NaN, "d", null, null, "7.0000", "0.01", null, null, null, null, null, null, null], + [5, null, NaN, -1.25, null, "z", "3.00", null, null, "1970-01-01 00:00:03", "1970-01-01 00:00:00.000003", 11, null, [], [3], []], + [2, 40, -2.5, 2.75, "c", "w", "0.75", "1.2500", "99.99", "1970-01-01 00:00:04", "1970-01-01 00:00:00.000004", 9, null, [4], [4], [["d", [[3, null]]]]] + ])"; + auto schema = arrow::schema(fields); + std::shared_ptr fs = std::make_shared(); + std::string file_name = dir_->Str() + "/multiple_row_groups.parquet"; + ASSERT_OK_AND_ASSIGN(std::shared_ptr out, + fs->Create(file_name, /*overwrite=*/false)); + std::shared_ptr pool = GetDefaultPool(); + ::parquet::WriterProperties::Builder builder; + builder.enable_store_decimal_as_integer(); + builder.max_row_group_length(2); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr format_writer, + ParquetFormatWriter::Create(out, schema, builder.build(), + DEFAULT_PARQUET_WRITER_MAX_MEMORY_USE, GetArrowPool(pool))); + std::shared_ptr array = + arrow::ipc::internal::json::ArrayFromJSON(arrow::struct_(fields), data_str).ValueOrDie(); + ArrowArray c_array; + ASSERT_TRUE(arrow::ExportArray(*array, &c_array).ok()); + ASSERT_OK(format_writer->AddBatch(&c_array)); + ASSERT_NOK_WITH_MSG(format_writer->ExtractWrittenFileStats(pool), + "Cannot extract stats before the parquet file is finished"); + ASSERT_OK(format_writer->Finish()); + ASSERT_OK(out->Flush()); + ASSERT_OK(out->Close()); + + std::unique_ptr<::parquet::ParquetFileReader> parquet_reader = + ::parquet::ParquetFileReader::OpenFile(file_name); + ASSERT_EQ(3, parquet_reader->metadata()->num_row_groups()); + + ParquetStatsExtractor extractor(schema); + ASSERT_OK_AND_ASSIGN(auto result, extractor.ExtractWithFileInfo(fs, file_name, pool)); + ASSERT_EQ(5, result.second.GetRowCount()); + const ColumnStatsVector& column_stats = result.first; + ASSERT_EQ(fields.size(), column_stats.size()); + ASSERT_EQ("min 1, max 5, null count 1", column_stats[0]->ToString()); + ASSERT_EQ("min -2.5, max 1.5, null count 0", column_stats[2]->ToString()); + ASSERT_EQ("min a, max d, null count 1", column_stats[4]->ToString()); + ASSERT_EQ("min null, max null, null count 5", column_stats[12]->ToString()); + ASSERT_EQ(FieldType::MAP, column_stats[15]->GetFieldType()); + ASSERT_EQ("min null, max null, null count null", column_stats[15]->ToString()); + CheckWrittenFileStats(*format_writer, column_stats); +} + +TEST_F(ParquetStatsExtractorTest, TestWrittenFileStatsOfEmptyFile) { + arrow::FieldVector fields = { + arrow::field("f_int", arrow::int32()), + arrow::field("f_string", arrow::utf8()), + arrow::field("f_list", arrow::list(arrow::int32())), + }; + auto schema = arrow::schema(fields); + std::shared_ptr fs = std::make_shared(); + std::string file_name = dir_->Str() + "/empty.parquet"; + ASSERT_OK_AND_ASSIGN(std::shared_ptr out, + fs->Create(file_name, /*overwrite=*/false)); + std::shared_ptr pool = GetDefaultPool(); + ::parquet::WriterProperties::Builder builder; + ASSERT_OK_AND_ASSIGN( + std::unique_ptr format_writer, + ParquetFormatWriter::Create(out, schema, builder.build(), + DEFAULT_PARQUET_WRITER_MAX_MEMORY_USE, GetArrowPool(pool))); + ASSERT_OK(format_writer->Finish()); + ASSERT_OK(out->Flush()); + ASSERT_OK(out->Close()); + + ParquetStatsExtractor extractor(schema); + ASSERT_OK_AND_ASSIGN(auto result, extractor.ExtractWithFileInfo(fs, file_name, pool)); + ASSERT_EQ(0, result.second.GetRowCount()); + const ColumnStatsVector& column_stats = result.first; + ASSERT_EQ(fields.size(), column_stats.size()); + for (const auto& stats : column_stats) { + ASSERT_EQ("min null, max null, null count null", stats->ToString()); + } + CheckWrittenFileStats(*format_writer, column_stats); +} } // namespace paimon::parquet::test