diff --git a/docs/source/user_guide/manifest_cache.rst b/docs/source/user_guide/manifest_cache.rst index e1c095a77..f6dd034af 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,17 @@ 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. + +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/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/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/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_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 d525f62a6..c0a71a2f8 100644 --- a/src/paimon/core/manifest/manifest_file.cpp +++ b/src/paimon/core/manifest/manifest_file.cpp @@ -44,6 +44,7 @@ #include "paimon/format/writer_builder.h" #include "paimon/predicate/predicate_builder.h" #include "paimon/status.h" +#include "paimon/utils/row_range_index.h" namespace arrow { class DataType; @@ -133,6 +134,64 @@ 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 { + return ReadArrowBatches( + file_name, file_size, + [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. + // 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(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"); + } + 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) { + 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))) { + continue; + } + } + PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry, serializer_->FromRow(row)); + 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..fa839b3f3 100644 --- a/src/paimon/core/manifest/manifest_list_test.cpp +++ b/src/paimon/core/manifest/manifest_list_test.cpp @@ -22,9 +22,11 @@ #include #include +#include "arrow/io/memory.h" +#include "arrow/ipc/api.h" #include "arrow/type.h" #include "gtest/gtest.h" -#include "paimon/core/core_options.h" +#include "paimon/common/utils/path_util.h" #include "paimon/core/manifest/manifest_file_meta.h" #include "paimon/core/snapshot.h" #include "paimon/core/stats/simple_stats.h" @@ -34,6 +36,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 +100,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 +113,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 +130,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..a4f60fea1 --- /dev/null +++ b/src/paimon/core/manifest/manifest_row_range_test.cpp @@ -0,0 +1,705 @@ +/* + * 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 "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/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/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" +#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/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) { + 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(FileKind::Add(), 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 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; + 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, 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(); + 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)); + } +} + +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, 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: + 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); +} + +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) { + constexpr int32_t kVersionedFileFieldIndex = ManifestEntry::kFileFieldIndex + 1; + 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[kVersionedFileFieldIndex]); + auto file_fields = checked_pointer_cast(file->type())->fields(); + auto file_columns = file->fields(); + if (mode == 0) { + file_fields.erase(file_fields.begin() + DataFileMeta::kFirstRowIdFieldIndex); + file_columns.erase(file_columns.begin() + DataFileMeta::kFirstRowIdFieldIndex); + } else if (mode == 1) { + 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()); + 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[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(); + const std::string name = fmt::format("manifest-evolved-{}", mode); + ASSERT_OK(WriteArray(dir->Str(), name, evolved)); + 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); + } + } + } +} + +} // 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/utils/objects_file.h b/src/paimon/core/utils/objects_file.h index e65ee75ab..2a99e98d1 100644 --- a/src/paimon/core/utils/objects_file.h +++ b/src/paimon/core/utils/objects_file.h @@ -19,6 +19,7 @@ #pragma once #include +#include #include #include #include @@ -27,6 +28,8 @@ #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" @@ -44,8 +47,6 @@ #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" namespace paimon { @@ -92,8 +93,7 @@ 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, @@ -113,8 +113,15 @@ class ObjectsFile { std::string compression_; std::shared_ptr cache_; - Result ReadFileSegment(const std::string& file_path, - const std::optional& file_size) const; + 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 +190,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)); @@ -256,22 +239,126 @@ Result> ObjectsFile::OpenForRead( } 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::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"); + } + // 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(), data->capacity()) {} + 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); +} - 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 +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); + } + const std::string path = path_factory_->ToPath(file_name); + auto key = CacheKey::ForKind(path, /*position=*/0, /*length=*/-1, CacheKind::MANIFEST); + std::optional read_status; + std::vector> batches; + auto cached = cache_->Get(key, [&](const std::shared_ptr&) { + return SerializeArrowBatches(file_name, file_size, &batches, &read_status); + }); + const bool usable = cached.ok() && cached.value() && cached.value()->GetSegment().Data(); + 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. + PAIMON_RETURN_NOT_OK(read_status.value()); + for (const auto& batch : batches) { + PAIMON_RETURN_NOT_OK(consumer(batch)); + } + return Status::OK(); } - return MemorySegment::Wrap(bytes); + 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()); + // Cached batches were already aligned by ManifestMetaReader before serialization. + PAIMON_RETURN_NOT_OK(consumer(array)); + } + return Status::OK(); } template