From 4f32e09056f520bc9eab2db31fa610c705527b30 Mon Sep 17 00:00:00 2001 From: gripleaf <425797155@qq.com> Date: Thu, 8 Oct 2026 14:23:27 +0800 Subject: [PATCH 1/5] perf(manifest): prune entries by row range and cache decoded Arrow IPC Filter aligned manifest rows by row ID range before constructing file metadata, while preserving version checks and the existing entry filters. Replace raw manifest byte caching with a shared Arrow IPC representation in ObjectsFile for manifest files, manifest lists, and index manifests. Reuse whole-file cache keys, coalesce concurrent loads, and preserve allocator ownership and per-reader filtering. Add regression coverage and metrics for pruning, cache reuse, schema compatibility, concurrent readers, eviction, and error handling. Document the full-file decoding cost of cold cached reads. --- docs/source/user_guide/metrics.rst | 33 + include/paimon/table/source/scan_metrics.h | 10 + src/paimon/CMakeLists.txt | 1 + .../index_manifest_file_handler_test.cpp | 53 +- src/paimon/core/manifest/manifest_file.cpp | 83 +++ src/paimon/core/manifest/manifest_file.h | 10 + .../core/manifest/manifest_file_test.cpp | 6 +- .../core/manifest/manifest_list_test.cpp | 79 +- .../core/manifest/manifest_row_range_test.cpp | 673 ++++++++++++++++++ .../append_only_file_store_scan_test.cpp | 45 +- src/paimon/core/operation/file_store_scan.cpp | 9 + src/paimon/core/operation/file_store_scan.h | 1 + src/paimon/core/utils/manifest_meta_reader.h | 5 +- src/paimon/core/utils/objects_file.h | 246 +++++-- 14 files changed, 1178 insertions(+), 76 deletions(-) create mode 100644 src/paimon/core/manifest/manifest_row_range_test.cpp diff --git a/docs/source/user_guide/metrics.rst b/docs/source/user_guide/metrics.rst index 72ee2a176..7a3997fb1 100644 --- a/docs/source/user_guide/metrics.rst +++ b/docs/source/user_guide/metrics.rst @@ -114,3 +114,36 @@ These metrics are C++-only and have no counterparts in Java Paimon. "io.async.pending", "gauge", "requests", "Asynchronous callbacks not yet completed" "io.async.latency.count", "counter", "requests", "Completed asynchronous callback latency samples" "io.async.latency.sum-us", "counter", "microseconds", "Sum of asynchronous callback latency" + +Manifest reads +~~~~~~~~~~~~~~ + +Data-evolution scans with row ID ranges prune non-overlapping entries before +materializing file metadata. Unknown ranges are retained. + +When a caller cache is provided, ordinary, bucket and row-range manifest reads +share one Arrow IPC representation of each immutable manifest, replacing the raw +manifest byte cache. Manifest lists and index manifests use the same decoded IPC +cache through ``ObjectsFile``, with the existing whole-file keys and cache budget. +Cache keys do not contain query filters or mutable snapshot +selection. Concurrent loads of the same path through the same cache are coalesced +across readers; each reader applies its own filters. + +Cold reads decode the complete manifest to populate the shared cache and consume +the original batches without reading the IPC stream back. Warm reads skip the +source format decoder. Without a cache, or when the cache rejects manifest reads, +bucket reads retain source-level selective decoding when supported. Cold bucket +reads with a cache may therefore decode more entries than uncached selective +reads. Cache entries retain their buffers and allocator until the last reader +releases them, including after eviction. + +The scan metrics additionally expose cumulative ``rowRangeManifestEntriesScanned``, +``rowRangeManifestEntriesPruned`` and ``rowRangeManifestEntriesMaterialized``. +Materialized entries are counted before the ordinary entry filter. +``manifestArrowCacheHits``, ``manifestArrowCacheMisses`` and +``manifestArrowCacheFallbacks`` describe the decoded cache path for all three manifest read modes. +Manifest list and index manifest readers expose the same counters through their own read metrics. +Readers reusing a concurrent load count as cache hits. +``rowRangeManifestReadDuration`` is a per-file duration histogram in milliseconds, +including cache access, filtering and materialization, also on failed reads. +Parallel file durations overlap and must not be added to obtain request latency. diff --git a/include/paimon/table/source/scan_metrics.h b/include/paimon/table/source/scan_metrics.h index 15bc41ff1..4296d5353 100644 --- a/include/paimon/table/source/scan_metrics.h +++ b/include/paimon/table/source/scan_metrics.h @@ -25,6 +25,16 @@ namespace paimon { /// Metric names for scan planning operations. class PAIMON_EXPORT ScanMetrics { public: + // Cumulative work for row-range manifest reads; duration histogram samples are per file (ms). + static constexpr char ROW_RANGE_MANIFEST_ENTRIES_SCANNED[] = "rowRangeManifestEntriesScanned"; + static constexpr char ROW_RANGE_MANIFEST_ENTRIES_PRUNED[] = "rowRangeManifestEntriesPruned"; + static constexpr char ROW_RANGE_MANIFEST_ENTRIES_MATERIALIZED[] = + "rowRangeManifestEntriesMaterialized"; + static constexpr char ROW_RANGE_MANIFEST_READ_DURATION[] = "rowRangeManifestReadDuration"; + static constexpr char MANIFEST_ARROW_CACHE_HITS[] = "manifestArrowCacheHits"; + static constexpr char MANIFEST_ARROW_CACHE_MISSES[] = "manifestArrowCacheMisses"; + static constexpr char MANIFEST_ARROW_CACHE_FALLBACKS[] = "manifestArrowCacheFallbacks"; + static constexpr char LAST_SCAN_DURATION[] = "lastScanDuration"; // Histogram metric for scan plan duration (milliseconds). static constexpr char SCAN_DURATION[] = "scanDuration"; diff --git a/src/paimon/CMakeLists.txt b/src/paimon/CMakeLists.txt index eb17538de..ca2630e4a 100644 --- a/src/paimon/CMakeLists.txt +++ b/src/paimon/CMakeLists.txt @@ -834,6 +834,7 @@ if(PAIMON_BUILD_TESTS) core/manifest/manifest_entry_serializer_test.cpp core/manifest/manifest_file_meta_serializer_test.cpp core/manifest/manifest_file_test.cpp + core/manifest/manifest_row_range_test.cpp core/manifest/manifest_list_test.cpp core/manifest/partition_entry_test.cpp core/manifest/file_entry_test.cpp diff --git a/src/paimon/core/manifest/index_manifest_file_handler_test.cpp b/src/paimon/core/manifest/index_manifest_file_handler_test.cpp index 42048dc75..e0ce213a3 100644 --- a/src/paimon/core/manifest/index_manifest_file_handler_test.cpp +++ b/src/paimon/core/manifest/index_manifest_file_handler_test.cpp @@ -25,7 +25,10 @@ #include #include "arrow/api.h" +#include "arrow/io/memory.h" +#include "arrow/ipc/api.h" #include "gtest/gtest.h" +#include "paimon/common/utils/path_util.h" #include "paimon/core/core_options.h" #include "paimon/core/deletionvectors/deletion_vectors_index_file.h" #include "paimon/core/index/index_file_meta.h" @@ -37,6 +40,7 @@ #include "paimon/format/file_format_factory.h" #include "paimon/memory/memory_pool.h" #include "paimon/testing/utils/binary_row_generator.h" +#include "paimon/testing/utils/counting_cache_test_utils.h" #include "paimon/testing/utils/testharness.h" namespace paimon::test { @@ -50,7 +54,8 @@ class IndexManifestFileHandlerTest : public testing::Test { } Result> CreateManifestFile( - int32_t bucket_mode, const std::string& file_format_identifier = "avro") const { + int32_t bucket_mode, const std::string& file_format_identifier = "avro", + const std::shared_ptr& cache = nullptr) const { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr file_format, FileFormatFactory::Get(file_format_identifier, {})); auto schema = arrow::schema({arrow::field("f0", arrow::int32())}); @@ -63,6 +68,7 @@ class IndexManifestFileHandlerTest : public testing::Test { /*global_index_external_path=*/std::nullopt, /*index_file_in_data_file_dir=*/false, pool_)); PAIMON_ASSIGN_OR_RAISE(CoreOptions options, CoreOptions::FromMap({})); + options.WithCache(cache); return IndexManifestFile::Create(dir_->GetFileSystem(), file_format, "zstd", path_factory, bucket_mode, pool_, options); } @@ -100,6 +106,51 @@ class IndexManifestFileHandlerTest : public testing::Test { std::unique_ptr dir_; }; +TEST_F(IndexManifestFileHandlerTest, DecodedCacheReusesIpcAcrossReaders) { + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, + CreateManifestFile(2, "avro", cache)); + const std::vector expected = { + MakeDvEntry(FileKind::Add(), BinaryRow::EmptyRow(), 0, "dv-0", {"data-a", "data-b"}, 10), + MakeEntry(FileKind::Add(), BinaryRow::EmptyRow(), 1, "BTREE", "index-1", 20)}; + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, writer->WriteWithoutRolling(expected)); + std::vector cold; + ASSERT_OK(writer->Read(written.first, nullptr, written.second, &cold)); + ASSERT_EQ(expected, cold); + auto key = CacheKey::ForKind( + PathUtil::JoinPath(FileStorePathFactory::ManifestPath(dir_->Str()), written.first), + /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr cached, + cache->Get(key, + [](const std::shared_ptr&) -> Result> { + return Status::Invalid("index manifest was not cached"); + })); + ASSERT_TRUE(cached); + const auto& segment = cached->GetSegment(); + auto buffer = std::make_shared(reinterpret_cast(segment.Data()), + segment.Size()); + auto ipc = arrow::ipc::RecordBatchStreamReader::Open( + std::make_shared(buffer)); + ASSERT_TRUE(ipc.ok()) << ipc.status().ToString(); + auto batch = ipc.ValueOrDie()->Next(); + ASSERT_TRUE(batch.ok()) << batch.status().ToString(); + ASSERT_TRUE(batch.ValueOrDie()); + ASSERT_EQ(expected.size(), batch.ValueOrDie()->num_rows()); + writer->DeleteQuietly(written.first); + writer.reset(); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateManifestFile(2, "avro", cache)); + std::vector warm; + ASSERT_OK(reader->Read(written.first, nullptr, written.second, &warm)); + ASSERT_EQ(expected, warm); + ASSERT_EQ(1, cache->SupplierCallCount()); + cache->Invalidate(key); + ASSERT_EQ(0, cache->Size()); + ASSERT_NOK(reader->Read(written.first, nullptr, written.second, &warm)); +} + TEST_F(IndexManifestFileHandlerTest, GlobalCombinerDeletesThenAddsByFileName) { ASSERT_OK_AND_ASSIGN(auto index_manifest_file, CreateManifestFile(/*bucket_mode=*/4)); diff --git a/src/paimon/core/manifest/manifest_file.cpp b/src/paimon/core/manifest/manifest_file.cpp index d525f62a6..b39a2cd1b 100644 --- a/src/paimon/core/manifest/manifest_file.cpp +++ b/src/paimon/core/manifest/manifest_file.cpp @@ -30,11 +30,13 @@ #include "paimon/common/reader/late_materializing_file_batch_reader.h" #include "paimon/common/utils/arrow/status_utils.h" #include "paimon/common/utils/checked_cast.h" +#include "paimon/common/utils/scope_guard.h" #include "paimon/core/io/rolling_file_writer.h" #include "paimon/core/manifest/manifest_entry.h" #include "paimon/core/manifest/manifest_entry_serializer.h" #include "paimon/core/manifest/manifest_entry_writer_factory.h" #include "paimon/core/manifest/manifest_file_meta.h" +#include "paimon/core/utils/duration.h" #include "paimon/core/utils/file_store_path_factory.h" #include "paimon/core/utils/object_serializer.h" #include "paimon/core/utils/path_factory.h" @@ -44,6 +46,8 @@ #include "paimon/format/writer_builder.h" #include "paimon/predicate/predicate_builder.h" #include "paimon/status.h" +#include "paimon/table/source/scan_metrics.h" +#include "paimon/utils/row_range_index.h" namespace arrow { class DataType; @@ -57,6 +61,7 @@ namespace { constexpr int32_t kVersionFieldIndex = 0; constexpr int32_t kBucketFieldIndex = 3; constexpr int32_t kTotalBucketsFieldIndex = 4; + } // namespace ManifestFile::ManifestFile(const std::shared_ptr& file_system, @@ -133,6 +138,84 @@ Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t buc }); } +Status ManifestFile::ReadRowRangeEntries( + const std::string& file_name, const RowRangeIndex& row_ranges, + const std::function(const ManifestEntry&)>& filter, + std::optional file_size, std::vector* entries) const { + uint64_t scanned = 0; + uint64_t pruned = 0; + uint64_t materialized = 0; + Duration duration; + auto metrics = std::make_shared(); + // Merge once per file rather than locking a shared metric for every entry. ScopeGuard + // also accounts for work completed before an I/O or filter error. + ScopeGuard record_metrics([&]() { + metrics->SetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_SCANNED, scanned); + metrics->SetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_PRUNED, pruned); + metrics->SetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_MATERIALIZED, materialized); + metrics->ObserveHistogram(ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION, + static_cast(duration.Get())); + read_metrics_->Merge(metrics); + }); + return ReadArrowBatches( + file_name, file_size, + [this, &row_ranges, &filter, &scanned, &pruned, &materialized, + entries](const std::shared_ptr& batch) -> Status { + // ManifestMetaReader has aligned both the entry and its nested file schema. Probe + // the two range columns without allocating DataFileMeta, stats or binary keys. + constexpr int32_t kFileFieldIndex = 5; + constexpr int32_t kRowCountFieldIndex = 2; + constexpr int32_t kFirstRowIdFieldIndex = 18; + const auto& file_column = batch->field(kFileFieldIndex); + if (file_column->type_id() != arrow::Type::STRUCT) { + return Status::Invalid("Manifest entry file metadata must be a struct"); + } + auto files = checked_pointer_cast(file_column); + const auto& count_column = files->field(kRowCountFieldIndex); + const auto& first_column = files->field(kFirstRowIdFieldIndex); + if (count_column->type_id() != arrow::Type::INT64 || + first_column->type_id() != arrow::Type::INT64) { + return Status::Invalid("Manifest entry row range must contain int64 fields"); + } + auto counts = checked_pointer_cast(count_column); + auto first_ids = checked_pointer_cast(first_column); + ColumnarRow row(batch->fields(), pool_, /*row_id=*/0); + for (int64_t i = 0; i < batch->length(); ++i) { + ++scanned; + row.SetRowId(i); + PAIMON_RETURN_NOT_OK( + ManifestEntrySerializer::ValidateVersion(row.GetInt(kVersionFieldIndex))); + if (files->IsNull(i)) { + return Status::Invalid( + "ManifestEntry convert from row failed, with null DataFileMeta"); + } + if (!first_ids->IsNull(i) && !counts->IsNull(i)) { + const int64_t first = first_ids->Value(i); + const int64_t count = counts->Value(i); + // Missing or invalid ranges cannot prove an entry is irrelevant. In + // particular, do not overflow when computing the inclusive upper bound. + if (first >= 0 && count > 0 && + first <= std::numeric_limits::max() - (count - 1) && + !row_ranges.Intersects(first, first + (count - 1))) { + ++pruned; + continue; + } + } + PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry, serializer_->FromRow(row)); + ++materialized; + if (filter) { + PAIMON_ASSIGN_OR_RAISE(bool keep, filter(entry)); + if (!keep) { + continue; + } + } + entries->push_back(std::move(entry)); + } + return Status::OK(); + }, + /*prepare_reader=*/nullptr); +} + Status ManifestFile::PrepareBucketRead(int32_t bucket, const std::optional& expected_total_buckets, std::unique_ptr* reader) const { diff --git a/src/paimon/core/manifest/manifest_file.h b/src/paimon/core/manifest/manifest_file.h index 3840368f3..c37882b54 100644 --- a/src/paimon/core/manifest/manifest_file.h +++ b/src/paimon/core/manifest/manifest_file.h @@ -19,6 +19,7 @@ #pragma once #include +#include #include #include #include @@ -46,6 +47,7 @@ class WriterBuilder; class PathFactory; class ManifestFileMeta; class ManifestEntry; +class RowRangeIndex; class MemoryPool; /// This file includes several `ManifestEntry`s, representing the additional changes since last @@ -76,6 +78,14 @@ class ManifestFile : public ObjectsFile { std::optional file_size, std::vector* entries) const; + /// Read entries intersecting row ID ranges before constructing their file metadata. + /// Unknown row ranges are retained. Add and Delete entries use the same selection, and + /// the ordinary entry filter still runs on every retained entry. + Status ReadRowRangeEntries(const std::string& file_name, const RowRangeIndex& row_ranges, + const std::function(const ManifestEntry&)>& filter, + std::optional file_size, + std::vector* entries) const; + private: Status PrepareBucketRead(int32_t bucket, const std::optional& expected_total_buckets, std::unique_ptr* reader) const; diff --git a/src/paimon/core/manifest/manifest_file_test.cpp b/src/paimon/core/manifest/manifest_file_test.cpp index df5bcbd85..b0ce4b9e4 100644 --- a/src/paimon/core/manifest/manifest_file_test.cpp +++ b/src/paimon/core/manifest/manifest_file_test.cpp @@ -314,7 +314,7 @@ TEST_F(ManifestFileTest, TestManifestCacheIsDisabledWithoutInjectedCache) { ASSERT_EQ(0, counting_file_system->get_file_status_count); } -TEST_F(ManifestFileTest, TestManifestCacheReusesCachedBytes) { +TEST_F(ManifestFileTest, TestManifestCacheReusesDecodedBatches) { auto pool = GetDefaultPool(); auto counting_file_system = std::make_shared(); auto manifest_cache = @@ -500,7 +500,7 @@ TEST_F(ManifestFileTest, TestInferredBucketProbeSkipsArrowMaterialization) { manifest_file->WriteWithoutRolling({excluded, selected})); std::vector warm; ASSERT_OK(manifest_file->Read(written.first, nullptr, /*file_size=*/std::nullopt, &warm)); - // Cached bytes, if available, belong to a separate pool from reader allocations. + // Warm decoded batches, if available, belong to a separate pool from reader allocations. for (bool inferred : {false, true}) { std::shared_ptr read_pool = GetMemoryPool(); ASSERT_OK_AND_ASSIGN( @@ -571,7 +571,7 @@ TEST_F(ManifestFileTest, TestReadBucketEntriesSkipsDeserializingOtherBuckets) { WrittenFile written_file, manifest_file->WriteWithoutRolling({invalid_other_bucket, valid_target_bucket})); - // Exercise both a cold cache and reuse of the retained manifest bytes. + // Exercise both a cold cache and reuse of the decoded manifest batches. for (int32_t read = 0; read < 2; ++read) { std::vector bucket_entries; ASSERT_OK(manifest_file->ReadBucketEntries(written_file.first, /*bucket=*/0, diff --git a/src/paimon/core/manifest/manifest_list_test.cpp b/src/paimon/core/manifest/manifest_list_test.cpp index d8d255380..2a83a5514 100644 --- a/src/paimon/core/manifest/manifest_list_test.cpp +++ b/src/paimon/core/manifest/manifest_list_test.cpp @@ -22,8 +22,11 @@ #include #include +#include "arrow/io/memory.h" +#include "arrow/ipc/api.h" #include "arrow/type.h" #include "gtest/gtest.h" +#include "paimon/common/utils/path_util.h" #include "paimon/core/core_options.h" #include "paimon/core/manifest/manifest_file_meta.h" #include "paimon/core/snapshot.h" @@ -34,6 +37,7 @@ #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/counting_cache_test_utils.h" #include "paimon/testing/utils/testharness.h" namespace paimon::test { @@ -97,7 +101,8 @@ class ManifestListTest : public testing::Test { std::unique_ptr CreateManifestList( const std::shared_ptr& file_system, const std::string& file_format_str, - const std::string& root_path, const std::shared_ptr& pool) const { + const std::string& root_path, const std::shared_ptr& pool, + const std::shared_ptr& cache = nullptr) const { EXPECT_OK_AND_ASSIGN(std::shared_ptr file_format, FileFormatFactory::Get(file_format_str, {})); auto unused_schema = arrow::schema(arrow::FieldVector({arrow::field("f0", arrow::utf8())})); @@ -109,10 +114,9 @@ class ManifestListTest : public testing::Test { /*legacy_partition_name_enabled=*/true, /*external_paths=*/{}, /*global_index_external_path=*/std::nullopt, /*index_file_in_data_file_dir=*/false, pool)); - EXPECT_OK_AND_ASSIGN(CoreOptions options, CoreOptions::FromMap({})); - EXPECT_OK_AND_ASSIGN(auto manifest_list, - ManifestList::Create(file_system, file_format, "zstd", path_factory, - options.GetCache(), pool)); + EXPECT_OK_AND_ASSIGN( + auto manifest_list, + ManifestList::Create(file_system, file_format, "zstd", path_factory, cache, pool)); return manifest_list; } @@ -127,6 +131,71 @@ class ManifestListTest : public testing::Test { } }; +TEST_F(ManifestListTest, DecodedCacheReusesIpcAcrossReaders) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto pool = GetDefaultPool(); + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + auto writer = CreateManifestList(dir->GetFileSystem(), "avro", dir->Str(), pool, cache); + const std::vector expected = {MakeMeta("manifest-a", 100, 2, 0), + MakeMeta("manifest-b", 200, 0, 1)}; + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, writer->Write(expected)); + std::vector cold; + ASSERT_OK(writer->Read(written.first, nullptr, written.second, &cold)); + ASSERT_EQ(expected, cold); + auto key = CacheKey::ForKind( + PathUtil::JoinPath(FileStorePathFactory::ManifestPath(dir->Str()), written.first), + /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr cached, + cache->Get(key, + [](const std::shared_ptr&) -> Result> { + return Status::Invalid("manifest list was not cached"); + })); + ASSERT_TRUE(cached); + const auto& segment = cached->GetSegment(); + auto buffer = std::make_shared(reinterpret_cast(segment.Data()), + segment.Size()); + // The existing cache key contains IPC, not source-format bytes. + auto ipc = arrow::ipc::RecordBatchStreamReader::Open( + std::make_shared(buffer)); + ASSERT_TRUE(ipc.ok()) << ipc.status().ToString(); + auto batch = ipc.ValueOrDie()->Next(); + ASSERT_TRUE(batch.ok()) << batch.status().ToString(); + ASSERT_TRUE(batch.ValueOrDie()); + ASSERT_EQ(expected.size(), batch.ValueOrDie()->num_rows()); + writer->DeleteQuietly(written.first); + writer.reset(); + auto reader = CreateManifestList(dir->GetFileSystem(), "avro", dir->Str(), pool, cache); + std::vector warm; + ASSERT_OK(reader->Read(written.first, nullptr, written.second, &warm)); + ASSERT_EQ(expected, warm); + ASSERT_EQ(1, cache->SupplierCallCount()); + cache->Invalidate(key); + ASSERT_EQ(0, cache->Size()); + ASSERT_NOK(reader->Read(written.first, nullptr, written.second, &warm)); +} + +TEST_F(ManifestListTest, OrcDecodedCachePreservesLegacyMetadata) { + auto pool = GetDefaultPool(); + auto fs = std::make_shared(); + const std::string path = GetDataDir() + "/orc/append_09.db/append_09"; + const std::string file = "manifest-list-f2d59cb8-3ec6-4860-b34b-050b1a533416-2"; + auto uncached = CreateManifestList(fs, "orc", path, pool); + std::vector expected; + ASSERT_OK(uncached->Read(file, nullptr, std::nullopt, &expected)); + ASSERT_EQ(4, expected.size()); + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + for (int32_t attempt = 0; attempt < 2; ++attempt) { + auto reader = CreateManifestList(fs, "orc", path, pool, cache); + std::vector actual; + ASSERT_OK(reader->Read(file, nullptr, std::nullopt, &actual)); + ASSERT_EQ(expected, actual); + ASSERT_EQ(1, cache->SupplierCallCount()); + } +} + TEST_F(ManifestListTest, TestSimple) { auto pool = GetDefaultPool(); auto manifest_file_metas = diff --git a/src/paimon/core/manifest/manifest_row_range_test.cpp b/src/paimon/core/manifest/manifest_row_range_test.cpp new file mode 100644 index 000000000..5a04e5efc --- /dev/null +++ b/src/paimon/core/manifest/manifest_row_range_test.cpp @@ -0,0 +1,673 @@ +/* + * 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. + */ + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "arrow/api.h" +#include "arrow/c/bridge.h" +#include "fmt/format.h" +#include "gtest/gtest.h" +#include "paimon/common/utils/checked_cast.h" +#include "paimon/common/utils/path_util.h" +#include "paimon/common/utils/scope_guard.h" +#include "paimon/core/io/data_file_meta.h" +#include "paimon/core/io/meta_to_arrow_array_converter.h" +#include "paimon/core/manifest/manifest_entry_serializer.h" +#include "paimon/core/manifest/manifest_file.h" +#include "paimon/core/operation/file_store_scan.h" +#include "paimon/core/utils/file_store_path_factory.h" +#include "paimon/format/file_format.h" +#include "paimon/format/file_format_factory.h" +#include "paimon/format/format_writer.h" +#include "paimon/format/writer_builder.h" +#include "paimon/fs/local/local_file_system.h" +#include "paimon/table/source/scan_metrics.h" +#include "paimon/testing/utils/counting_cache_test_utils.h" +#include "paimon/testing/utils/testharness.h" +#include "paimon/utils/row_range_index.h" + +namespace paimon::test { +class RowRangeManifestFileTest : public ::testing::Test { + protected: + Result> CreateManifest( + const std::string& path, const std::shared_ptr& fs, bool cache_enabled, + const std::shared_ptr& cache = nullptr, + const std::shared_ptr& pool = GetDefaultPool()) { + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr format, + FileFormatFactory::Get("avro", {})); + auto schema = arrow::schema(arrow::FieldVector{}); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr path_factory, + FileStorePathFactory::Create(path, schema, {}, "", "avro", "data-", + true, {}, std::nullopt, false, pool)); + PAIMON_ASSIGN_OR_RAISE(CoreOptions options, + CoreOptions::FromMap({{Options::READ_BATCH_SIZE, "2"}})); + if (cache) { + options.WithCache(cache); + } else if (cache_enabled) { + options.WithCache(std::make_shared(64 * 1024 * 1024)); + } + PAIMON_RETURN_NOT_OK(fs->Mkdirs(FileStorePathFactory::ManifestPath(path))); + return ManifestFile::Create(fs, format, "null", path_factory, 1024, pool, options, schema); + } + + Result Entry(const std::string& name, std::optional first, + int64_t count, const FileKind& kind = FileKind::Add()) { + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr meta, + DataFileMeta::ForAppend(name, 100, count, SimpleStats::EmptyStats(), + 0, 0, 0, FileSource::Append(), std::nullopt, + std::nullopt, first, std::nullopt)); + return ManifestEntry(kind, BinaryRow::EmptyRow(), 0, 1, meta); + } + + Status WriteArray(const std::string& path, const std::string& name, + const std::shared_ptr& array) { + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr format, + FileFormatFactory::Get("avro", {})); + ArrowSchema schema; + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportType(*array->type(), &schema)); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr builder, + format->CreateWriterBuilder(&schema, 2)); + LocalFileSystem fs; + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr output, + fs.Create(FileStorePathFactory::ManifestPath(path) + "/" + name, false)); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr writer, + builder->Build(output, "null")); + ArrowArray batch; + PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportArray(*array, &batch)); + PAIMON_RETURN_NOT_OK(writer->AddBatch(&batch)); + PAIMON_RETURN_NOT_OK(writer->Flush()); + PAIMON_RETURN_NOT_OK(writer->Finish()); + return output->Close(); + } +}; + +TEST_F(RowRangeManifestFileTest, ArrowCacheReuseEvictionAndConcurrentReaders) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + const auto cache_kind = CacheKind::MANIFEST; + auto cache = std::make_shared(cache_kind, 1024 * 1024); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + ASSERT_OK_AND_ASSIGN(ManifestEntry a, Entry("a.parquet", 100, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry b, Entry("b.parquet", 110, 10)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({a, b})); + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(110, 110)})); + std::vector expected{b}; + auto read = [&]() -> Status { + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr reader, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + std::vector entries; + PAIMON_RETURN_NOT_OK( + reader->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, &entries)); + if (entries != expected) { + return Status::Invalid("cached manifest result differs from source"); + } + return Status::OK(); + }; + ASSERT_OK(read()); + // The manifest has one decoded entry; no raw-byte copy is retained. + ASSERT_EQ(1, cache->SupplierCallCount(cache_kind)); + // Readers own their decode state; the cache only shares immutable bytes. + std::vector> readers; + for (int32_t i = 0; i < 8; ++i) { + readers.push_back(std::async(std::launch::async, read)); + } + for (auto& reader : readers) { + ASSERT_OK(reader.get()); + } + ASSERT_EQ(1, cache->SupplierCallCount(cache_kind)); + cache->InvalidateAll(); + ASSERT_OK(read()); + ASSERT_EQ(2, cache->SupplierCallCount(cache_kind)); + // A different immutable manifest must not reuse the previous file's cached content. + ASSERT_OK_AND_ASSIGN(WrittenFile next, manifest->WriteWithoutRolling({a})); + std::vector entries; + ASSERT_OK(manifest->ReadRowRangeEntries(next.first, ranges, nullptr, next.second, &entries)); + ASSERT_TRUE(entries.empty()); + ASSERT_EQ(3, cache->SupplierCallCount(cache_kind)); + + auto unsupported_cache = + std::make_shared(CacheKind::DATA_FILE_FOOTER, 1024 * 1024); + ASSERT_OK_AND_ASSIGN(std::unique_ptr unsupported, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, unsupported_cache)); + ASSERT_OK( + unsupported->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, &entries)); + ASSERT_EQ(expected, entries); + entries.clear(); + auto small_cache = std::make_shared(1); + ASSERT_OK_AND_ASSIGN(std::unique_ptr uncached, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, small_cache)); + ASSERT_OK( + uncached->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, &entries)); + ASSERT_EQ(expected, entries); + ASSERT_LE(small_cache->GetCurrentWeight(), 1); + small_cache->InvalidateAll(); + uncached.reset(); + // Materialized entries must remain valid after both cache eviction and reader destruction. + ASSERT_EQ(expected, entries); +} + +TEST_F(RowRangeManifestFileTest, OrcCachePreservesManifestMetadata) { + auto pool = GetDefaultPool(); + auto fs = std::make_shared(); + ASSERT_OK_AND_ASSIGN(std::shared_ptr format, FileFormatFactory::Get("orc", {})); + auto schema = arrow::schema(arrow::FieldVector{}); + const std::string path = GetDataDir() + "/orc/append_09.db/append_09"; + ASSERT_OK_AND_ASSIGN(std::shared_ptr paths, + FileStorePathFactory::Create(path, schema, {}, "", "orc", "data-", true, + {}, std::nullopt, false, pool)); + ASSERT_OK_AND_ASSIGN(CoreOptions options, + CoreOptions::FromMap({{Options::READ_BATCH_SIZE, "2"}})); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr uncached, + ManifestFile::Create(fs, format, "zstd", paths, 1024, pool, options, schema)); + const std::string name = "manifest-3ea5ee21-d399-4f1c-a749-2fc63dbf0852-1"; + std::vector expected; + ASSERT_OK(uncached->Read(name, nullptr, std::nullopt, &expected)); + ASSERT_EQ(5, expected.size()); + ASSERT_EQ(1721643142456LL, expected[0].File()->creation_time.GetMillisecond()); + // This legacy manifest has unknown row IDs, so every entry must be retained. + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(0, 0)})); + std::vector without_cache; + ASSERT_OK(uncached->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &without_cache)); + ASSERT_EQ(expected, without_cache); + + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + options.WithCache(cache); + for (int32_t attempt = 0; attempt < 2; ++attempt) { + SCOPED_TRACE(attempt); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr reader, + ManifestFile::Create(fs, format, "zstd", paths, 1024, pool, options, schema)); + std::vector actual; + ASSERT_OK(reader->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual)); + ASSERT_EQ(expected, actual); + ASSERT_EQ(1, cache->SupplierCallCount(CacheKind::MANIFEST)); + ASSERT_OK_AND_ASSIGN(uint64_t count, + reader->GetReadMetrics()->GetCounter( + attempt == 0 ? ScanMetrics::MANIFEST_ARROW_CACHE_MISSES + : ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); + ASSERT_EQ(1, count); + } +} + +TEST_F(RowRangeManifestFileTest, EmptyManifestCacheReuse) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({})); + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(0, 0)})); + for (int32_t attempt = 0; attempt < 2; ++attempt) { + std::vector entries; + ASSERT_OK(manifest->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, + &entries)); + ASSERT_TRUE(entries.empty()); + ASSERT_EQ(1, cache->SupplierCallCount(CacheKind::MANIFEST)); + } +} + +TEST_F(RowRangeManifestFileTest, AllReadPathsShareOneDecodedEntry) { + for (int32_t first_reader : {0, 1, 2}) { + SCOPED_TRACE(first_reader); + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + ASSERT_OK_AND_ASSIGN(ManifestEntry a, Entry("a.parquet", 100, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry b, Entry("b.parquet", 110, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry c, Entry("c.parquet", 120, 10)); + std::vector source = {a, b, c}; + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling(source)); + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(110, 110)})); + auto read = [&](int32_t mode) { + std::vector actual; + if (mode == 0) { + ASSERT_OK(manifest->Read(written.first, nullptr, written.second, &actual)); + ASSERT_EQ(source, actual); + } else if (mode == 1) { + ASSERT_OK(manifest->ReadBucketEntries(written.first, 0, std::nullopt, + written.second, &actual)); + ASSERT_EQ(source, actual); + } else { + ASSERT_OK(manifest->ReadRowRangeEntries(written.first, ranges, nullptr, + written.second, &actual)); + ASSERT_EQ(std::vector{b}, actual); + } + }; + read(first_reader); + ASSERT_EQ(1, cache->Size()); + ASSERT_EQ(1, cache->SupplierCallCount()); + auto key = CacheKey::ForKind( + PathUtil::JoinPath(FileStorePathFactory::ManifestPath(dir->Str()), written.first), + /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr cached, + cache->Get( + key, [](const std::shared_ptr&) -> Result> { + return Status::Invalid("manifest is missing under the existing cache key"); + })); + ASSERT_TRUE(cached); + // All later paths must work entirely from the same decoded cache entry. + manifest->DeleteQuietly(written.first); + for (int32_t mode : {0, 1, 2}) { + read(mode); + } + ASSERT_EQ(1, cache->Size()); + ASSERT_EQ(1, cache->SupplierCallCount()); + // Existing callers can still invalidate a manifest by its whole-file key. + cache->Invalidate(key); + ASSERT_EQ(0, cache->Size()); + } +} + +TEST_F(RowRangeManifestFileTest, ColdFilterFailureDoesNotPoisonCacheOrRepeatCallbacks) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + ASSERT_OK_AND_ASSIGN(ManifestEntry a, Entry("a.parquet", 100, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry b, Entry("b.parquet", 110, 10)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({a, b})); + int32_t calls = 0; + auto filter = [&](const ManifestEntry&) -> Result { + ++calls; + return Status::IOError("consumer failed"); + }; + std::vector entries; + ASSERT_NOK_WITH_MSG(manifest->Read(written.first, filter, written.second, &entries), + "consumer failed"); + ASSERT_EQ(1, calls); + ASSERT_TRUE(entries.empty()); + manifest->DeleteQuietly(written.first); + ASSERT_OK(manifest->Read(written.first, nullptr, written.second, &entries)); + ASSERT_EQ(std::vector({a, b}), entries); + ASSERT_EQ(1, cache->SupplierCallCount()); +} + +TEST_F(RowRangeManifestFileTest, CachedBufferRetainsAllocatorUntilEviction) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto cache = std::make_shared(CacheKind::MANIFEST, 1024 * 1024); + std::weak_ptr weak_pool; + ASSERT_OK_AND_ASSIGN(ManifestEntry entry, Entry("a.parquet", 100, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry next, Entry("b.parquet", 110, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry last, Entry("c.parquet", 120, 10)); + const std::vector expected = {entry, next, last}; + using WrittenFile = std::pair; + WrittenFile written; + { + std::shared_ptr pool = GetMemoryPool(); + weak_pool = pool; + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache, pool)); + ASSERT_OK_AND_ASSIGN(written, manifest->WriteWithoutRolling(expected)); + std::vector entries; + ASSERT_OK(manifest->Read(written.first, nullptr, written.second, &entries)); + } + ASSERT_FALSE(weak_pool.expired()); + std::vector entries; + { + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + reader->DeleteQuietly(written.first); + auto filter = [&](const ManifestEntry&) -> Result { + cache->InvalidateAll(); + EXPECT_FALSE(weak_pool.expired()); + return true; + }; + // Eviction during the first batch must not invalidate later IPC batches. + ASSERT_OK(reader->Read(written.first, filter, written.second, &entries)); + } + ASSERT_TRUE(weak_pool.expired()); + ASSERT_EQ(expected, entries); +} + +TEST_F(RowRangeManifestFileTest, CacheAdmissionFailureReusesAlreadyDecodedBatches) { + class FailingCache : public CountingRoutingCache { + public: + FailingCache() : CountingRoutingCache(CacheKind::MANIFEST, 1024 * 1024) {} + + Result> Get( + const std::shared_ptr& key, + std::function>(const std::shared_ptr&)> + supplier) override { + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr value, supplier(key)); + after_load(); + return Status::IOError("cache admission failed"); + } + + std::function after_load; + }; + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto cache = std::make_shared(); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + ASSERT_OK_AND_ASSIGN(ManifestEntry entry, Entry("a.parquet", 100, 10)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({entry})); + // Reopening the source after an optional cache failure would now fail. + cache->after_load = [&]() { manifest->DeleteQuietly(written.first); }; + int32_t calls = 0; + auto filter = [&](const ManifestEntry&) -> Result { + ++calls; + return true; + }; + std::vector entries; + ASSERT_OK(manifest->Read(written.first, filter, written.second, &entries)); + ASSERT_EQ(1, calls); + ASSERT_EQ(std::vector{entry}, entries); + ASSERT_OK_AND_ASSIGN(uint64_t fallbacks, manifest->GetReadMetrics()->GetCounter( + ScanMetrics::MANIFEST_ARROW_CACHE_FALLBACKS)); + ASSERT_EQ(1, fallbacks); +} + +class BlockingManifestCache : public CountingRoutingCache { + public: + BlockingManifestCache() : CountingRoutingCache(CacheKind::MANIFEST, 1024 * 1024) {} + + Result> Get( + const std::shared_ptr& key, + std::function>(const std::shared_ptr&)> + supplier) override { + return CountingRoutingCache::Get(key, [&](const std::shared_ptr& supplier_key) { + { + std::unique_lock lock(mutex_); + ++loads_; + cv_.notify_all(); + cv_.wait(lock, [&]() { return released_; }); + } + return supplier(supplier_key); + }); + } + + bool WaitForLoads(int32_t count) { + std::unique_lock lock(mutex_); + return cv_.wait_for(lock, std::chrono::seconds(10), [&]() { return loads_ >= count; }); + } + + void Release() { + std::lock_guard lock(mutex_); + released_ = true; + cv_.notify_all(); + } + + private: + std::mutex mutex_; + std::condition_variable cv_; + int32_t loads_ = 0; + bool released_ = false; +}; + +TEST_F(RowRangeManifestFileTest, ConcurrentColdReadersShareLoadAcrossInstances) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto cache = std::make_shared(); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + ASSERT_OK_AND_ASSIGN(ManifestEntry entry, Entry("a.parquet", 100, 10)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({entry})); + ASSERT_OK_AND_ASSIGN(WrittenFile other, manifest->WriteWithoutRolling({entry})); + std::vector> instances; + for (int32_t i = 0; i < 9; ++i) { + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); + instances.push_back(std::move(reader)); + } + std::vector> readers; + // Release blocked suppliers before future destructors wait, including on assertion failure. + ScopeGuard release([&]() { cache->Release(); }); + std::promise start; + auto gate = start.get_future().share(); + for (int32_t i = 0; i < 9; ++i) { + readers.push_back(std::async(std::launch::async, [&, i]() -> Status { + gate.wait(); + const auto& file = i == 8 ? other : written; + std::vector entries; + PAIMON_RETURN_NOT_OK(instances[i]->Read(file.first, nullptr, file.second, &entries)); + if (entries != std::vector{entry}) { + return Status::Invalid("concurrent manifest result differs"); + } + return Status::OK(); + })); + } + start.set_value(); + // An unrelated manifest must load while the first one is blocked. + ASSERT_TRUE(cache->WaitForLoads(2)); + cache->Release(); + for (auto& reader : readers) { + ASSERT_OK(reader.get()); + } + ASSERT_EQ(2, cache->SupplierCallCount()); + ASSERT_EQ(2, cache->Size()); +} + +TEST_F(RowRangeManifestFileTest, MetricsCountPruningCacheReuseAndFilterErrors) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true)); + ASSERT_OK_AND_ASSIGN(ManifestEntry a, Entry("a.parquet", 100, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry b, Entry("b.parquet", 110, 10)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({a, b})); + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(110, 110)})); + for (int32_t i = 0; i < 2; ++i) { + std::vector entries; + ASSERT_OK(manifest->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, + &entries)); + ASSERT_EQ(std::vector{b}, entries); + } + auto metrics = manifest->GetReadMetrics(); + ASSERT_OK_AND_ASSIGN(uint64_t scanned, + metrics->GetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_SCANNED)); + ASSERT_OK_AND_ASSIGN(uint64_t pruned, + metrics->GetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_PRUNED)); + ASSERT_OK_AND_ASSIGN(uint64_t materialized, + metrics->GetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_MATERIALIZED)); + ASSERT_EQ(4, scanned); + ASSERT_EQ(2, pruned); + ASSERT_EQ(2, materialized); + ASSERT_OK_AND_ASSIGN(uint64_t hits, + metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); + ASSERT_OK_AND_ASSIGN(uint64_t misses, + metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_MISSES)); + ASSERT_EQ(1, hits); + ASSERT_EQ(1, misses); + ASSERT_OK_AND_ASSIGN(HistogramStats stats, + metrics->GetHistogramStats(ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION)); + ASSERT_EQ(2, stats.count); + std::vector entries; + ASSERT_NOK(manifest->ReadRowRangeEntries( + written.first, ranges, + [](const ManifestEntry&) -> Result { return Status::IOError("filter failure"); }, + written.second, &entries)); + ASSERT_OK_AND_ASSIGN(HistogramStats after_error, + manifest->GetReadMetrics()->GetHistogramStats( + ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION)); + ASSERT_EQ(3, after_error.count); + // Previously returned snapshots remain unchanged after subsequent reads. + ASSERT_OK_AND_ASSIGN(HistogramStats old_snapshot, + metrics->GetHistogramStats(ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION)); + ASSERT_EQ(2, old_snapshot.count); +} + +TEST_F(RowRangeManifestFileTest, BoundariesUnknownRangesAndDeleteMerging) { + for (bool cache_enabled : {false, true}) { + SCOPED_TRACE(cache_enabled); + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), cache_enabled)); + ASSERT_OK_AND_ASSIGN(ManifestEntry a, Entry("a.parquet", 100, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry b, Entry("b.parquet", 110, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry unknown, Entry("unknown.parquet", std::nullopt, 10)); + ManifestEntry deleted(FileKind::Delete(), a.Partition(), a.Bucket(), a.TotalBuckets(), + a.File()); + std::vector source = {a, b, deleted, unknown}; + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling(source)); + const std::vector> queries = {{Range(100, 100)}, + {Range(109, 109)}, + {Range(110, 110)}, + {Range(119, 119)}, + {Range(99, 99)}, + {Range(120, 120)}, + {Range(109, 110)}, + {Range(100, 100), Range(119, 119)}, + {}}; + for (const auto& query : queries) { + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create(query)); + auto filter = [&ranges](const ManifestEntry& entry) -> Result { + const auto& meta = entry.File(); + return !meta->first_row_id || + ranges.Intersects(*meta->first_row_id, + *meta->first_row_id + meta->row_count - 1); + }; + std::vector ordinary; + ASSERT_OK(manifest->Read(written.first, filter, written.second, &ordinary)); + std::vector selected; + ASSERT_OK(manifest->ReadRowRangeEntries(written.first, ranges, filter, written.second, + &selected)); + ASSERT_EQ(ordinary, selected); + std::vector ordinary_live; + std::vector selected_live; + ASSERT_OK(FileStoreScan::MergeLiveEntries(ordinary, &ordinary_live)); + ASSERT_OK(FileStoreScan::MergeLiveEntries(selected, &selected_live)); + ASSERT_EQ(ordinary_live, selected_live); + for (const auto& entry : selected_live) { + ASSERT_NE("a.parquet", entry.File()->file_name); + } + } + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(100, 100)})); + std::vector entries; + ASSERT_NOK_WITH_MSG(manifest->ReadRowRangeEntries( + written.first, ranges, + [](const ManifestEntry&) -> Result { + return Status::IOError("filter error"); + }, + written.second, &entries), + "filter error"); + ASSERT_NOK(manifest->ReadRowRangeEntries("missing-manifest", ranges, nullptr, std::nullopt, + &entries)); + } +} + +TEST_F(RowRangeManifestFileTest, UncertainAndOverflowingRangesAreRetained) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), false)); + std::vector source; + constexpr int64_t kMax = std::numeric_limits::max(); + for (const auto& range : std::vector>{ + {kMax, 2}, {kMax, 1}, {0, 0}, {-1, 10}, {0, -1}}) { + ASSERT_OK_AND_ASSIGN(ManifestEntry entry, + Entry(fmt::format("file-{}-{}.parquet", range.first, range.second), + range.first, range.second)); + source.push_back(entry); + } + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling(source)); + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(0, 0)})); + std::vector actual; + ASSERT_OK( + manifest->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, &actual)); + source.erase(source.begin() + 1); + ASSERT_EQ(source, actual); +} + +TEST_F(RowRangeManifestFileTest, SchemaEvolutionAndVersionValidation) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto pool = GetDefaultPool(); + ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, + CreateManifest(dir->Str(), dir->GetFileSystem(), true)); + ASSERT_OK_AND_ASSIGN(ManifestEntry entry, Entry("a.parquet", 100, 10)); + ManifestEntrySerializer serializer(pool); + ASSERT_OK_AND_ASSIGN(BinaryRow row, serializer.ToRow(entry)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr converter, + MetaToArrowArrayConverter::Create(serializer.GetDataType(), pool)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr array, converter->NextBatch({row})); + auto batch = checked_pointer_cast(array); + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(1000, 1000)})); + for (int32_t mode : {0, 1, 2, 3}) { + SCOPED_TRACE(mode); + auto fields = serializer.GetDataType()->fields(); + auto columns = batch->fields(); + auto file = checked_pointer_cast(columns[5]); + auto file_fields = checked_pointer_cast(file->type())->fields(); + auto file_columns = file->fields(); + if (mode == 0) { + file_fields.erase(file_fields.begin() + 18); + file_columns.erase(file_columns.begin() + 18); + } else if (mode == 1) { + std::swap(file_fields[2], file_fields[18]); + std::swap(file_columns[2], file_columns[18]); + } else { + arrow::Int32Builder versions; + ASSERT_TRUE(versions.Append(mode == 2 ? 1 : 999).ok()); + ASSERT_TRUE(versions.Finish(&columns[0]).ok()); + } + auto updated_file = arrow::StructArray::Make(file_columns, file_fields); + ASSERT_TRUE(updated_file.ok()) << updated_file.status().ToString(); + columns[5] = updated_file.ValueOrDie(); + fields[5] = fields[5]->WithType(arrow::struct_(file_fields)); + std::swap(fields[1], fields[5]); + std::swap(columns[1], columns[5]); + auto evolved_result = arrow::StructArray::Make(columns, fields); + ASSERT_TRUE(evolved_result.ok()) << evolved_result.status().ToString(); + std::shared_ptr evolved = evolved_result.ValueOrDie(); + const std::string name = fmt::format("manifest-evolved-{}", mode); + ASSERT_OK(WriteArray(dir->Str(), name, evolved)); + std::vector actual; + if (mode >= 2) { + ASSERT_NOK_WITH_MSG( + manifest->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual), + (mode == 2 ? "not compatible" : "Unsupported version: 999")); + } else { + ASSERT_OK(manifest->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual)); + ASSERT_EQ(mode == 0 ? 1 : 0, actual.size()); + if (mode == 0) { + ASSERT_FALSE(actual[0].File()->first_row_id.has_value()); + } + } + } +} + +} // namespace paimon::test diff --git a/src/paimon/core/operation/append_only_file_store_scan_test.cpp b/src/paimon/core/operation/append_only_file_store_scan_test.cpp index 908f9bbb1..1bb5bda01 100644 --- a/src/paimon/core/operation/append_only_file_store_scan_test.cpp +++ b/src/paimon/core/operation/append_only_file_store_scan_test.cpp @@ -235,7 +235,6 @@ TEST_F(AppendBucketPruningTest, ReadsSingleBucketManifestOnce) { {}, std::nullopt, false, pool_)); ASSERT_OK(fs->Mkdirs(FileStorePathFactory::ManifestPath(dir->Str()))); ASSERT_OK_AND_ASSIGN(CoreOptions options, CoreOptions::FromMap({})); - options.WithCache(std::make_shared(16 * 1024 * 1024)); ASSERT_OK_AND_ASSIGN(std::shared_ptr manifest, ManifestFile::Create(fs, format, "null", paths, 1024 * 1024, pool_, options, arrow::schema({}))); @@ -252,24 +251,34 @@ TEST_F(AppendBucketPruningTest, ReadsSingleBucketManifestOnce) { ASSERT_OK_AND_ASSIGN(BinaryRow row, serializer.ToRow(entries[0])); ASSERT_OK_AND_ASSIGN(std::shared_ptr data, converter->NextBatch({row})); auto reader_builder = std::make_shared(data); - manifest->reader_builder_ = reader_builder; options_[Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED] = "true"; - for (bool inferred : {false, true}) { - SCOPED_TRACE(inferred); - ASSERT_OK_AND_ASSIGN( - std::unique_ptr scan, - CreateScan(KeyEquals(), inferred ? std::nullopt : std::optional(bucket))); - scan->manifest_file_ = manifest; - for (bool known_bounds : {false, true}) { - SCOPED_TRACE(known_bounds); - auto meta = metas[0]; - meta.min_bucket_ = known_bounds ? std::optional(bucket) : std::nullopt; - meta.max_bucket_ = meta.min_bucket_; - reader_builder->rows_read = 0; - std::vector actual; - ASSERT_OK(scan->ReadAndMergeBucketFileEntries({meta}, bucket, &actual)); - ASSERT_EQ(actual, entries); - ASSERT_EQ(reader_builder->rows_read.load(), known_bounds ? 1 : 2); + for (bool cache_enabled : {false, true}) { + SCOPED_TRACE(cache_enabled); + options.WithCache(cache_enabled ? std::make_shared(16 * 1024 * 1024) : nullptr); + ASSERT_OK_AND_ASSIGN(manifest, ManifestFile::Create(fs, format, "null", paths, 1024 * 1024, + pool_, options, arrow::schema({}))); + manifest->reader_builder_ = reader_builder; + bool cache_warm = false; + for (bool inferred : {false, true}) { + SCOPED_TRACE(inferred); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr scan, + CreateScan(KeyEquals(), inferred ? std::nullopt : std::optional(bucket))); + scan->manifest_file_ = manifest; + for (bool known_bounds : {false, true}) { + SCOPED_TRACE(known_bounds); + auto meta = metas[0]; + meta.min_bucket_ = known_bounds ? std::optional(bucket) : std::nullopt; + meta.max_bucket_ = meta.min_bucket_; + reader_builder->rows_read = 0; + std::vector actual; + ASSERT_OK(scan->ReadAndMergeBucketFileEntries({meta}, bucket, &actual)); + ASSERT_EQ(actual, entries); + const int64_t expected_rows = + cache_enabled ? (cache_warm ? 0 : 1) : (known_bounds ? 1 : 2); + ASSERT_EQ(reader_builder->rows_read.load(), expected_rows); + cache_warm = cache_enabled; + } } } } diff --git a/src/paimon/core/operation/file_store_scan.cpp b/src/paimon/core/operation/file_store_scan.cpp index fa22ac565..f1c9f7861 100644 --- a/src/paimon/core/operation/file_store_scan.cpp +++ b/src/paimon/core/operation/file_store_scan.cpp @@ -613,6 +613,15 @@ bool FileStoreScan::FilterManifestByRowRanges(const ManifestFileMeta& manifest) Status FileStoreScan::ReadManifestFileMeta(const ManifestFileMeta& manifest, std::vector* entries) const { + if (row_range_index_ && core_options_.DataEvolutionEnabled() && + core_options_.ScanManifestEntryLazyDecodeEnabled()) { + return manifest_file_->ReadRowRangeEntries( + manifest.FileName(), row_range_index_.value(), + [this](const ManifestEntry& entry) -> Result { + return FilterManifestEntry(entry); + }, + manifest.FileSize(), entries); + } std::vector unfiltered_entries; PAIMON_RETURN_NOT_OK(manifest_file_->Read( manifest.FileName(), diff --git a/src/paimon/core/operation/file_store_scan.h b/src/paimon/core/operation/file_store_scan.h index 5fe128da1..f8e85c32e 100644 --- a/src/paimon/core/operation/file_store_scan.h +++ b/src/paimon/core/operation/file_store_scan.h @@ -174,6 +174,7 @@ class FileStoreScan { std::shared_ptr GetScanMetrics() const { auto snapshot = std::make_shared(); snapshot->Overwrite(metrics_); + snapshot->Merge(manifest_file_->GetReadMetrics()); return snapshot; } diff --git a/src/paimon/core/utils/manifest_meta_reader.h b/src/paimon/core/utils/manifest_meta_reader.h index 354132cf3..95eaaeea0 100644 --- a/src/paimon/core/utils/manifest_meta_reader.h +++ b/src/paimon/core/utils/manifest_meta_reader.h @@ -60,12 +60,13 @@ class PAIMON_EXPORT ManifestMetaReader : public BatchReader { reader_->Close(); } - private: - // fill non exist field and correct int precision + /// Align fields by name, fill missing fields and correct integer precision. + /// Other physical types, including timestamp units, are preserved. static Result> AlignArrayWithSchema( const std::shared_ptr& src_array, const std::shared_ptr& target_type, arrow::MemoryPool* pool); + private: std::unique_ptr reader_; std::shared_ptr target_type_; std::shared_ptr pool_; diff --git a/src/paimon/core/utils/objects_file.h b/src/paimon/core/utils/objects_file.h index e65ee75ab..55d048622 100644 --- a/src/paimon/core/utils/objects_file.h +++ b/src/paimon/core/utils/objects_file.h @@ -19,7 +19,11 @@ #pragma once #include +#include +#include +#include #include +#include #include #include #include @@ -27,9 +31,12 @@ #include "arrow/c/bridge.h" #include "arrow/c/helpers.h" +#include "arrow/io/memory.h" +#include "arrow/ipc/api.h" #include "fmt/format.h" #include "paimon/cache/cache.h" #include "paimon/common/data/columnar/columnar_row.h" +#include "paimon/common/metrics/metrics_impl.h" #include "paimon/common/utils/arrow/arrow_utils.h" #include "paimon/common/utils/arrow/mem_utils.h" #include "paimon/common/utils/arrow/status_utils.h" @@ -44,9 +51,8 @@ #include "paimon/format/reader_builder.h" #include "paimon/format/writer_builder.h" #include "paimon/fs/file_system.h" -#include "paimon/io/byte_array_input_stream.h" -#include "paimon/memory/bytes.h" #include "paimon/record_batch.h" +#include "paimon/table/source/scan_metrics.h" namespace paimon { /// A file which contains several `T`s, provides read and write. @@ -83,6 +89,13 @@ class ObjectsFile { Result> WriteWithoutRolling(const std::vector& records); + /// Cumulative cache and derived reader work, including concurrent reads. + std::shared_ptr GetReadMetrics() const { + auto snapshot = std::make_shared(); + snapshot->Overwrite(read_metrics_); + return snapshot; + } + protected: Status ValidateWrite() const { if (file_format_identifier_ != "avro") { @@ -92,13 +105,13 @@ class ObjectsFile { return Status::OK(); } - // Optional preparation may wrap the reader and reread the file, using either cached bytes - // or the underlying file stream. + // Cached batches are query-independent. Optional preparation only applies to uncached reads. Status ReadArrowBatches( const std::string& file_name, std::optional file_size, const std::function&)>& consumer, const std::function*)>& prepare_reader) const; + std::shared_ptr read_metrics_ = std::make_shared(); std::shared_ptr path_factory_; std::shared_ptr pool_; std::shared_ptr arrow_pool_; @@ -113,8 +126,19 @@ class ObjectsFile { std::string compression_; std::shared_ptr cache_; - Result ReadFileSegment(const std::string& file_path, - const std::optional& file_size) const; + static Result> LoadCacheEntry( + const std::shared_ptr& cache, const std::string& path, + const std::function>()>& load); + + Result> SerializeArrowBatches( + const std::string& file_name, std::optional file_size, + std::vector>* batches, + std::optional* read_status) const; + + Status ReadUncachedArrowBatches( + const std::string& file_name, std::optional file_size, + const std::function&)>& consumer, + const std::function*)>& prepare_reader) const; /// Opens the file for reading, handing over the length when the caller already has it. Result> OpenForRead(const std::string& file_path, @@ -183,37 +207,13 @@ Status ObjectsFile::Read(const std::string& file_name, } template -Status ObjectsFile::ReadArrowBatches( +Status ObjectsFile::ReadUncachedArrowBatches( const std::string& file_name, std::optional file_size, const std::function&)>& consumer, const std::function*)>& prepare_reader) const { std::string file_path = path_factory_->ToPath(file_name); - std::shared_ptr file_input_stream; - std::shared_ptr cached_bytes; - if (cache_) { - // Use a whole-file key so cache hits do not need a metadata lookup just to discover file - // length. - auto cache_key = - CacheKey::ForKind(file_path, /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); - auto supplier = - [this, &file_path, - &file_size](const std::shared_ptr&) -> Result> { - PAIMON_ASSIGN_OR_RAISE(MemorySegment segment, ReadFileSegment(file_path, file_size)); - return std::make_shared(segment, CacheCallback()); - }; - Result> cache_result = cache_->Get(cache_key, supplier); - if (cache_result.ok() && cache_result.value() && - cache_result.value()->GetSegment().Data() != nullptr) { - cached_bytes = cache_result.value()->GetSegment().GetOrCreateHeapMemory(pool_.get()); - file_input_stream = - std::make_shared(cached_bytes->data(), cached_bytes->size()); - } - } - if (!file_input_stream) { - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr unique_file_input_stream, - OpenForRead(file_path, file_size)); - file_input_stream = std::shared_ptr(std::move(unique_file_input_stream)); - } + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr file_input_stream, + OpenForRead(file_path, file_size)); PAIMON_ASSIGN_OR_RAISE(std::unique_ptr batch_reader, reader_builder_->Build(file_input_stream)); @@ -255,23 +255,175 @@ Result> ObjectsFile::OpenForRead( return file_system_->Open(file_path); } +// Only in-progress loads are retained here. Readers sharing a cache and path share one load, +// including across ObjectsFile instances; unrelated files load without holding the registry lock. template -Result ObjectsFile::ReadFileSegment( - const std::string& file_path, const std::optional& file_size) const { - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr input_stream, - OpenForRead(file_path, file_size)); - PAIMON_ASSIGN_OR_RAISE(int64_t input_length, input_stream->Length()); +Result> ObjectsFile::LoadCacheEntry( + const std::shared_ptr& cache, const std::string& path, + const std::function>()>& load) { + using LoadKey = std::pair, std::string>; + using CacheResult = Result>; + static std::mutex mutex; + static std::map> pending; + const LoadKey key(cache, path); + std::shared_ptr> promise; + std::shared_future future; + { + std::lock_guard lock(mutex); + auto iter = pending.find(key); + if (iter != pending.end()) { + future = iter->second; + } else { + promise = std::make_shared>(); + future = promise->get_future().share(); + pending.emplace(key, future); + } + } + if (!promise) { + return future.get(); + } + ScopeGuard remove_pending([&]() { + std::lock_guard lock(mutex); + pending.erase(key); + }); + CacheResult result = load(); + promise->set_value(result); + return result; +} - PAIMON_RETURN_NOT_OK(input_stream->Seek(0, FS_SEEK_SET)); - auto bytes = std::make_shared(input_length, pool_.get()); - PAIMON_ASSIGN_OR_RAISE(int64_t actual_read_size, - input_stream->Read(bytes->data(), input_length)); - if (actual_read_size != input_length) { - return Status::IOError(fmt::format( - "Unexpected EOF while reading manifest file {}, expected {} bytes, got {} bytes", - file_path, input_length, actual_read_size)); +template +Result> ObjectsFile::SerializeArrowBatches( + const std::string& file_name, std::optional file_size, + std::vector>* batches, + std::optional* read_status) const { + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( + std::shared_ptr output, + arrow::io::BufferOutputStream::Create(4096, arrow_pool_.get())); + auto write_options = arrow::ipc::IpcWriteOptions::Defaults(); + write_options.memory_pool = arrow_pool_.get(); + write_options.use_threads = false; + std::shared_ptr writer; + Status cache_status; + *read_status = ReadUncachedArrowBatches( + file_name, file_size, + [&](const std::shared_ptr& batch) -> Status { + // The loading caller consumes these original batches, without an IPC round trip. + batches->push_back(batch); + if (!cache_status.ok()) { + return Status::OK(); + } + auto write_batch = [&]() -> Status { + // Preserve physical types, including ORC nanosecond timestamps. + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( + std::shared_ptr record_batch, + arrow::RecordBatch::FromStructArray(batch, arrow_pool_.get())); + if (!writer) { + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( + writer, arrow::ipc::MakeStreamWriter(output, record_batch->schema(), + write_options)); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW(writer->WriteRecordBatch(*record_batch)); + return Status::OK(); + }; + // Optional serialization failures must not interrupt the source read. + cache_status = write_batch(); + return Status::OK(); + }, + /*prepare_reader=*/nullptr); + PAIMON_RETURN_NOT_OK(read_status->value()); + PAIMON_RETURN_NOT_OK(cache_status); + if (!writer) { + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( + writer, + arrow::ipc::MakeStreamWriter( + output, arrow::schema(serializer_->GetDataType()->fields()), write_options)); + } + PAIMON_RETURN_NOT_OK_FROM_ARROW(writer->Close()); + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr buffer, output->Finish()); + if (buffer->size() > std::numeric_limits::max()) { + return Status::Invalid("Manifest Arrow cache entry exceeds the memory segment limit"); } - return MemorySegment::Wrap(bytes); + // Retain the IPC buffer and its allocator without copying the complete stream into Bytes. + // The aliasing CacheValue keeps both alive, including after eviction with active readers. + struct OwnedBuffer { + OwnedBuffer(const std::shared_ptr& owner, + const std::shared_ptr& data) + : pool(owner), + buffer(data), + value(MemorySegment::WrapView(reinterpret_cast(data->data()), + static_cast(data->size())), + CacheCallback()) {} + std::shared_ptr pool; + std::shared_ptr buffer; + CacheValue value; + }; + auto owner = std::make_shared(arrow_pool_, buffer); + auto* value = &owner->value; + return std::shared_ptr(std::move(owner), value); +} + +template +Status ObjectsFile::ReadArrowBatches( + const std::string& file_name, std::optional file_size, + const std::function&)>& consumer, + const std::function*)>& prepare_reader) const { + if (!cache_) { + return ReadUncachedArrowBatches(file_name, file_size, consumer, prepare_reader); + } + auto metrics = std::make_shared(); + ScopeGuard record_metrics([&]() { read_metrics_->Merge(metrics); }); + const std::string path = path_factory_->ToPath(file_name); + auto key = CacheKey::ForKind(path, /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); + bool loaded = false; + std::optional read_status; + std::vector> batches; + auto cached = LoadCacheEntry(cache_, path, [&]() { + return cache_->Get(key, [&](const std::shared_ptr&) { + loaded = true; + return SerializeArrowBatches(file_name, file_size, &batches, &read_status); + }); + }); + const bool usable = cached.ok() && cached.value() && cached.value()->GetSegment().Data(); + metrics->SetCounter(!usable ? ScanMetrics::MANIFEST_ARROW_CACHE_FALLBACKS + : loaded ? ScanMetrics::MANIFEST_ARROW_CACHE_MISSES + : ScanMetrics::MANIFEST_ARROW_CACHE_HITS, + 1); + if (read_status) { + // Publish the query-independent cache before invoking user filters. A filter error + // neither invalidates the cache nor affects another reader waiting for this file. + PAIMON_RETURN_NOT_OK(read_status.value()); + for (const auto& batch : batches) { + PAIMON_RETURN_NOT_OK(consumer(batch)); + } + return Status::OK(); + } + if (!usable) { + return ReadUncachedArrowBatches(file_name, file_size, consumer, prepare_reader); + } + const auto& segment = cached.value()->GetSegment(); + auto buffer = std::make_shared(reinterpret_cast(segment.Data()), + segment.Size()); + auto input = std::make_shared(buffer); + auto read_options = arrow::ipc::IpcReadOptions::Defaults(); + read_options.memory_pool = arrow_pool_.get(); + read_options.use_threads = false; + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( + std::shared_ptr reader, + arrow::ipc::RecordBatchStreamReader::Open(input, read_options)); + while (true) { + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr batch, + reader->Next()); + if (!batch) { + break; + } + PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr array, + batch->ToStructArray()); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr aligned, + ManifestMetaReader::AlignArrayWithSchema( + array, serializer_->GetDataType(), arrow_pool_.get())); + PAIMON_RETURN_NOT_OK(consumer(checked_pointer_cast(aligned))); + } + return Status::OK(); } template From a5a919e761bb1db7f5986b9f05b6726e7aaeb289 Mon Sep 17 00:00:00 2001 From: gripleaf <425797155@qq.com> Date: Thu, 8 Oct 2026 21:28:29 +0800 Subject: [PATCH 2/5] refactor(manifest): simplify cache loading and metrics Delegate concurrent load coordination to Cache::Get instead of keeping a process-wide pending-load registry in ObjectsFile. Remove the dedicated cold-load coalescing test while retaining concurrent cache-read coverage. Keep only the decoded manifest cache hit and miss counters. Remove the row-range work counters, per-file histogram, and cache fallback counter, and trim the metrics documentation to the exported scan metrics. Preserve cache-error fallback and verify metric snapshot isolation. --- docs/source/user_guide/metrics.rst | 39 ++---- include/paimon/table/source/scan_metrics.h | 8 +- src/paimon/core/manifest/manifest_file.cpp | 23 +--- .../core/manifest/manifest_row_range_test.cpp | 117 ++---------------- src/paimon/core/utils/objects_file.h | 64 ++-------- 5 files changed, 31 insertions(+), 220 deletions(-) diff --git a/docs/source/user_guide/metrics.rst b/docs/source/user_guide/metrics.rst index 7a3997fb1..eed72f572 100644 --- a/docs/source/user_guide/metrics.rst +++ b/docs/source/user_guide/metrics.rst @@ -118,32 +118,13 @@ These metrics are C++-only and have no counterparts in Java Paimon. Manifest reads ~~~~~~~~~~~~~~ -Data-evolution scans with row ID ranges prune non-overlapping entries before -materializing file metadata. Unknown ranges are retained. - -When a caller cache is provided, ordinary, bucket and row-range manifest reads -share one Arrow IPC representation of each immutable manifest, replacing the raw -manifest byte cache. Manifest lists and index manifests use the same decoded IPC -cache through ``ObjectsFile``, with the existing whole-file keys and cache budget. -Cache keys do not contain query filters or mutable snapshot -selection. Concurrent loads of the same path through the same cache are coalesced -across readers; each reader applies its own filters. - -Cold reads decode the complete manifest to populate the shared cache and consume -the original batches without reading the IPC stream back. Warm reads skip the -source format decoder. Without a cache, or when the cache rejects manifest reads, -bucket reads retain source-level selective decoding when supported. Cold bucket -reads with a cache may therefore decode more entries than uncached selective -reads. Cache entries retain their buffers and allocator until the last reader -releases them, including after eviction. - -The scan metrics additionally expose cumulative ``rowRangeManifestEntriesScanned``, -``rowRangeManifestEntriesPruned`` and ``rowRangeManifestEntriesMaterialized``. -Materialized entries are counted before the ordinary entry filter. -``manifestArrowCacheHits``, ``manifestArrowCacheMisses`` and -``manifestArrowCacheFallbacks`` describe the decoded cache path for all three manifest read modes. -Manifest list and index manifest readers expose the same counters through their own read metrics. -Readers reusing a concurrent load count as cache hits. -``rowRangeManifestReadDuration`` is a per-file duration histogram in milliseconds, -including cache access, filtering and materialization, also on failed reads. -Parallel file durations overlap and must not be added to obtain request latency. +``manifestArrowCacheHits`` and ``manifestArrowCacheMisses`` are cumulative +counters for successful decoded-cache accesses by the scan's manifest-file +reader. Ordinary, bucket and row-range reads share the same query-independent +Arrow IPC entry under the existing whole-file cache key. Cache failures do not +update these counters. + +Warm reads skip the source format decoder. Cold reads populate the cache by +decoding the complete manifest, so a cold bucket read may decode more entries +than an uncached selective read. Concurrent load coordination is delegated to +the caller-provided cache. diff --git a/include/paimon/table/source/scan_metrics.h b/include/paimon/table/source/scan_metrics.h index 4296d5353..113ff1056 100644 --- a/include/paimon/table/source/scan_metrics.h +++ b/include/paimon/table/source/scan_metrics.h @@ -25,15 +25,9 @@ namespace paimon { /// Metric names for scan planning operations. class PAIMON_EXPORT ScanMetrics { public: - // Cumulative work for row-range manifest reads; duration histogram samples are per file (ms). - static constexpr char ROW_RANGE_MANIFEST_ENTRIES_SCANNED[] = "rowRangeManifestEntriesScanned"; - static constexpr char ROW_RANGE_MANIFEST_ENTRIES_PRUNED[] = "rowRangeManifestEntriesPruned"; - static constexpr char ROW_RANGE_MANIFEST_ENTRIES_MATERIALIZED[] = - "rowRangeManifestEntriesMaterialized"; - static constexpr char ROW_RANGE_MANIFEST_READ_DURATION[] = "rowRangeManifestReadDuration"; + // Cumulative decoded manifest cache hits and misses. static constexpr char MANIFEST_ARROW_CACHE_HITS[] = "manifestArrowCacheHits"; static constexpr char MANIFEST_ARROW_CACHE_MISSES[] = "manifestArrowCacheMisses"; - static constexpr char MANIFEST_ARROW_CACHE_FALLBACKS[] = "manifestArrowCacheFallbacks"; static constexpr char LAST_SCAN_DURATION[] = "lastScanDuration"; // Histogram metric for scan plan duration (milliseconds). diff --git a/src/paimon/core/manifest/manifest_file.cpp b/src/paimon/core/manifest/manifest_file.cpp index b39a2cd1b..0d1a03d34 100644 --- a/src/paimon/core/manifest/manifest_file.cpp +++ b/src/paimon/core/manifest/manifest_file.cpp @@ -30,13 +30,11 @@ #include "paimon/common/reader/late_materializing_file_batch_reader.h" #include "paimon/common/utils/arrow/status_utils.h" #include "paimon/common/utils/checked_cast.h" -#include "paimon/common/utils/scope_guard.h" #include "paimon/core/io/rolling_file_writer.h" #include "paimon/core/manifest/manifest_entry.h" #include "paimon/core/manifest/manifest_entry_serializer.h" #include "paimon/core/manifest/manifest_entry_writer_factory.h" #include "paimon/core/manifest/manifest_file_meta.h" -#include "paimon/core/utils/duration.h" #include "paimon/core/utils/file_store_path_factory.h" #include "paimon/core/utils/object_serializer.h" #include "paimon/core/utils/path_factory.h" @@ -46,7 +44,6 @@ #include "paimon/format/writer_builder.h" #include "paimon/predicate/predicate_builder.h" #include "paimon/status.h" -#include "paimon/table/source/scan_metrics.h" #include "paimon/utils/row_range_index.h" namespace arrow { @@ -142,24 +139,9 @@ Status ManifestFile::ReadRowRangeEntries( const std::string& file_name, const RowRangeIndex& row_ranges, const std::function(const ManifestEntry&)>& filter, std::optional file_size, std::vector* entries) const { - uint64_t scanned = 0; - uint64_t pruned = 0; - uint64_t materialized = 0; - Duration duration; - auto metrics = std::make_shared(); - // Merge once per file rather than locking a shared metric for every entry. ScopeGuard - // also accounts for work completed before an I/O or filter error. - ScopeGuard record_metrics([&]() { - metrics->SetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_SCANNED, scanned); - metrics->SetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_PRUNED, pruned); - metrics->SetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_MATERIALIZED, materialized); - metrics->ObserveHistogram(ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION, - static_cast(duration.Get())); - read_metrics_->Merge(metrics); - }); return ReadArrowBatches( file_name, file_size, - [this, &row_ranges, &filter, &scanned, &pruned, &materialized, + [this, &row_ranges, &filter, entries](const std::shared_ptr& batch) -> Status { // ManifestMetaReader has aligned both the entry and its nested file schema. Probe // the two range columns without allocating DataFileMeta, stats or binary keys. @@ -181,7 +163,6 @@ Status ManifestFile::ReadRowRangeEntries( auto first_ids = checked_pointer_cast(first_column); ColumnarRow row(batch->fields(), pool_, /*row_id=*/0); for (int64_t i = 0; i < batch->length(); ++i) { - ++scanned; row.SetRowId(i); PAIMON_RETURN_NOT_OK( ManifestEntrySerializer::ValidateVersion(row.GetInt(kVersionFieldIndex))); @@ -197,12 +178,10 @@ Status ManifestFile::ReadRowRangeEntries( if (first >= 0 && count > 0 && first <= std::numeric_limits::max() - (count - 1) && !row_ranges.Intersects(first, first + (count - 1))) { - ++pruned; continue; } } PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry, serializer_->FromRow(row)); - ++materialized; if (filter) { PAIMON_ASSIGN_OR_RAISE(bool keep, filter(entry)); if (!keep) { diff --git a/src/paimon/core/manifest/manifest_row_range_test.cpp b/src/paimon/core/manifest/manifest_row_range_test.cpp index 5a04e5efc..15b328b79 100644 --- a/src/paimon/core/manifest/manifest_row_range_test.cpp +++ b/src/paimon/core/manifest/manifest_row_range_test.cpp @@ -17,12 +17,9 @@ * under the License. */ -#include -#include #include #include #include -#include #include #include #include @@ -34,7 +31,6 @@ #include "gtest/gtest.h" #include "paimon/common/utils/checked_cast.h" #include "paimon/common/utils/path_util.h" -#include "paimon/common/utils/scope_guard.h" #include "paimon/core/io/data_file_meta.h" #include "paimon/core/io/meta_to_arrow_array_converter.h" #include "paimon/core/manifest/manifest_entry_serializer.h" @@ -160,6 +156,7 @@ TEST_F(RowRangeManifestFileTest, ArrowCacheReuseEvictionAndConcurrentReaders) { ASSERT_OK( unsupported->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, &entries)); ASSERT_EQ(expected, entries); + ASSERT_TRUE(unsupported->GetReadMetrics()->GetAllCounters().empty()); entries.clear(); auto small_cache = std::make_shared(1); ASSERT_OK_AND_ASSIGN(std::unique_ptr uncached, @@ -391,93 +388,10 @@ TEST_F(RowRangeManifestFileTest, CacheAdmissionFailureReusesAlreadyDecodedBatche ASSERT_OK(manifest->Read(written.first, filter, written.second, &entries)); ASSERT_EQ(1, calls); ASSERT_EQ(std::vector{entry}, entries); - ASSERT_OK_AND_ASSIGN(uint64_t fallbacks, manifest->GetReadMetrics()->GetCounter( - ScanMetrics::MANIFEST_ARROW_CACHE_FALLBACKS)); - ASSERT_EQ(1, fallbacks); + ASSERT_TRUE(manifest->GetReadMetrics()->GetAllCounters().empty()); } -class BlockingManifestCache : public CountingRoutingCache { - public: - BlockingManifestCache() : CountingRoutingCache(CacheKind::MANIFEST, 1024 * 1024) {} - - Result> Get( - const std::shared_ptr& key, - std::function>(const std::shared_ptr&)> - supplier) override { - return CountingRoutingCache::Get(key, [&](const std::shared_ptr& supplier_key) { - { - std::unique_lock lock(mutex_); - ++loads_; - cv_.notify_all(); - cv_.wait(lock, [&]() { return released_; }); - } - return supplier(supplier_key); - }); - } - - bool WaitForLoads(int32_t count) { - std::unique_lock lock(mutex_); - return cv_.wait_for(lock, std::chrono::seconds(10), [&]() { return loads_ >= count; }); - } - - void Release() { - std::lock_guard lock(mutex_); - released_ = true; - cv_.notify_all(); - } - - private: - std::mutex mutex_; - std::condition_variable cv_; - int32_t loads_ = 0; - bool released_ = false; -}; - -TEST_F(RowRangeManifestFileTest, ConcurrentColdReadersShareLoadAcrossInstances) { - auto dir = UniqueTestDirectory::Create(); - ASSERT_TRUE(dir); - auto cache = std::make_shared(); - ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, - CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); - ASSERT_OK_AND_ASSIGN(ManifestEntry entry, Entry("a.parquet", 100, 10)); - using WrittenFile = std::pair; - ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({entry})); - ASSERT_OK_AND_ASSIGN(WrittenFile other, manifest->WriteWithoutRolling({entry})); - std::vector> instances; - for (int32_t i = 0; i < 9; ++i) { - ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, - CreateManifest(dir->Str(), dir->GetFileSystem(), true, cache)); - instances.push_back(std::move(reader)); - } - std::vector> readers; - // Release blocked suppliers before future destructors wait, including on assertion failure. - ScopeGuard release([&]() { cache->Release(); }); - std::promise start; - auto gate = start.get_future().share(); - for (int32_t i = 0; i < 9; ++i) { - readers.push_back(std::async(std::launch::async, [&, i]() -> Status { - gate.wait(); - const auto& file = i == 8 ? other : written; - std::vector entries; - PAIMON_RETURN_NOT_OK(instances[i]->Read(file.first, nullptr, file.second, &entries)); - if (entries != std::vector{entry}) { - return Status::Invalid("concurrent manifest result differs"); - } - return Status::OK(); - })); - } - start.set_value(); - // An unrelated manifest must load while the first one is blocked. - ASSERT_TRUE(cache->WaitForLoads(2)); - cache->Release(); - for (auto& reader : readers) { - ASSERT_OK(reader.get()); - } - ASSERT_EQ(2, cache->SupplierCallCount()); - ASSERT_EQ(2, cache->Size()); -} - -TEST_F(RowRangeManifestFileTest, MetricsCountPruningCacheReuseAndFilterErrors) { +TEST_F(RowRangeManifestFileTest, CacheMetricsCountReadsAndPreserveSnapshots) { auto dir = UniqueTestDirectory::Create(); ASSERT_TRUE(dir); ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, @@ -494,37 +408,24 @@ TEST_F(RowRangeManifestFileTest, MetricsCountPruningCacheReuseAndFilterErrors) { ASSERT_EQ(std::vector{b}, entries); } auto metrics = manifest->GetReadMetrics(); - ASSERT_OK_AND_ASSIGN(uint64_t scanned, - metrics->GetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_SCANNED)); - ASSERT_OK_AND_ASSIGN(uint64_t pruned, - metrics->GetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_PRUNED)); - ASSERT_OK_AND_ASSIGN(uint64_t materialized, - metrics->GetCounter(ScanMetrics::ROW_RANGE_MANIFEST_ENTRIES_MATERIALIZED)); - ASSERT_EQ(4, scanned); - ASSERT_EQ(2, pruned); - ASSERT_EQ(2, materialized); ASSERT_OK_AND_ASSIGN(uint64_t hits, metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); ASSERT_OK_AND_ASSIGN(uint64_t misses, metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_MISSES)); ASSERT_EQ(1, hits); ASSERT_EQ(1, misses); - ASSERT_OK_AND_ASSIGN(HistogramStats stats, - metrics->GetHistogramStats(ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION)); - ASSERT_EQ(2, stats.count); std::vector entries; ASSERT_NOK(manifest->ReadRowRangeEntries( written.first, ranges, [](const ManifestEntry&) -> Result { return Status::IOError("filter failure"); }, written.second, &entries)); - ASSERT_OK_AND_ASSIGN(HistogramStats after_error, - manifest->GetReadMetrics()->GetHistogramStats( - ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION)); - ASSERT_EQ(3, after_error.count); + ASSERT_OK_AND_ASSIGN(uint64_t after_error, manifest->GetReadMetrics()->GetCounter( + ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); + ASSERT_EQ(2, after_error); // Previously returned snapshots remain unchanged after subsequent reads. - ASSERT_OK_AND_ASSIGN(HistogramStats old_snapshot, - metrics->GetHistogramStats(ScanMetrics::ROW_RANGE_MANIFEST_READ_DURATION)); - ASSERT_EQ(2, old_snapshot.count); + ASSERT_OK_AND_ASSIGN(uint64_t old_snapshot, + metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); + ASSERT_EQ(1, old_snapshot); } TEST_F(RowRangeManifestFileTest, BoundariesUnknownRangesAndDeleteMerging) { diff --git a/src/paimon/core/utils/objects_file.h b/src/paimon/core/utils/objects_file.h index 55d048622..a493174b9 100644 --- a/src/paimon/core/utils/objects_file.h +++ b/src/paimon/core/utils/objects_file.h @@ -19,11 +19,8 @@ #pragma once #include -#include #include -#include #include -#include #include #include #include @@ -89,7 +86,7 @@ class ObjectsFile { Result> WriteWithoutRolling(const std::vector& records); - /// Cumulative cache and derived reader work, including concurrent reads. + /// Cumulative decoded cache hits and misses, including concurrent reads. std::shared_ptr GetReadMetrics() const { auto snapshot = std::make_shared(); snapshot->Overwrite(read_metrics_); @@ -126,10 +123,6 @@ class ObjectsFile { std::string compression_; std::shared_ptr cache_; - static Result> LoadCacheEntry( - const std::shared_ptr& cache, const std::string& path, - const std::function>()>& load); - Result> SerializeArrowBatches( const std::string& file_name, std::optional file_size, std::vector>* batches, @@ -255,42 +248,6 @@ Result> ObjectsFile::OpenForRead( return file_system_->Open(file_path); } -// Only in-progress loads are retained here. Readers sharing a cache and path share one load, -// including across ObjectsFile instances; unrelated files load without holding the registry lock. -template -Result> ObjectsFile::LoadCacheEntry( - const std::shared_ptr& cache, const std::string& path, - const std::function>()>& load) { - using LoadKey = std::pair, std::string>; - using CacheResult = Result>; - static std::mutex mutex; - static std::map> pending; - const LoadKey key(cache, path); - std::shared_ptr> promise; - std::shared_future future; - { - std::lock_guard lock(mutex); - auto iter = pending.find(key); - if (iter != pending.end()) { - future = iter->second; - } else { - promise = std::make_shared>(); - future = promise->get_future().share(); - pending.emplace(key, future); - } - } - if (!promise) { - return future.get(); - } - ScopeGuard remove_pending([&]() { - std::lock_guard lock(mutex); - pending.erase(key); - }); - CacheResult result = load(); - promise->set_value(result); - return result; -} - template Result> ObjectsFile::SerializeArrowBatches( const std::string& file_name, std::optional file_size, @@ -377,20 +334,19 @@ Status ObjectsFile::ReadArrowBatches( bool loaded = false; std::optional read_status; std::vector> batches; - auto cached = LoadCacheEntry(cache_, path, [&]() { - return cache_->Get(key, [&](const std::shared_ptr&) { - loaded = true; - return SerializeArrowBatches(file_name, file_size, &batches, &read_status); - }); + auto cached = cache_->Get(key, [&](const std::shared_ptr&) { + loaded = true; + return SerializeArrowBatches(file_name, file_size, &batches, &read_status); }); const bool usable = cached.ok() && cached.value() && cached.value()->GetSegment().Data(); - metrics->SetCounter(!usable ? ScanMetrics::MANIFEST_ARROW_CACHE_FALLBACKS - : loaded ? ScanMetrics::MANIFEST_ARROW_CACHE_MISSES - : ScanMetrics::MANIFEST_ARROW_CACHE_HITS, - 1); + if (usable) { + metrics->SetCounter(loaded ? ScanMetrics::MANIFEST_ARROW_CACHE_MISSES + : ScanMetrics::MANIFEST_ARROW_CACHE_HITS, + 1); + } if (read_status) { // Publish the query-independent cache before invoking user filters. A filter error - // neither invalidates the cache nor affects another reader waiting for this file. + // does not invalidate the cache or affect other readers of this file. PAIMON_RETURN_NOT_OK(read_status.value()); for (const auto& batch : batches) { PAIMON_RETURN_NOT_OK(consumer(batch)); From df2b29c820005c9257ace25b7560cb60d0268a2f Mon Sep 17 00:00:00 2001 From: gripleaf <425797155@qq.com> Date: Fri, 9 Oct 2026 10:52:29 +0800 Subject: [PATCH 3/5] refactor(manifest): remove redundant metrics and schema alignment --- docs/source/user_guide/metrics.rst | 14 ---- include/paimon/table/source/scan_metrics.h | 4 - src/paimon/core/manifest/manifest_file.cpp | 1 - .../core/manifest/manifest_list_test.cpp | 1 - .../core/manifest/manifest_row_range_test.cpp | 79 +++++-------------- src/paimon/core/operation/file_store_scan.h | 1 - src/paimon/core/utils/manifest_meta_reader.h | 5 +- src/paimon/core/utils/objects_file.h | 25 +----- 8 files changed, 25 insertions(+), 105 deletions(-) diff --git a/docs/source/user_guide/metrics.rst b/docs/source/user_guide/metrics.rst index eed72f572..72ee2a176 100644 --- a/docs/source/user_guide/metrics.rst +++ b/docs/source/user_guide/metrics.rst @@ -114,17 +114,3 @@ These metrics are C++-only and have no counterparts in Java Paimon. "io.async.pending", "gauge", "requests", "Asynchronous callbacks not yet completed" "io.async.latency.count", "counter", "requests", "Completed asynchronous callback latency samples" "io.async.latency.sum-us", "counter", "microseconds", "Sum of asynchronous callback latency" - -Manifest reads -~~~~~~~~~~~~~~ - -``manifestArrowCacheHits`` and ``manifestArrowCacheMisses`` are cumulative -counters for successful decoded-cache accesses by the scan's manifest-file -reader. Ordinary, bucket and row-range reads share the same query-independent -Arrow IPC entry under the existing whole-file cache key. Cache failures do not -update these counters. - -Warm reads skip the source format decoder. Cold reads populate the cache by -decoding the complete manifest, so a cold bucket read may decode more entries -than an uncached selective read. Concurrent load coordination is delegated to -the caller-provided cache. diff --git a/include/paimon/table/source/scan_metrics.h b/include/paimon/table/source/scan_metrics.h index 113ff1056..15bc41ff1 100644 --- a/include/paimon/table/source/scan_metrics.h +++ b/include/paimon/table/source/scan_metrics.h @@ -25,10 +25,6 @@ namespace paimon { /// Metric names for scan planning operations. class PAIMON_EXPORT ScanMetrics { public: - // Cumulative decoded manifest cache hits and misses. - static constexpr char MANIFEST_ARROW_CACHE_HITS[] = "manifestArrowCacheHits"; - static constexpr char MANIFEST_ARROW_CACHE_MISSES[] = "manifestArrowCacheMisses"; - static constexpr char LAST_SCAN_DURATION[] = "lastScanDuration"; // Histogram metric for scan plan duration (milliseconds). static constexpr char SCAN_DURATION[] = "scanDuration"; diff --git a/src/paimon/core/manifest/manifest_file.cpp b/src/paimon/core/manifest/manifest_file.cpp index 0d1a03d34..7ca856ce2 100644 --- a/src/paimon/core/manifest/manifest_file.cpp +++ b/src/paimon/core/manifest/manifest_file.cpp @@ -58,7 +58,6 @@ namespace { constexpr int32_t kVersionFieldIndex = 0; constexpr int32_t kBucketFieldIndex = 3; constexpr int32_t kTotalBucketsFieldIndex = 4; - } // namespace ManifestFile::ManifestFile(const std::shared_ptr& file_system, diff --git a/src/paimon/core/manifest/manifest_list_test.cpp b/src/paimon/core/manifest/manifest_list_test.cpp index 2a83a5514..fa839b3f3 100644 --- a/src/paimon/core/manifest/manifest_list_test.cpp +++ b/src/paimon/core/manifest/manifest_list_test.cpp @@ -27,7 +27,6 @@ #include "arrow/type.h" #include "gtest/gtest.h" #include "paimon/common/utils/path_util.h" -#include "paimon/core/core_options.h" #include "paimon/core/manifest/manifest_file_meta.h" #include "paimon/core/snapshot.h" #include "paimon/core/stats/simple_stats.h" diff --git a/src/paimon/core/manifest/manifest_row_range_test.cpp b/src/paimon/core/manifest/manifest_row_range_test.cpp index 15b328b79..71f6468c0 100644 --- a/src/paimon/core/manifest/manifest_row_range_test.cpp +++ b/src/paimon/core/manifest/manifest_row_range_test.cpp @@ -42,7 +42,6 @@ #include "paimon/format/format_writer.h" #include "paimon/format/writer_builder.h" #include "paimon/fs/local/local_file_system.h" -#include "paimon/table/source/scan_metrics.h" #include "paimon/testing/utils/counting_cache_test_utils.h" #include "paimon/testing/utils/testharness.h" #include "paimon/utils/row_range_index.h" @@ -72,12 +71,12 @@ class RowRangeManifestFileTest : public ::testing::Test { } Result Entry(const std::string& name, std::optional first, - int64_t count, const FileKind& kind = FileKind::Add()) { + int64_t count) { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr meta, DataFileMeta::ForAppend(name, 100, count, SimpleStats::EmptyStats(), 0, 0, 0, FileSource::Append(), std::nullopt, std::nullopt, first, std::nullopt)); - return ManifestEntry(kind, BinaryRow::EmptyRow(), 0, 1, meta); + return ManifestEntry(FileKind::Add(), BinaryRow::EmptyRow(), 0, 1, meta); } Status WriteArray(const std::string& path, const std::string& name, @@ -128,7 +127,7 @@ TEST_F(RowRangeManifestFileTest, ArrowCacheReuseEvictionAndConcurrentReaders) { return Status::OK(); }; ASSERT_OK(read()); - // The manifest has one decoded entry; no raw-byte copy is retained. + // The first read loads the manifest once; concurrent warm reads reuse it. ASSERT_EQ(1, cache->SupplierCallCount(cache_kind)); // Readers own their decode state; the cache only shares immutable bytes. std::vector> readers; @@ -156,7 +155,6 @@ TEST_F(RowRangeManifestFileTest, ArrowCacheReuseEvictionAndConcurrentReaders) { ASSERT_OK( unsupported->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, &entries)); ASSERT_EQ(expected, entries); - ASSERT_TRUE(unsupported->GetReadMetrics()->GetAllCounters().empty()); entries.clear(); auto small_cache = std::make_shared(1); ASSERT_OK_AND_ASSIGN(std::unique_ptr uncached, @@ -207,11 +205,6 @@ TEST_F(RowRangeManifestFileTest, OrcCachePreservesManifestMetadata) { ASSERT_OK(reader->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual)); ASSERT_EQ(expected, actual); ASSERT_EQ(1, cache->SupplierCallCount(CacheKind::MANIFEST)); - ASSERT_OK_AND_ASSIGN(uint64_t count, - reader->GetReadMetrics()->GetCounter( - attempt == 0 ? ScanMetrics::MANIFEST_ARROW_CACHE_MISSES - : ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); - ASSERT_EQ(1, count); } } @@ -388,44 +381,6 @@ TEST_F(RowRangeManifestFileTest, CacheAdmissionFailureReusesAlreadyDecodedBatche ASSERT_OK(manifest->Read(written.first, filter, written.second, &entries)); ASSERT_EQ(1, calls); ASSERT_EQ(std::vector{entry}, entries); - ASSERT_TRUE(manifest->GetReadMetrics()->GetAllCounters().empty()); -} - -TEST_F(RowRangeManifestFileTest, CacheMetricsCountReadsAndPreserveSnapshots) { - auto dir = UniqueTestDirectory::Create(); - ASSERT_TRUE(dir); - ASSERT_OK_AND_ASSIGN(std::unique_ptr manifest, - CreateManifest(dir->Str(), dir->GetFileSystem(), true)); - ASSERT_OK_AND_ASSIGN(ManifestEntry a, Entry("a.parquet", 100, 10)); - ASSERT_OK_AND_ASSIGN(ManifestEntry b, Entry("b.parquet", 110, 10)); - using WrittenFile = std::pair; - ASSERT_OK_AND_ASSIGN(WrittenFile written, manifest->WriteWithoutRolling({a, b})); - ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, RowRangeIndex::Create({Range(110, 110)})); - for (int32_t i = 0; i < 2; ++i) { - std::vector entries; - ASSERT_OK(manifest->ReadRowRangeEntries(written.first, ranges, nullptr, written.second, - &entries)); - ASSERT_EQ(std::vector{b}, entries); - } - auto metrics = manifest->GetReadMetrics(); - ASSERT_OK_AND_ASSIGN(uint64_t hits, - metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); - ASSERT_OK_AND_ASSIGN(uint64_t misses, - metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_MISSES)); - ASSERT_EQ(1, hits); - ASSERT_EQ(1, misses); - std::vector entries; - ASSERT_NOK(manifest->ReadRowRangeEntries( - written.first, ranges, - [](const ManifestEntry&) -> Result { return Status::IOError("filter failure"); }, - written.second, &entries)); - ASSERT_OK_AND_ASSIGN(uint64_t after_error, manifest->GetReadMetrics()->GetCounter( - ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); - ASSERT_EQ(2, after_error); - // Previously returned snapshots remain unchanged after subsequent reads. - ASSERT_OK_AND_ASSIGN(uint64_t old_snapshot, - metrics->GetCounter(ScanMetrics::MANIFEST_ARROW_CACHE_HITS)); - ASSERT_EQ(1, old_snapshot); } TEST_F(RowRangeManifestFileTest, BoundariesUnknownRangesAndDeleteMerging) { @@ -556,16 +511,24 @@ TEST_F(RowRangeManifestFileTest, SchemaEvolutionAndVersionValidation) { std::shared_ptr evolved = evolved_result.ValueOrDie(); const std::string name = fmt::format("manifest-evolved-{}", mode); ASSERT_OK(WriteArray(dir->Str(), name, evolved)); - std::vector actual; - if (mode >= 2) { - ASSERT_NOK_WITH_MSG( - manifest->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual), - (mode == 2 ? "not compatible" : "Unsupported version: 999")); - } else { - ASSERT_OK(manifest->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual)); - ASSERT_EQ(mode == 0 ? 1 : 0, actual.size()); - if (mode == 0) { - ASSERT_FALSE(actual[0].File()->first_row_id.has_value()); + for (int32_t attempt = 0; attempt < 2; ++attempt) { + SCOPED_TRACE(attempt); + std::vector actual; + if (mode >= 2) { + ASSERT_NOK_WITH_MSG( + manifest->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual), + (mode == 2 ? "not compatible" : "Unsupported version: 999")); + } else { + ASSERT_OK( + manifest->ReadRowRangeEntries(name, ranges, nullptr, std::nullopt, &actual)); + ASSERT_EQ(mode == 0 ? 1 : 0, actual.size()); + if (mode == 0) { + ASSERT_FALSE(actual[0].File()->first_row_id.has_value()); + } + } + if (attempt == 0) { + // The second read must preserve schema evolution and validation from IPC alone. + manifest->DeleteQuietly(name); } } } diff --git a/src/paimon/core/operation/file_store_scan.h b/src/paimon/core/operation/file_store_scan.h index f8e85c32e..5fe128da1 100644 --- a/src/paimon/core/operation/file_store_scan.h +++ b/src/paimon/core/operation/file_store_scan.h @@ -174,7 +174,6 @@ class FileStoreScan { std::shared_ptr GetScanMetrics() const { auto snapshot = std::make_shared(); snapshot->Overwrite(metrics_); - snapshot->Merge(manifest_file_->GetReadMetrics()); return snapshot; } diff --git a/src/paimon/core/utils/manifest_meta_reader.h b/src/paimon/core/utils/manifest_meta_reader.h index 95eaaeea0..354132cf3 100644 --- a/src/paimon/core/utils/manifest_meta_reader.h +++ b/src/paimon/core/utils/manifest_meta_reader.h @@ -60,13 +60,12 @@ class PAIMON_EXPORT ManifestMetaReader : public BatchReader { reader_->Close(); } - /// Align fields by name, fill missing fields and correct integer precision. - /// Other physical types, including timestamp units, are preserved. + private: + // fill non exist field and correct int precision static Result> AlignArrayWithSchema( const std::shared_ptr& src_array, const std::shared_ptr& target_type, arrow::MemoryPool* pool); - private: std::unique_ptr reader_; std::shared_ptr target_type_; std::shared_ptr pool_; diff --git a/src/paimon/core/utils/objects_file.h b/src/paimon/core/utils/objects_file.h index a493174b9..32e2ecf3c 100644 --- a/src/paimon/core/utils/objects_file.h +++ b/src/paimon/core/utils/objects_file.h @@ -33,7 +33,6 @@ #include "fmt/format.h" #include "paimon/cache/cache.h" #include "paimon/common/data/columnar/columnar_row.h" -#include "paimon/common/metrics/metrics_impl.h" #include "paimon/common/utils/arrow/arrow_utils.h" #include "paimon/common/utils/arrow/mem_utils.h" #include "paimon/common/utils/arrow/status_utils.h" @@ -49,7 +48,6 @@ #include "paimon/format/writer_builder.h" #include "paimon/fs/file_system.h" #include "paimon/record_batch.h" -#include "paimon/table/source/scan_metrics.h" namespace paimon { /// A file which contains several `T`s, provides read and write. @@ -86,13 +84,6 @@ class ObjectsFile { Result> WriteWithoutRolling(const std::vector& records); - /// Cumulative decoded cache hits and misses, including concurrent reads. - std::shared_ptr GetReadMetrics() const { - auto snapshot = std::make_shared(); - snapshot->Overwrite(read_metrics_); - return snapshot; - } - protected: Status ValidateWrite() const { if (file_format_identifier_ != "avro") { @@ -108,7 +99,6 @@ class ObjectsFile { const std::function&)>& consumer, const std::function*)>& prepare_reader) const; - std::shared_ptr read_metrics_ = std::make_shared(); std::shared_ptr path_factory_; std::shared_ptr pool_; std::shared_ptr arrow_pool_; @@ -327,23 +317,14 @@ Status ObjectsFile::ReadArrowBatches( if (!cache_) { return ReadUncachedArrowBatches(file_name, file_size, consumer, prepare_reader); } - auto metrics = std::make_shared(); - ScopeGuard record_metrics([&]() { read_metrics_->Merge(metrics); }); const std::string path = path_factory_->ToPath(file_name); auto key = CacheKey::ForKind(path, /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); - bool loaded = false; std::optional read_status; std::vector> batches; auto cached = cache_->Get(key, [&](const std::shared_ptr&) { - loaded = true; return SerializeArrowBatches(file_name, file_size, &batches, &read_status); }); const bool usable = cached.ok() && cached.value() && cached.value()->GetSegment().Data(); - if (usable) { - metrics->SetCounter(loaded ? ScanMetrics::MANIFEST_ARROW_CACHE_MISSES - : ScanMetrics::MANIFEST_ARROW_CACHE_HITS, - 1); - } if (read_status) { // Publish the query-independent cache before invoking user filters. A filter error // does not invalidate the cache or affect other readers of this file. @@ -374,10 +355,8 @@ Status ObjectsFile::ReadArrowBatches( } PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr array, batch->ToStructArray()); - PAIMON_ASSIGN_OR_RAISE(std::shared_ptr aligned, - ManifestMetaReader::AlignArrayWithSchema( - array, serializer_->GetDataType(), arrow_pool_.get())); - PAIMON_RETURN_NOT_OK(consumer(checked_pointer_cast(aligned))); + // Cached batches were already aligned by ManifestMetaReader before serialization. + PAIMON_RETURN_NOT_OK(consumer(array)); } return Status::OK(); } From e681f768cee21d07da2501ca0bd75e6e30a5f04c Mon Sep 17 00:00:00 2001 From: gripleaf <425797155@qq.com> Date: Fri, 9 Oct 2026 16:02:11 +0800 Subject: [PATCH 4/5] test(manifest): cover scan pruning modes and update cache docs --- docs/source/user_guide/manifest_cache.rst | 37 +++--- .../core/manifest/manifest_row_range_test.cpp | 111 ++++++++++++++++++ 2 files changed, 132 insertions(+), 16 deletions(-) diff --git a/docs/source/user_guide/manifest_cache.rst b/docs/source/user_guide/manifest_cache.rst index e1c095a77..13943c654 100644 --- a/docs/source/user_guide/manifest_cache.rst +++ b/docs/source/user_guide/manifest_cache.rst @@ -21,18 +21,24 @@ Manifest Cache Overview -------- -Paimon C++ caches raw manifest file bytes at the ``ObjectsFile::Read()`` -layer. The cache uses the public ``Cache`` abstraction and is enabled through -``ScanContextBuilder::WithCache()``. The cache covers data manifests, manifest -lists, and index manifests because they all read through ``ObjectsFile``. +Paimon C++ caches decoded, schema-aligned manifest batches as Arrow IPC streams +in ``ObjectsFile``. The cache uses the public ``Cache`` abstraction and is +enabled through ``ScanContextBuilder::WithCache()``. It covers data manifests, +manifest lists, and index manifests because they all read through ``ObjectsFile``. For repeated ``get``, ``scan``, or batch ``get/scan -f`` requests in the same process, the same snapshot often reads the same manifest files repeatedly. On a -cache hit, the read path skips remote filesystem ``open/read``, builds an -in-memory input stream from cached bytes, and still runs the format reader, -Arrow decoding, and object deserialization. This design primarily reduces -remote IO latency and bandwidth while keeping cache weight aligned with the -actual cached bytes. +cache hit, the read path skips remote filesystem ``open/read`` and Avro/ORC +source-format decoding. It still reads the IPC stream, deserializes objects, +and applies scan filters. Bucket and row-range reads can select entries before +constructing their file metadata. Ordinary, bucket, and row-range reads share +the same complete, query-independent cache entry. + +A cold cache load decodes the complete manifest and serializes its aligned +batches into IPC. The loading reader consumes the original batches without an +IPC round trip. Cold bucket reads may therefore decode more entries than +uncached selective reads. Without a cache, reads retain the existing +source-format reader path. Configuration ------------- @@ -86,11 +92,10 @@ Example: Passing ``nullptr`` or omitting ``ScanContextBuilder::WithCache()`` leaves manifest caching disabled. -Future Optimizations --------------------- +Cache Implementation Responsibilities +------------------------------------- -- Add hit, miss, bypass, and eviction metrics to read trace or metrics. -- Add single-flight loading for high-concurrency misses on the same manifest - path. -- Evaluate a decoded-records second-level cache, configurable as a - CPU-vs-memory tradeoff. +Embedding applications can implement hit/miss and eviction statistics in their +``Cache`` implementation. Coordination of concurrent loads for the same key +also belongs to that implementation; ``ObjectsFile`` does not deduplicate +concurrent cache misses. diff --git a/src/paimon/core/manifest/manifest_row_range_test.cpp b/src/paimon/core/manifest/manifest_row_range_test.cpp index 71f6468c0..62354bd8f 100644 --- a/src/paimon/core/manifest/manifest_row_range_test.cpp +++ b/src/paimon/core/manifest/manifest_row_range_test.cpp @@ -17,8 +17,10 @@ * under the License. */ +#include #include #include +#include #include #include #include @@ -35,6 +37,8 @@ #include "paimon/core/io/meta_to_arrow_array_converter.h" #include "paimon/core/manifest/manifest_entry_serializer.h" #include "paimon/core/manifest/manifest_file.h" +#include "paimon/core/manifest/manifest_list.h" +#include "paimon/core/operation/data_evolution_file_store_scan.h" #include "paimon/core/operation/file_store_scan.h" #include "paimon/core/utils/file_store_path_factory.h" #include "paimon/format/file_format.h" @@ -169,6 +173,113 @@ TEST_F(RowRangeManifestFileTest, ArrowCacheReuseEvictionAndConcurrentReaders) { ASSERT_EQ(expected, entries); } +TEST_F(RowRangeManifestFileTest, ScanPlanPreservesResultsAcrossLazyDecodeAndCacheModes) { + auto pool = GetDefaultPool(); + std::shared_ptr executor = CreateDefaultExecutor(); + auto filters = std::make_shared( + nullptr, std::vector>{}, std::nullopt); + auto schema = arrow::schema({arrow::field("value", arrow::int32())}); + ASSERT_OK_AND_ASSIGN(std::shared_ptr table_schema, + TableSchema::Create(0, schema, {}, {}, {})); + ASSERT_OK_AND_ASSIGN(ManifestEntry a, Entry("a.parquet", 100, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry b, Entry("b.parquet", 110, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry c, Entry("c.parquet", 120, 10)); + ASSERT_OK_AND_ASSIGN(ManifestEntry unknown, Entry("unknown.parquet", std::nullopt, 10)); + ManifestEntry deleted(FileKind::Delete(), a.Partition(), a.Bucket(), a.TotalBuckets(), + a.File()); + struct Query { + std::optional> ranges; + std::vector expected; + int32_t retained_entries; + }; + const std::vector queries = { + {std::vector{Range(100, 100)}, {unknown}, 3}, + {std::vector{Range(110, 110)}, {b, unknown}, 2}, + {std::vector{Range(120, 129)}, {unknown, c}, 2}, + {std::vector{Range(109, 110), Range(129, 129)}, {b, unknown, c}, 5}, + {std::vector{Range(1000, 1000)}, {unknown}, 1}, + {std::vector{}, {unknown}, 1}, + {std::nullopt, {b, unknown, c}, 5}}; + for (bool cache_enabled : {false, true}) { + SCOPED_TRACE(cache_enabled); + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = dir->GetFileSystem(); + auto cache = cache_enabled + ? std::make_shared(CacheKind::MANIFEST, 1024 * 1024) + : nullptr; + ASSERT_OK_AND_ASSIGN(std::shared_ptr manifest, + CreateManifest(dir->Str(), fs, cache_enabled, cache)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile base, manifest->WriteWithoutRolling({a, b, unknown})); + ASSERT_OK_AND_ASSIGN(WrittenFile delta, manifest->WriteWithoutRolling({deleted, c})); + // Unknown manifest bounds force entry-level pruning even for disjoint queries. + ManifestFileMeta base_meta(base.first, base.second, 3, 0, SimpleStats::EmptyStats(), 0, 0, + 0, 0, 0, std::nullopt, std::nullopt); + ManifestFileMeta delta_meta(delta.first, delta.second, 1, 1, SimpleStats::EmptyStats(), 0, + 0, 0, 0, 0, std::nullopt, std::nullopt); + ASSERT_OK_AND_ASSIGN(std::shared_ptr format, + FileFormatFactory::Get("avro", {})); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr paths, + FileStorePathFactory::Create(dir->Str(), schema, {}, "", "avro", "data-", true, {}, + std::nullopt, false, pool)); + ASSERT_OK_AND_ASSIGN(std::shared_ptr lists, + ManifestList::Create(fs, format, "null", paths, cache, pool)); + ASSERT_OK_AND_ASSIGN(WrittenFile base_list, lists->Write({base_meta})); + ASSERT_OK_AND_ASSIGN(WrittenFile delta_list, lists->Write({delta_meta})); + Snapshot snapshot(1, 0, base_list.first, base_list.second, delta_list.first, + delta_list.second, std::nullopt, std::nullopt, std::nullopt, "test", 0, + Snapshot::CommitKind::Append(), 0, 30, 0, std::nullopt, std::nullopt, + std::nullopt, std::nullopt, std::nullopt); + for (bool lazy_decode : {false, true}) { + SCOPED_TRACE(lazy_decode); + const int64_t previous_loads = cache ? cache->SupplierCallCount() : 0; + if (cache) { + cache->InvalidateAll(); + } + ASSERT_OK_AND_ASSIGN( + CoreOptions options, + CoreOptions::FromMap({{Options::DATA_EVOLUTION_ENABLED, "true"}, + {Options::SCAN_MANIFEST_ENTRY_LAZY_DECODE_ENABLED, + lazy_decode ? "true" : "false"}})); + options.WithCache(cache); + for (int32_t attempt = 0; attempt < 2; ++attempt) { + SCOPED_TRACE(attempt); + for (size_t i = 0; i < queries.size(); ++i) { + SCOPED_TRACE(i); + const auto& query = queries[i]; + ASSERT_OK_AND_ASSIGN(std::unique_ptr scan, + DataEvolutionFileStoreScan::Create( + nullptr, nullptr, lists, manifest, table_schema, + schema, filters, options, executor, pool)); + scan->WithSnapshot(snapshot); + if (query.ranges) { + ASSERT_OK_AND_ASSIGN(RowRangeIndex ranges, + RowRangeIndex::Create(query.ranges.value())); + scan->WithRowRangeIndex(ranges); + } + std::atomic filter_calls{0}; + scan->WithLevelFilter([&](int32_t) { + ++filter_calls; + return true; + }); + ASSERT_OK_AND_ASSIGN(std::shared_ptr plan, + scan->CreatePlan()); + ASSERT_EQ(query.expected, plan->Files()); + // The regular filter sees only retained entries when early pruning is on. + // This also verifies that both the Add and Delete of a reach the merge. + ASSERT_EQ(lazy_decode ? query.retained_entries : 5, filter_calls.load()); + if (cache) { + // Two manifest lists and two manifests load once per cache reset. + ASSERT_EQ(previous_loads + 4, cache->SupplierCallCount()); + } + } + } + } + } +} + TEST_F(RowRangeManifestFileTest, OrcCachePreservesManifestMetadata) { auto pool = GetDefaultPool(); auto fs = std::make_shared(); From 0a3bb7d52be9d06d173383a80730697968dcb20e Mon Sep 17 00:00:00 2001 From: gripleaf <425797155@qq.com> Date: Sat, 10 Oct 2026 15:40:17 +0800 Subject: [PATCH 5/5] fix(manifest): charge retained capacity and centralize field positions --- docs/source/user_guide/manifest_cache.rst | 7 ++ include/paimon/cache/cache.h | 9 +++ src/paimon/common/io/cache/cache.cpp | 12 ++- src/paimon/common/io/cache/lru_cache.cpp | 2 +- src/paimon/common/io/cache/lru_cache.h | 2 +- src/paimon/common/io/cache/lru_cache_test.cpp | 34 +++++++++ src/paimon/core/io/data_file_meta.h | 4 + .../core/io/data_file_meta_serializer.cpp | 12 +-- src/paimon/core/manifest/manifest_entry.h | 3 + .../manifest/manifest_entry_serializer.cpp | 4 +- src/paimon/core/manifest/manifest_file.cpp | 10 +-- .../core/manifest/manifest_row_range_test.cpp | 75 ++++++++++++++++--- src/paimon/core/utils/objects_file.h | 2 +- 13 files changed, 149 insertions(+), 27 deletions(-) diff --git a/docs/source/user_guide/manifest_cache.rst b/docs/source/user_guide/manifest_cache.rst index 13943c654..f6dd034af 100644 --- a/docs/source/user_guide/manifest_cache.rst +++ b/docs/source/user_guide/manifest_cache.rst @@ -99,3 +99,10 @@ Embedding applications can implement hit/miss and eviction statistics in their ``Cache`` implementation. Coordination of concurrent loads for the same key also belongs to that implementation; ``ObjectsFile`` does not deduplicate concurrent cache misses. + +Cache implementations should use ``CacheValue::GetMemoryUsage()`` for admission +and eviction accounting. For manifest IPC, this reports the retained Arrow +buffer capacity, including unused space from growth. ``GetSegment().Size()`` +continues to describe the valid IPC byte length. The built-in ``LruCache`` uses +the memory usage value; existing cache values without an explicit allocation +size continue to charge their segment length. diff --git a/include/paimon/cache/cache.h b/include/paimon/cache/cache.h index b282da74b..e05b555e5 100644 --- a/include/paimon/cache/cache.h +++ b/include/paimon/cache/cache.h @@ -88,12 +88,20 @@ class PAIMON_EXPORT Cache { class PAIMON_EXPORT CacheValue { public: + /// Charge the segment's logical size when no separate allocation size is supplied. CacheValue(const MemorySegment& segment, CacheCallback callback); + /// @param memory_usage Retained allocation size in bytes, including unused capacity. + /// The charge is at least the segment's logical size. + CacheValue(const MemorySegment& segment, CacheCallback callback, int64_t memory_usage); + ~CacheValue(); const MemorySegment& GetSegment() const; + /// Bytes to charge against the cache budget. May exceed GetSegment().Size(). + int64_t GetMemoryUsage() const; + void OnEvict(const std::shared_ptr& key) const; bool operator==(const CacheValue& other) const; @@ -101,6 +109,7 @@ class PAIMON_EXPORT CacheValue { private: MemorySegment segment_; CacheCallback callback_; + int64_t memory_usage_; }; } // namespace paimon diff --git a/src/paimon/common/io/cache/cache.cpp b/src/paimon/common/io/cache/cache.cpp index 37a26e917..c4c9e5418 100644 --- a/src/paimon/common/io/cache/cache.cpp +++ b/src/paimon/common/io/cache/cache.cpp @@ -19,12 +19,18 @@ #include "paimon/cache/cache.h" +#include #include namespace paimon { CacheValue::CacheValue(const MemorySegment& segment, CacheCallback callback) - : segment_(segment), callback_(std::move(callback)) {} + : CacheValue(segment, std::move(callback), segment.Size()) {} + +CacheValue::CacheValue(const MemorySegment& segment, CacheCallback callback, int64_t memory_usage) + : segment_(segment), + callback_(std::move(callback)), + memory_usage_(std::max(segment.Size(), memory_usage)) {} CacheValue::~CacheValue() = default; @@ -32,6 +38,10 @@ const MemorySegment& CacheValue::GetSegment() const { return segment_; } +int64_t CacheValue::GetMemoryUsage() const { + return memory_usage_; +} + void CacheValue::OnEvict(const std::shared_ptr& key) const { if (callback_) { callback_(key); diff --git a/src/paimon/common/io/cache/lru_cache.cpp b/src/paimon/common/io/cache/lru_cache.cpp index fbfca83ca..892ddb6f1 100644 --- a/src/paimon/common/io/cache/lru_cache.cpp +++ b/src/paimon/common/io/cache/lru_cache.cpp @@ -26,7 +26,7 @@ LruCache::LruCache(int64_t max_weight) .expire_after_access_ms = -1, .weigh_func = [](const std::shared_ptr& /*key*/, const std::shared_ptr& value) -> int64_t { - return value ? value->GetSegment().Size() : 0; + return value ? value->GetMemoryUsage() : 0; }, .removal_callback = [](const std::shared_ptr& key, const std::shared_ptr& value, diff --git a/src/paimon/common/io/cache/lru_cache.h b/src/paimon/common/io/cache/lru_cache.h index 745ca7f5d..3bf2b10a4 100644 --- a/src/paimon/common/io/cache/lru_cache.h +++ b/src/paimon/common/io/cache/lru_cache.h @@ -32,7 +32,7 @@ namespace paimon { /// LRU Cache implementation with weight-based eviction for block cache. /// /// Wraps GenericLruCache with CacheKey/CacheValue types. Capacity is measured -/// in bytes (sum of MemorySegment sizes). When an entry is evicted, its +/// in bytes (sum of CacheValue::GetMemoryUsage()). When an entry is evicted, its /// CacheCallback is invoked to notify the upper layer. /// /// @note Thread-safe: all public methods are protected by the underlying GenericLruCache lock. diff --git a/src/paimon/common/io/cache/lru_cache_test.cpp b/src/paimon/common/io/cache/lru_cache_test.cpp index 29eb10fef..6a9cab5b2 100644 --- a/src/paimon/common/io/cache/lru_cache_test.cpp +++ b/src/paimon/common/io/cache/lru_cache_test.cpp @@ -113,6 +113,40 @@ TEST_F(LruCacheTest, TestPutInsertAndUpdate) { ASSERT_EQ(result->GetSegment().Get(0), 'B'); } +TEST_F(LruCacheTest, TestRetainedMemoryAccounting) { + LruCache cache(192); + auto segment = MemorySegment::AllocateHeapMemory(64, pool_.get()); + auto value = std::make_shared(segment, CacheCallback(), 128); + ASSERT_EQ(64, value->GetSegment().Size()); + ASSERT_EQ(128, value->GetMemoryUsage()); + ASSERT_EQ(64, CacheValue(segment, CacheCallback()).GetMemoryUsage()); + ASSERT_EQ(64, CacheValue(segment, CacheCallback(), 32).GetMemoryUsage()); + + std::vector> evicted; + auto first = std::make_shared( + segment, [&](const std::shared_ptr& key) { evicted.push_back(key); }, 128); + auto first_key = MakeKey(0); + ASSERT_OK(cache.Put(first_key, first)); + ASSERT_EQ(128, cache.GetCurrentWeight()); + ASSERT_OK(cache.Put(MakeKey(1), value)); + ASSERT_EQ(1, cache.Size()); + ASSERT_EQ(128, cache.GetCurrentWeight()); + ASSERT_EQ(std::vector>{first_key}, evicted); + cache.InvalidateAll(); + ASSERT_EQ(0, cache.GetCurrentWeight()); + + // The payload fits, but its retained allocation exceeds the cache budget. + LruCache small_cache(100); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr loaded, + small_cache.Get(MakeKey(0), + [&](const std::shared_ptr&) + -> Result> { return value; })); + ASSERT_EQ(value, loaded); + ASSERT_EQ(0, small_cache.Size()); + ASSERT_EQ(0, small_cache.GetCurrentWeight()); +} + /// Verifies weight-based eviction: when total weight exceeds max, LRU entries are evicted. TEST_F(LruCacheTest, TestWeightBasedEviction) { // Cache can hold at most 200 bytes diff --git a/src/paimon/core/io/data_file_meta.h b/src/paimon/core/io/data_file_meta.h index 436d73785..ee3d4a9f5 100644 --- a/src/paimon/core/io/data_file_meta.h +++ b/src/paimon/core/io/data_file_meta.h @@ -43,6 +43,10 @@ class Bytes; /// Metadata of a data file. struct DataFileMeta { + /// Field positions in DataType(), shared by serialization and manifest pruning. + static constexpr int32_t kRowCountFieldIndex = 2; + static constexpr int32_t kFirstRowIdFieldIndex = 18; + static const BinaryRow& EmptyMinKey(); static const BinaryRow& EmptyMaxKey(); static constexpr int32_t DUMMY_LEVEL = 0; diff --git a/src/paimon/core/io/data_file_meta_serializer.cpp b/src/paimon/core/io/data_file_meta_serializer.cpp index 9f88ebec0..8fa7eba0d 100644 --- a/src/paimon/core/io/data_file_meta_serializer.cpp +++ b/src/paimon/core/io/data_file_meta_serializer.cpp @@ -46,7 +46,7 @@ Result DataFileMetaSerializer::ToRow(const std::shared_ptrfile_name, pool_.get())); writer.WriteLong(1, meta->file_size); - writer.WriteLong(2, meta->row_count); + writer.WriteLong(DataFileMeta::kRowCountFieldIndex, meta->row_count); auto min_key_bytes = SerializationUtils::SerializeBinaryRow(meta->min_key, pool_.get()); writer.WriteBinary(3, *min_key_bytes); auto max_key_bytes = SerializationUtils::SerializeBinaryRow(meta->max_key, pool_.get()); @@ -86,9 +86,9 @@ Result DataFileMetaSerializer::ToRow(const std::shared_ptrexternal_path.value(), pool_.get())); } if (meta->first_row_id == std::nullopt) { - writer.SetNullAt(18); + writer.SetNullAt(DataFileMeta::kFirstRowIdFieldIndex); } else { - writer.WriteLong(18, meta->first_row_id.value()); + writer.WriteLong(DataFileMeta::kFirstRowIdFieldIndex, meta->first_row_id.value()); } if (meta->write_cols == std::nullopt) { writer.SetNullAt(19); @@ -110,7 +110,7 @@ Result> DataFileMetaSerializer::FromRow( const InternalRow& row) const { auto file_name = row.GetString(0); auto file_size = row.GetLong(1); - auto row_count = row.GetLong(2); + auto row_count = row.GetLong(DataFileMeta::kRowCountFieldIndex); auto min_key = row.GetBinary(3); auto max_key = row.GetBinary(4); auto key_stats_row = row.GetRow(5, 3); @@ -155,8 +155,8 @@ Result> DataFileMetaSerializer::FromRow( external_path = row.GetString(17).ToString(); } std::optional first_row_id; - if (!row.IsNullAt(18)) { - first_row_id = row.GetLong(18); + if (!row.IsNullAt(DataFileMeta::kFirstRowIdFieldIndex)) { + first_row_id = row.GetLong(DataFileMeta::kFirstRowIdFieldIndex); } std::optional> write_cols; diff --git a/src/paimon/core/manifest/manifest_entry.h b/src/paimon/core/manifest/manifest_entry.h index ca6e8ffad..f75461134 100644 --- a/src/paimon/core/manifest/manifest_entry.h +++ b/src/paimon/core/manifest/manifest_entry.h @@ -38,6 +38,9 @@ namespace paimon { /// Entry of a manifest file, representing an addition / deletion of a data file. class ManifestEntry : public FileEntry { public: + /// Field position in DataType(), before the serializer prepends _VERSION. + static constexpr int32_t kFileFieldIndex = 4; + static const std::shared_ptr& DataType(); static int64_t RecordCount(const std::vector& manifest_entries); static std::optional NullableRecordCount( diff --git a/src/paimon/core/manifest/manifest_entry_serializer.cpp b/src/paimon/core/manifest/manifest_entry_serializer.cpp index 2389cd8d7..5e1a2be08 100644 --- a/src/paimon/core/manifest/manifest_entry_serializer.cpp +++ b/src/paimon/core/manifest/manifest_entry_serializer.cpp @@ -58,7 +58,7 @@ Result ManifestEntrySerializer::ConvertFrom(int32_t version, SerializationUtils::DeserializeBinaryRow(partition_bytes)); auto bucket = row.GetInt(2); auto total_buckets = row.GetInt(3); - auto file = row.GetRow(4, data_file_meta_serializer_.NumFields()); + auto file = row.GetRow(ManifestEntry::kFileFieldIndex, data_file_meta_serializer_.NumFields()); if (!file) { return Status::Invalid("ManifestEntry convert from row failed, with null DataFileMeta"); } @@ -80,7 +80,7 @@ Result ManifestEntrySerializer::ToRow(const ManifestEntry& record) co writer.WriteInt(4, record.TotalBuckets()); PAIMON_ASSIGN_OR_RAISE(BinaryRow data_file_meta_row, data_file_meta_serializer_.ToRow(record.File())); - writer.WriteRow(5, data_file_meta_row); + writer.WriteRow(ManifestEntry::kFileFieldIndex + 1, data_file_meta_row); writer.Complete(); return row; } diff --git a/src/paimon/core/manifest/manifest_file.cpp b/src/paimon/core/manifest/manifest_file.cpp index 7ca856ce2..c0a71a2f8 100644 --- a/src/paimon/core/manifest/manifest_file.cpp +++ b/src/paimon/core/manifest/manifest_file.cpp @@ -144,16 +144,14 @@ Status ManifestFile::ReadRowRangeEntries( entries](const std::shared_ptr& batch) -> Status { // ManifestMetaReader has aligned both the entry and its nested file schema. Probe // the two range columns without allocating DataFileMeta, stats or binary keys. - constexpr int32_t kFileFieldIndex = 5; - constexpr int32_t kRowCountFieldIndex = 2; - constexpr int32_t kFirstRowIdFieldIndex = 18; - const auto& file_column = batch->field(kFileFieldIndex); + // The serialized entry has a leading _VERSION field. + const auto& file_column = batch->field(ManifestEntry::kFileFieldIndex + 1); if (file_column->type_id() != arrow::Type::STRUCT) { return Status::Invalid("Manifest entry file metadata must be a struct"); } auto files = checked_pointer_cast(file_column); - const auto& count_column = files->field(kRowCountFieldIndex); - const auto& first_column = files->field(kFirstRowIdFieldIndex); + const auto& count_column = files->field(DataFileMeta::kRowCountFieldIndex); + const auto& first_column = files->field(DataFileMeta::kFirstRowIdFieldIndex); if (count_column->type_id() != arrow::Type::INT64 || first_column->type_id() != arrow::Type::INT64) { return Status::Invalid("Manifest entry row range must contain int64 fields"); diff --git a/src/paimon/core/manifest/manifest_row_range_test.cpp b/src/paimon/core/manifest/manifest_row_range_test.cpp index 62354bd8f..a4f60fea1 100644 --- a/src/paimon/core/manifest/manifest_row_range_test.cpp +++ b/src/paimon/core/manifest/manifest_row_range_test.cpp @@ -457,6 +457,59 @@ TEST_F(RowRangeManifestFileTest, CachedBufferRetainsAllocatorUntilEviction) { ASSERT_EQ(expected, entries); } +TEST_F(RowRangeManifestFileTest, CacheChargesRetainedArrowCapacity) { + auto dir = UniqueTestDirectory::Create(); + ASSERT_TRUE(dir); + auto fs = dir->GetFileSystem(); + ASSERT_OK_AND_ASSIGN(std::unique_ptr writer, + CreateManifest(dir->Str(), fs, false)); + // Exercise growth beyond the initial IPC buffer allocation. + ASSERT_OK_AND_ASSIGN(ManifestEntry entry, Entry(std::string(5000, 'a') + ".parquet", 100, 10)); + using WrittenFile = std::pair; + ASSERT_OK_AND_ASSIGN(WrittenFile written, writer->WriteWithoutRolling({entry})); + std::shared_ptr pool = GetMemoryPool(); + auto cache = std::make_shared(1024 * 1024); + { + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateManifest(dir->Str(), fs, true, cache, pool)); + std::vector entries; + ASSERT_OK(reader->Read(written.first, nullptr, written.second, &entries)); + ASSERT_EQ(std::vector{entry}, entries); + } + // Only the cached IPC allocation remains in this reader's pool. + ASSERT_EQ(1, cache->Size()); + ASSERT_EQ(pool->CurrentUsage(), cache->GetCurrentWeight()); + auto key = CacheKey::ForKind( + PathUtil::JoinPath(FileStorePathFactory::ManifestPath(dir->Str()), written.first), + /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr cached, + cache->Get(key, + [](const std::shared_ptr&) -> Result> { + return Status::Invalid("manifest was not cached"); + })); + ASSERT_TRUE(cached); + const int64_t logical_size = cached->GetSegment().Size(); + ASSERT_GT(cached->GetMemoryUsage(), logical_size); + ASSERT_EQ(pool->CurrentUsage(), cached->GetMemoryUsage()); + cached.reset(); + cache->InvalidateAll(); + ASSERT_EQ(0, pool->CurrentUsage()); + + // Admission must reject the allocation even though the valid IPC bytes would fit. + auto small_cache = std::make_shared(logical_size); + { + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + CreateManifest(dir->Str(), fs, true, small_cache, pool)); + std::vector entries; + ASSERT_OK(reader->Read(written.first, nullptr, written.second, &entries)); + ASSERT_EQ(std::vector{entry}, entries); + ASSERT_EQ(0, small_cache->Size()); + ASSERT_EQ(0, small_cache->GetCurrentWeight()); + } + ASSERT_EQ(0, pool->CurrentUsage()); +} + TEST_F(RowRangeManifestFileTest, CacheAdmissionFailureReusesAlreadyDecodedBatches) { class FailingCache : public CountingRoutingCache { public: @@ -580,6 +633,7 @@ TEST_F(RowRangeManifestFileTest, UncertainAndOverflowingRangesAreRetained) { } TEST_F(RowRangeManifestFileTest, SchemaEvolutionAndVersionValidation) { + constexpr int32_t kVersionedFileFieldIndex = ManifestEntry::kFileFieldIndex + 1; auto dir = UniqueTestDirectory::Create(); ASSERT_TRUE(dir); auto pool = GetDefaultPool(); @@ -597,15 +651,17 @@ TEST_F(RowRangeManifestFileTest, SchemaEvolutionAndVersionValidation) { SCOPED_TRACE(mode); auto fields = serializer.GetDataType()->fields(); auto columns = batch->fields(); - auto file = checked_pointer_cast(columns[5]); + auto file = checked_pointer_cast(columns[kVersionedFileFieldIndex]); auto file_fields = checked_pointer_cast(file->type())->fields(); auto file_columns = file->fields(); if (mode == 0) { - file_fields.erase(file_fields.begin() + 18); - file_columns.erase(file_columns.begin() + 18); + file_fields.erase(file_fields.begin() + DataFileMeta::kFirstRowIdFieldIndex); + file_columns.erase(file_columns.begin() + DataFileMeta::kFirstRowIdFieldIndex); } else if (mode == 1) { - std::swap(file_fields[2], file_fields[18]); - std::swap(file_columns[2], file_columns[18]); + std::swap(file_fields[DataFileMeta::kRowCountFieldIndex], + file_fields[DataFileMeta::kFirstRowIdFieldIndex]); + std::swap(file_columns[DataFileMeta::kRowCountFieldIndex], + file_columns[DataFileMeta::kFirstRowIdFieldIndex]); } else { arrow::Int32Builder versions; ASSERT_TRUE(versions.Append(mode == 2 ? 1 : 999).ok()); @@ -613,10 +669,11 @@ TEST_F(RowRangeManifestFileTest, SchemaEvolutionAndVersionValidation) { } auto updated_file = arrow::StructArray::Make(file_columns, file_fields); ASSERT_TRUE(updated_file.ok()) << updated_file.status().ToString(); - columns[5] = updated_file.ValueOrDie(); - fields[5] = fields[5]->WithType(arrow::struct_(file_fields)); - std::swap(fields[1], fields[5]); - std::swap(columns[1], columns[5]); + columns[kVersionedFileFieldIndex] = updated_file.ValueOrDie(); + fields[kVersionedFileFieldIndex] = + fields[kVersionedFileFieldIndex]->WithType(arrow::struct_(file_fields)); + std::swap(fields[1], fields[kVersionedFileFieldIndex]); + std::swap(columns[1], columns[kVersionedFileFieldIndex]); auto evolved_result = arrow::StructArray::Make(columns, fields); ASSERT_TRUE(evolved_result.ok()) << evolved_result.status().ToString(); std::shared_ptr evolved = evolved_result.ValueOrDie(); diff --git a/src/paimon/core/utils/objects_file.h b/src/paimon/core/utils/objects_file.h index 32e2ecf3c..2a99e98d1 100644 --- a/src/paimon/core/utils/objects_file.h +++ b/src/paimon/core/utils/objects_file.h @@ -299,7 +299,7 @@ Result> ObjectsFile::SerializeArrowBatches( buffer(data), value(MemorySegment::WrapView(reinterpret_cast(data->data()), static_cast(data->size())), - CacheCallback()) {} + CacheCallback(), data->capacity()) {} std::shared_ptr pool; std::shared_ptr buffer; CacheValue value;