Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 28 additions & 16 deletions docs/source/user_guide/manifest_cache.rst
Original file line number Diff line number Diff line change
Expand Up @@ -21,18 +21,24 @@ Manifest Cache
Overview
--------

Paimon C++ caches raw manifest file bytes at the ``ObjectsFile<T>::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<T>``.
Paimon C++ caches decoded, schema-aligned manifest batches as Arrow IPC streams
in ``ObjectsFile<T>``. 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<T>``.

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
-------------
Expand Down Expand Up @@ -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<T>`` 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.
9 changes: 9 additions & 0 deletions include/paimon/cache/cache.h
Original file line number Diff line number Diff line change
Expand Up @@ -88,19 +88,28 @@ 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<CacheKey>& key) const;

bool operator==(const CacheValue& other) const;

private:
MemorySegment segment_;
CacheCallback callback_;
int64_t memory_usage_;
};

} // namespace paimon
1 change: 1 addition & 0 deletions src/paimon/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 11 additions & 1 deletion src/paimon/common/io/cache/cache.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,19 +19,29 @@

#include "paimon/cache/cache.h"

#include <algorithm>
#include <utility>

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<int64_t>(segment.Size(), memory_usage)) {}

CacheValue::~CacheValue() = default;

const MemorySegment& CacheValue::GetSegment() const {
return segment_;
}

int64_t CacheValue::GetMemoryUsage() const {
return memory_usage_;
}

void CacheValue::OnEvict(const std::shared_ptr<CacheKey>& key) const {
if (callback_) {
callback_(key);
Expand Down
2 changes: 1 addition & 1 deletion src/paimon/common/io/cache/lru_cache.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ LruCache::LruCache(int64_t max_weight)
.expire_after_access_ms = -1,
.weigh_func = [](const std::shared_ptr<CacheKey>& /*key*/,
const std::shared_ptr<CacheValue>& value) -> int64_t {
return value ? value->GetSegment().Size() : 0;
return value ? value->GetMemoryUsage() : 0;
},
.removal_callback =
[](const std::shared_ptr<CacheKey>& key, const std::shared_ptr<CacheValue>& value,
Expand Down
2 changes: 1 addition & 1 deletion src/paimon/common/io/cache/lru_cache.h
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
34 changes: 34 additions & 0 deletions src/paimon/common/io/cache/lru_cache_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<CacheValue>(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<std::shared_ptr<CacheKey>> evicted;
auto first = std::make_shared<CacheValue>(
segment, [&](const std::shared_ptr<CacheKey>& 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<std::shared_ptr<CacheKey>>{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<CacheValue> loaded,
small_cache.Get(MakeKey(0),
[&](const std::shared_ptr<CacheKey>&)
-> Result<std::shared_ptr<CacheValue>> { 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
Expand Down
4 changes: 4 additions & 0 deletions src/paimon/core/io/data_file_meta.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
12 changes: 6 additions & 6 deletions src/paimon/core/io/data_file_meta_serializer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ Result<BinaryRow> DataFileMetaSerializer::ToRow(const std::shared_ptr<DataFileMe
BinaryRowWriter writer(&row, 32 * 1024, pool_.get());
writer.WriteString(0, BinaryString::FromString(meta->file_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());
Expand Down Expand Up @@ -86,9 +86,9 @@ Result<BinaryRow> DataFileMetaSerializer::ToRow(const std::shared_ptr<DataFileMe
writer.WriteString(17, BinaryString::FromString(meta->external_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);
Expand All @@ -110,7 +110,7 @@ Result<std::shared_ptr<DataFileMeta>> 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);
Expand Down Expand Up @@ -155,8 +155,8 @@ Result<std::shared_ptr<DataFileMeta>> DataFileMetaSerializer::FromRow(
external_path = row.GetString(17).ToString();
}
std::optional<int64_t> 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<std::vector<std::string>> write_cols;
Expand Down
53 changes: 52 additions & 1 deletion src/paimon/core/manifest/index_manifest_file_handler_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,10 @@
#include <vector>

#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"
Expand All @@ -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 {
Expand All @@ -50,7 +54,8 @@ class IndexManifestFileHandlerTest : public testing::Test {
}

Result<std::unique_ptr<IndexManifestFile>> 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>& cache = nullptr) const {
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<FileFormat> file_format,
FileFormatFactory::Get(file_format_identifier, {}));
auto schema = arrow::schema({arrow::field("f0", arrow::int32())});
Expand All @@ -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);
}
Expand Down Expand Up @@ -100,6 +106,51 @@ class IndexManifestFileHandlerTest : public testing::Test {
std::unique_ptr<UniqueTestDirectory> dir_;
};

TEST_F(IndexManifestFileHandlerTest, DecodedCacheReusesIpcAcrossReaders) {
auto cache = std::make_shared<CountingRoutingCache>(CacheKind::MANIFEST, 1024 * 1024);
ASSERT_OK_AND_ASSIGN(std::unique_ptr<IndexManifestFile> writer,
CreateManifestFile(2, "avro", cache));
const std::vector<IndexManifestEntry> 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<std::string, int64_t>;
ASSERT_OK_AND_ASSIGN(WrittenFile written, writer->WriteWithoutRolling(expected));
std::vector<IndexManifestEntry> 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<CacheValue> cached,
cache->Get(key,
[](const std::shared_ptr<CacheKey>&) -> Result<std::shared_ptr<CacheValue>> {
return Status::Invalid("index manifest was not cached");
}));
ASSERT_TRUE(cached);
const auto& segment = cached->GetSegment();
auto buffer = std::make_shared<arrow::Buffer>(reinterpret_cast<const uint8_t*>(segment.Data()),
segment.Size());
auto ipc = arrow::ipc::RecordBatchStreamReader::Open(
std::make_shared<arrow::io::BufferReader>(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<IndexManifestFile> reader,
CreateManifestFile(2, "avro", cache));
std::vector<IndexManifestEntry> 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));

Expand Down
3 changes: 3 additions & 0 deletions src/paimon/core/manifest/manifest_entry.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<arrow::DataType>& DataType();
static int64_t RecordCount(const std::vector<ManifestEntry>& manifest_entries);
static std::optional<int64_t> NullableRecordCount(
Expand Down
4 changes: 2 additions & 2 deletions src/paimon/core/manifest/manifest_entry_serializer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ Result<ManifestEntry> 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");
}
Expand All @@ -80,7 +80,7 @@ Result<BinaryRow> 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;
}
Expand Down
Loading
Loading