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
57 changes: 49 additions & 8 deletions include/paimon/utils/prefetch_cache_config.h
Original file line number Diff line number Diff line change
Expand Up @@ -34,10 +34,8 @@ namespace paimon {
/// ReadAheadCache to balance memory usage, I/O efficiency, and latency hiding.
class PAIMON_EXPORT CacheConfig {
public:
CacheConfig();
CacheConfig(uint64_t range_size_limit, uint64_t hole_size_limit, uint64_t pre_buffer_limit);

/// Returns the maximum allowed size (in bytes) for a single cached range.
/// Defaults to 32 MiB.
uint64_t GetRangeSizeLimit() const {
return range_size_limit_;
}
Expand All @@ -47,7 +45,8 @@ class PAIMON_EXPORT CacheConfig {
range_size_limit_ = range_size_limit;
}

/// Returns the maximum gap size (in bytes) considered mergeable between adjacent ranges.
/// Returns the maximum gap size (in bytes) considered mergeable between
/// adjacent ranges. Defaults to 8 KiB.
uint64_t GetHoleSizeLimit() const {
return hole_size_limit_;
}
Expand All @@ -57,7 +56,8 @@ class PAIMON_EXPORT CacheConfig {
hole_size_limit_ = hole_size_limit;
}

/// Returns the maximum size to pre-buffer ahead of the current read position.
/// Returns the maximum size to pre-buffer ahead of the current read
/// position. Defaults to 256 MiB.
uint64_t GetPreBufferLimit() const {
return pre_buffer_limit_;
}
Expand All @@ -67,10 +67,51 @@ class PAIMON_EXPORT CacheConfig {
pre_buffer_limit_ = pre_buffer_limit;
}

/// Returns the granularity (in bytes) of the block cache entries serving the
/// small reads that the prefetched ranges do not cover. Defaults to 64 KiB.
uint64_t GetBlockSize() const {
return block_size_;
}

/// Sets the granularity (in bytes) of the block cache entries.
void SetBlockSize(uint64_t block_size) {
block_size_ = block_size;
}

/// Returns the maximum total size (in bytes) of the block cache entries of
/// one file. Zero disables the block cache. Defaults to 1 MiB.
uint64_t GetBlockCacheLimit() const {
return block_cache_limit_;
}

/// Sets the maximum total size (in bytes) of the block cache entries of one
/// file. Zero disables the block cache.
void SetBlockCacheLimit(uint64_t block_cache_limit) {
block_cache_limit_ = block_cache_limit;
}

private:
uint64_t range_size_limit_;
uint64_t hole_size_limit_;
uint64_t pre_buffer_limit_;
// The defaults are aligned with the reader's request granularity and with
// realistic data file sizes:
// - range_size_limit matches the parquet reader's 32 MiB request blocks
// (Arrow ReadRangeCache's own range limit); a smaller limit cuts entries
// below the request size, so a request can never be served from one piece.
// - pre_buffer_limit must exceed the LARGEST single read a reader issues
// (coalesced column-chunk reads of ~128 MiB were observed): fetches are
// only dispatched up to this window, so a request reaching past it can
// never be served and falls back to a second fetch of the same bytes.
uint64_t range_size_limit_ = 32 * 1024 * 1024;
uint64_t hole_size_limit_ = 8 * 1024;
uint64_t pre_buffer_limit_ = 256 * 1024 * 1024;
// Blocks are aligned to the END of the file, so a block never reaches past
// EOF. 64 KiB is the granularity the reads no prefetched range covers are
// shared at: small enough that a metadata read at the tail of a file is
// served by one block instead of straddling two, large enough that a block
// fetch does not pull in much more than the reads ask for.
uint64_t block_size_ = 64 * 1024;
// One block is enough for the metadata tail of a file; the limit only
// bounds the pathological case, as blocks are never evicted.
uint64_t block_cache_limit_ = 1024 * 1024;
};

} // namespace paimon
2 changes: 2 additions & 0 deletions src/paimon/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,7 @@ set(PAIMON_COMMON_SRCS
common/data/shredding/shredding_file_reader.cpp
common/utils/delta_varint_compressor.cpp
common/utils/fields_comparator.cpp
common/utils/file_block_cache.cpp
common/utils/path_util.cpp
common/utils/range.cpp
common/utils/read_ahead_cache.cpp
Expand Down Expand Up @@ -667,6 +668,7 @@ if(PAIMON_BUILD_TESTS)
common/utils/roaring_bitmap64_test.cpp
common/utils/range_helper_test.cpp
common/utils/read_ahead_cache_test.cpp
common/utils/file_block_cache_test.cpp
common/io/cache/lru_cache_test.cpp
common/io/cache/cache_manager_test.cpp
common/utils/byte_range_combiner_test.cpp
Expand Down
25 changes: 18 additions & 7 deletions src/paimon/common/io/cache_input_stream_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,10 @@ namespace paimon::test {
class CacheInputStreamTest : public ::testing::Test {
public:
void SetUp() override {
pool_ = GetDefaultPool();
// A pool of its own, so that a cache buffer outliving the pool it was
// allocated from shows up instead of being covered by the global pool,
// which never goes away.
pool_ = std::shared_ptr<MemoryPool>(GetMemoryPool());
test_dir_ = UniqueTestDirectory::Create();
ASSERT_TRUE(test_dir_);
content_ = "abcdefghijklmnopqrstuvwxyz0123456789";
Expand All @@ -60,9 +63,14 @@ class CacheInputStreamTest : public ::testing::Test {

std::shared_ptr<ReadAheadCache> CreateCache(std::vector<ByteRange> ranges) {
auto stream = OpenFile();
CacheConfig config(/*range_size_limit=*/1024,
/*hole_size_limit=*/0, /*pre_buffer_limit=*/1024 * 1024);
auto cache = std::make_shared<ReadAheadCache>(std::move(stream), config, pool_);
CacheConfig config;
config.SetRangeSizeLimit(1024);
config.SetHoleSizeLimit(0);
config.SetPreBufferLimit(1024 * 1024);
// The file size is left unknown so the block cache stays off: these
// tests exercise the fallback of CacheInputStream on a cache miss.
auto cache =
std::make_shared<ReadAheadCache>(std::move(stream), config, /*file_size=*/0, pool_);
EXPECT_OK(cache->Init(std::move(ranges)));
return cache;
}
Expand Down Expand Up @@ -204,9 +212,12 @@ TEST_F(CacheInputStreamTest, TestReadAsyncCacheReadError) {
ASSERT_OK_AND_ASSIGN(auto fs, FileSystemFactory::Get("local", file_path_, {}));
ASSERT_OK_AND_ASSIGN(auto cache_stream, fs->Open(file_path_));
ASSERT_OK_AND_ASSIGN(auto underlying, fs->Open(file_path_));
CacheConfig config(/*range_size_limit=*/1024,
/*hole_size_limit=*/0, /*pre_buffer_limit=*/1024 * 1024);
auto cache = std::make_shared<ReadAheadCache>(std::move(cache_stream), config, pool_);
CacheConfig config;
config.SetRangeSizeLimit(1024);
config.SetHoleSizeLimit(0);
config.SetPreBufferLimit(1024 * 1024);
auto cache = std::make_shared<ReadAheadCache>(std::move(cache_stream), config,
/*file_size=*/0, pool_);
ASSERT_OK(cache->Init(std::vector<ByteRange>{{0, 10}}));

// Now activate IOHook so that the prefetch IO (triggered by cache_->Read -> PreBuffer)
Expand Down
60 changes: 60 additions & 0 deletions src/paimon/common/memory/bytes_utils.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

#pragma once

#include <cstddef>
#include <memory>
#include <utility>

#include "paimon/memory/bytes.h"
#include "paimon/memory/memory_pool.h"

namespace paimon {

/// Allocate a shared buffer that keeps the pool it was allocated from alive.
///
/// A `Bytes` holds a raw pointer to its pool and frees its allocation through
/// it, so it must not outlive the pool. Holding both in one owner is enough as
/// long as that owner drops the buffer, but not once the buffer is handed to an
/// asynchronous read: the callback of the read owns a reference to the buffer,
/// and a stream that destroys the callback after it has resolved the read - the
/// object store streams do - lets an IO thread drop the last reference to the
/// buffer after the owner and the pool are already gone. Binding the pool to
/// the buffer makes the buffer keep its own allocator alive, wherever its last
/// reference is dropped.
///
/// @param size Number of bytes to allocate.
/// @param pool Memory pool to allocate from, which the returned buffer keeps
/// alive.
inline std::shared_ptr<Bytes> AllocateBytesKeepingPoolAlive(
size_t size, const std::shared_ptr<MemoryPool>& pool) {
// The pool is declared before the buffer, so the holder destroys the buffer
// first and the pool it was allocated from second.
struct BytesWithMemoryPool {
std::shared_ptr<MemoryPool> pool;
std::shared_ptr<Bytes> bytes;
};
auto holder = std::make_shared<BytesWithMemoryPool>(
BytesWithMemoryPool{pool, std::make_shared<Bytes>(size, pool.get())});
Bytes* bytes = holder->bytes.get();
return std::shared_ptr<Bytes>(std::move(holder), bytes);
}

} // namespace paimon
57 changes: 57 additions & 0 deletions src/paimon/common/metrics/atomic_counter_pair.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

#pragma once

#include <atomic>
#include <cstdint>

namespace paimon {

/// How many requests of one kind were made and how many bytes they covered,
/// which is the shape every counter of the read path has. Grouping the two
/// keeps a counter and its byte counter from drifting apart.
///
/// The atomics are relaxed throughout: the counters are only reported, never
/// used to order anything, and the read path increments them on every request.
struct AtomicCounterPair {
std::atomic<uint64_t> count{0};
std::atomic<uint64_t> bytes{0};

/// Record one request covering `size` bytes.
void Add(uint64_t size) {
count.fetch_add(1, std::memory_order_relaxed);
bytes.fetch_add(size, std::memory_order_relaxed);
}

void Reset() {
count.store(0, std::memory_order_relaxed);
bytes.store(0, std::memory_order_relaxed);
}

uint64_t Count() const {
return count.load(std::memory_order_relaxed);
}

uint64_t Bytes() const {
return bytes.load(std::memory_order_relaxed);
}
};

} // namespace paimon
9 changes: 8 additions & 1 deletion src/paimon/common/reader/prefetch_file_batch_reader_impl.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,9 @@ Result<std::unique_ptr<PrefetchFileBatchReaderImpl>> PrefetchFileBatchReaderImpl
if (batch_size <= 0) {
return Status::Invalid("batch size should be greater than 0.");
}
if (data_file_size < 0) {
return Status::Invalid("data file size should not be negative.");
}
if (reader_builder == nullptr) {
return Status::Invalid("reader_builder should not be nullptr.");
}
Expand All @@ -236,7 +239,11 @@ Result<std::unique_ptr<PrefetchFileBatchReaderImpl>> PrefetchFileBatchReaderImpl
if (io_metrics) {
input_stream = std::make_shared<MetricsInputStream>(input_stream, io_metrics);
}
cache = std::make_shared<ReadAheadCache>(input_stream, cache_config, pool);
// The file size lets the cache align its blocks to the end of the file,
// where the metadata the readers read before any range is registered
// lives. A zero size means unknown and disables the block cache.
cache = std::make_shared<ReadAheadCache>(input_stream, cache_config,
static_cast<uint64_t>(data_file_size), pool);
}
std::vector<std::future<Result<std::unique_ptr<FileBatchReader>>>> futures;
for (uint32_t i = 0; i < prefetch_max_parallel_num; i++) {
Expand Down
18 changes: 14 additions & 4 deletions src/paimon/common/reader/prefetch_file_batch_reader_impl_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -664,10 +664,10 @@ TEST_F(PrefetchFileBatchReaderImplTest, WorkloopSetReadStatusWhenCacheInitFailed
int32_t batch_size = 5;
int32_t prefetch_max_parallel_num = 1;
MockFormatReaderBuilder reader_builder(data_array, data_type_, batch_size);
CacheConfig invalid_cache_config(
/*range_size_limit=*/4 * 1024,
/*hole_size_limit=*/8 * 1024,
/*pre_buffer_limit=*/128 * 1024);
CacheConfig invalid_cache_config;
invalid_cache_config.SetRangeSizeLimit(4 * 1024);
invalid_cache_config.SetHoleSizeLimit(8 * 1024);
invalid_cache_config.SetPreBufferLimit(128 * 1024);

ASSERT_OK_AND_ASSIGN(
auto reader,
Expand Down Expand Up @@ -930,6 +930,16 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestInvalidCase) {
/*initialize_read_ranges=*/true, /*read_ahead_cache_enabled=*/true, CacheConfig(),
/*enable_io_metrics=*/false, pool_, GetArrowPool(pool_)));
}
{
ASSERT_NOK_WITH_MSG(
PrefetchFileBatchReaderImpl::Create(
data_file_path, /*data_file_size=*/-1, &reader_builder, mock_fs_,
prefetch_max_parallel_num, batch_size, prefetch_max_parallel_num * 2,
/*enable_adaptive_prefetch_strategy=*/false, executor_,
/*initialize_read_ranges=*/true, /*read_ahead_cache_enabled=*/true, CacheConfig(),
/*enable_io_metrics=*/false, pool_, GetArrowPool(pool_)),
"data file size should not be negative");
}
{
ASSERT_NOK(PrefetchFileBatchReaderImpl::Create(
data_file_path, /*data_file_size=*/0, &reader_builder, mock_fs_,
Expand Down
Loading
Loading