From 05b13249fac917bfe00efc94f2dc3d55091a9e4e Mon Sep 17 00:00:00 2001 From: zhengyu Date: Thu, 16 Jul 2026 18:48:48 +0800 Subject: [PATCH] [fix](filecache) Support memory file cache resource management ### What problem does this PR solve? Issue Number: None Related PR: None Problem Summary: Memory file cache storage still followed disk-oriented LRU persistence behavior and skipped capacity-based resource limit monitoring. The adapted monitor also crashed when a memory cache had zero capacity, its pressure-mode flags were accessed concurrently without synchronization, and the production reset path validated the literal memory cache path with statfs. An all-cache reset could therefore partially mutate disk caches before failing on a memory or invalid target. Adapt the required SelectDB fixes to current Doris master, preserve the newer sharded memory storage duplicate-key handling already present on master, make cross-thread pressure signals atomic, handle zero-capacity statistics safely, and prevalidate storage-aware reset targets before applying any capacity changes. ### Release note Improve memory file cache resource monitoring and capacity reset behavior while avoiding disk-only persistence and validation work. ### Check List (For Author) - Test: Unit Test - `BlockFileCacheMemModeTest.*:FileCacheActionTest.reset_memory_capacity_without_memory_directory` (15 tests passed) - BE compilation with `./build.sh --be` - BE clang-format and clang-tidy checks - Behavior changed: Yes. Memory file cache uses capacity percentage for pressure modes, synchronizes pressure signals, safely handles zero capacity, skips disk-only LRU persistence, and resets capacity without filesystem validation. - Does this need documentation: No --- be/src/io/cache/block_file_cache.cpp | 256 +++++++--- be/src/io/cache/block_file_cache.h | 7 +- be/src/io/cache/block_file_cache_factory.cpp | 32 +- be/src/io/cache/lru_queue_recorder.cpp | 3 + ...block_file_cache_mem_storage_mode_test.cpp | 465 ++++++++++++++++++ .../service/http/file_cache_action_test.cpp | 27 + 6 files changed, 704 insertions(+), 86 deletions(-) create mode 100644 be/test/io/cache/block_file_cache_mem_storage_mode_test.cpp diff --git a/be/src/io/cache/block_file_cache.cpp b/be/src/io/cache/block_file_cache.cpp index b44679a94fa150..c31eeefe25de67 100644 --- a/be/src/io/cache/block_file_cache.cpp +++ b/be/src/io/cache/block_file_cache.cpp @@ -537,7 +537,8 @@ Status BlockFileCache::initialize() { Status BlockFileCache::initialize_unlocked(std::lock_guard& cache_lock) { DCHECK(!_is_initialized); _is_initialized = true; - if (config::file_cache_background_lru_dump_tail_record_num > 0) { + const bool memory_storage = is_memory_storage(); + if (!memory_storage && config::file_cache_background_lru_dump_tail_record_num > 0) { // requirements: // 1. restored data should not overwrite the last dump // 2. restore should happen before load and async load @@ -547,6 +548,9 @@ Status BlockFileCache::initialize_unlocked(std::lock_guard& cache_lo restore_lru_queues_from_disk(cache_lock); } RETURN_IF_ERROR(_storage->init(this)); + if (memory_storage) { + LOG(INFO) << "Skip LRU dump and restore for memory storage, path=" << _cache_base_path; + } if (auto* fs_storage = dynamic_cast(_storage.get())) { if (auto* meta_store = fs_storage->get_meta_store()) { @@ -561,10 +565,13 @@ Status BlockFileCache::initialize_unlocked(std::lock_guard& cache_lo _cache_background_block_lru_update_thread = std::thread(&BlockFileCache::run_background_block_lru_update, this); - // Initialize LRU dump thread and restore queues - _cache_background_lru_dump_thread = std::thread(&BlockFileCache::run_background_lru_dump, this); - _cache_background_lru_log_replay_thread = - std::thread(&BlockFileCache::run_background_lru_log_replay, this); + if (!memory_storage) { + // Initialize LRU dump thread and restore queues + _cache_background_lru_dump_thread = + std::thread(&BlockFileCache::run_background_lru_dump, this); + _cache_background_lru_log_replay_thread = + std::thread(&BlockFileCache::run_background_lru_log_replay, this); + } return Status::OK(); } @@ -1263,7 +1270,8 @@ bool BlockFileCache::try_reserve(const UInt128Wrapper& hash, const CacheContext& } // use this strategy in scenarios where there is insufficient disk capacity or insufficient number of inodes remaining // directly eliminate 5 times the size of the space - if (_disk_resource_limit_mode) { + const bool disk_resource_limit_mode = _disk_resource_limit_mode.load(std::memory_order_relaxed); + if (disk_resource_limit_mode) { size = 5 * size; } @@ -1292,12 +1300,12 @@ bool BlockFileCache::try_reserve(const UInt128Wrapper& hash, const CacheContext& size_t max_size = queue.get_max_size(); auto is_overflow = [&] { - return _disk_resource_limit_mode ? removed_size < size - : cur_cache_size + size - removed_size > _capacity || - (queue_size + size - removed_size > max_size) || - (query_context_cache_size + size - - (removed_size + ghost_remove_size) > - query_context->get_max_cache_size()); + return disk_resource_limit_mode ? removed_size < size + : cur_cache_size + size - removed_size > _capacity || + (queue_size + size - removed_size > max_size) || + (query_context_cache_size + size - + (removed_size + ghost_remove_size) > + query_context->get_max_cache_size()); }; /// Select the cache from the LRU queue held by query for expulsion. @@ -1522,7 +1530,7 @@ bool BlockFileCache::is_overflow(size_t removed_size, size_t need_size, size_t c ret = (removed_size < need_size); return ret; } - if (_disk_resource_limit_mode) { + if (_disk_resource_limit_mode.load(std::memory_order_relaxed)) { ret = (removed_size < need_size); } else { ret = (cur_cache_size + need_size - removed_size > _capacity); @@ -1980,13 +1988,41 @@ int disk_used_percentage(const std::string& path, std::pair* percent) return 0; } +namespace { + +const char* file_cache_storage_type_to_string(FileCacheStorageType type) { + switch (type) { + case FileCacheStorageType::DISK: + return "disk"; + case FileCacheStorageType::MEMORY: + return "memory"; + } + DORIS_CHECK(false) << "Unknown file cache storage type: " << type; + __builtin_unreachable(); +} + +} // namespace + +bool BlockFileCache::is_memory_storage() const { + return _storage->get_type() == FileCacheStorageType::MEMORY; +} + +size_t BlockFileCache::calc_size_percentage_unlocked(std::lock_guard&) const { + if (_capacity == 0) { + return 0; + } + return static_cast(static_cast(_cur_cache_size) / + static_cast(_capacity) * 100); +} + std::string BlockFileCache::reset_capacity(size_t new_capacity) { using namespace std::chrono; int64_t space_released = 0; size_t old_capacity = 0; + size_t size_percentage = 0; std::stringstream ss; ss << "finish reset_capacity, path=" << _cache_base_path; - auto adjust_start_time = steady_clock::time_point(); + auto adjust_start_time = steady_clock::now(); { SCOPED_CACHE_LOCK(_mutex, this); if (new_capacity < _capacity && new_capacity < _cur_cache_size) { @@ -1998,14 +2034,14 @@ std::string BlockFileCache::reset_capacity(size_t new_capacity) { if (need_remove_size <= 0) { break; } - need_remove_size -= entry_size; - space_released += entry_size; - queue_released += entry_size; auto* cell = get_cell(entry_key, entry_offset, cache_lock); if (!cell->releasable()) { cell->file_block->set_deleting(); continue; } + need_remove_size -= entry_size; + space_released += entry_size; + queue_released += entry_size; to_evict.push_back(cell); } for (auto& cell : to_evict) { @@ -2023,17 +2059,22 @@ std::string BlockFileCache::reset_capacity(size_t new_capacity) { ss << " index_queue released " << queue_released; queue_released = remove_blocks(_ttl_queue); ss << " ttl_queue released " << queue_released; - - _disk_resource_limit_mode = true; - _disk_limit_mode_metrics->set_value(1); ss << " total_space_released=" << space_released; } old_capacity = _capacity; _capacity = new_capacity; + size_percentage = calc_size_percentage_unlocked(cache_lock); _cache_capacity_metrics->set_value(_capacity); + _cur_cache_size_metrics->set_value(_cur_cache_size); + _disk_limit_mode_metrics->set_value( + _disk_resource_limit_mode.load(std::memory_order_relaxed)); + _need_evict_cache_in_advance_metrics->set_value( + _need_evict_cache_in_advance.load(std::memory_order_relaxed)); } - auto use_time = duration_cast(steady_clock::time_point() - adjust_start_time); + auto use_time = duration_cast(steady_clock::now() - adjust_start_time); LOG(INFO) << "Finish tag deleted block. path=" << _cache_base_path + << " storage_type=" << file_cache_storage_type_to_string(_storage->get_type()) + << " size_percent=" << size_percentage << " use_time=" << cast_set(use_time.count()); ss << " old_capacity=" << old_capacity << " new_capacity=" << new_capacity; LOG(INFO) << ss.str(); @@ -2041,13 +2082,47 @@ std::string BlockFileCache::reset_capacity(size_t new_capacity) { } void BlockFileCache::check_disk_resource_limit() { - if (_storage->get_type() != FileCacheStorageType::DISK) { + // ATTN: due to that can be changed dynamically, set it to default value if it's invalid + // FIXME: reject with config validator + if (config::file_cache_enter_disk_resource_limit_mode_percent <= + config::file_cache_exit_disk_resource_limit_mode_percent) { + LOG_WARNING("config error, set to default value") + .tag("enter", config::file_cache_enter_disk_resource_limit_mode_percent) + .tag("exit", config::file_cache_exit_disk_resource_limit_mode_percent); + config::file_cache_enter_disk_resource_limit_mode_percent = 88; + config::file_cache_exit_disk_resource_limit_mode_percent = 80; + } + size_t size_percentage = 0; + { + SCOPED_CACHE_LOCK(_mutex, this); + size_percentage = calc_size_percentage_unlocked(cache_lock); + } + auto is_insufficient = [](const int& percentage) { + return percentage >= config::file_cache_enter_disk_resource_limit_mode_percent; + }; + const bool previous_mode = _disk_resource_limit_mode.load(std::memory_order_relaxed); + bool current_mode = previous_mode; + if (is_memory_storage()) { + bool is_size_insufficient = is_insufficient(cast_set(size_percentage)); + if (is_size_insufficient) { + current_mode = true; + } else if (current_mode && + size_percentage < config::file_cache_exit_disk_resource_limit_mode_percent) { + current_mode = false; + } + _disk_resource_limit_mode.store(current_mode, std::memory_order_relaxed); + _disk_limit_mode_metrics->set_value(current_mode); + if (previous_mode != current_mode) { + LOG(INFO) << "Memory file cache resource limit mode changed: file_cache=" + << get_base_path() << " enabled=" << current_mode + << " size_percent=" << size_percentage; + } return; } - bool previous_mode = _disk_resource_limit_mode; if (_capacity > _cur_cache_size) { - _disk_resource_limit_mode = false; + current_mode = false; + _disk_resource_limit_mode.store(current_mode, std::memory_order_relaxed); _disk_limit_mode_metrics->set_value(0); } std::pair percent; @@ -2057,37 +2132,24 @@ void BlockFileCache::check_disk_resource_limit() { return; } auto [space_percentage, inode_percentage] = percent; - auto is_insufficient = [](const int& percentage) { - return percentage >= config::file_cache_enter_disk_resource_limit_mode_percent; - }; DCHECK_GE(space_percentage, 0); DCHECK_LE(space_percentage, 100); DCHECK_GE(inode_percentage, 0); DCHECK_LE(inode_percentage, 100); - // ATTN: due to that can be changed dynamically, set it to default value if it's invalid - // FIXME: reject with config validator - if (config::file_cache_enter_disk_resource_limit_mode_percent < - config::file_cache_exit_disk_resource_limit_mode_percent) { - LOG_WARNING("config error, set to default value") - .tag("enter", config::file_cache_enter_disk_resource_limit_mode_percent) - .tag("exit", config::file_cache_exit_disk_resource_limit_mode_percent); - config::file_cache_enter_disk_resource_limit_mode_percent = 88; - config::file_cache_exit_disk_resource_limit_mode_percent = 80; - } bool is_space_insufficient = is_insufficient(space_percentage); bool is_inode_insufficient = is_insufficient(inode_percentage); if (is_space_insufficient || is_inode_insufficient) { - _disk_resource_limit_mode = true; - _disk_limit_mode_metrics->set_value(1); - } else if (_disk_resource_limit_mode && + current_mode = true; + } else if (current_mode && (space_percentage < config::file_cache_exit_disk_resource_limit_mode_percent) && (inode_percentage < config::file_cache_exit_disk_resource_limit_mode_percent)) { - _disk_resource_limit_mode = false; - _disk_limit_mode_metrics->set_value(0); + current_mode = false; } - if (previous_mode != _disk_resource_limit_mode) { + _disk_resource_limit_mode.store(current_mode, std::memory_order_relaxed); + _disk_limit_mode_metrics->set_value(current_mode); + if (previous_mode != current_mode) { // add log for disk resource limit mode switching - if (_disk_resource_limit_mode) { + if (current_mode) { LOG(WARNING) << "Entering disk resource limit mode: file_cache=" << get_base_path() << " space_percent=" << space_percentage << " inode_percent=" << inode_percentage @@ -2101,7 +2163,7 @@ void BlockFileCache::check_disk_resource_limit() { << " inode_percent=" << inode_percentage << " exit threshold=" << config::file_cache_exit_disk_resource_limit_mode_percent; } - } else if (_disk_resource_limit_mode) { + } else if (current_mode) { // print log for disk resource limit mode running, but less frequently LOG_EVERY_N(WARNING, 10) << "file_cache=" << get_base_path() << " space_percent=" << space_percentage @@ -2113,7 +2175,41 @@ void BlockFileCache::check_disk_resource_limit() { } void BlockFileCache::check_need_evict_cache_in_advance() { - if (_storage->get_type() != FileCacheStorageType::DISK) { + // ATTN: due to that can be changed dynamically, set it to default value if it's invalid + // FIXME: reject with config validator + if (config::file_cache_enter_need_evict_cache_in_advance_percent <= + config::file_cache_exit_need_evict_cache_in_advance_percent) { + LOG_WARNING("config error, set to default value") + .tag("enter", config::file_cache_enter_need_evict_cache_in_advance_percent) + .tag("exit", config::file_cache_exit_need_evict_cache_in_advance_percent); + config::file_cache_enter_need_evict_cache_in_advance_percent = 78; + config::file_cache_exit_need_evict_cache_in_advance_percent = 75; + } + size_t size_percentage = 0; + { + SCOPED_CACHE_LOCK(_mutex, this); + size_percentage = calc_size_percentage_unlocked(cache_lock); + } + auto is_insufficient = [](const int& percentage) { + return percentage >= config::file_cache_enter_need_evict_cache_in_advance_percent; + }; + const bool previous_mode = _need_evict_cache_in_advance.load(std::memory_order_relaxed); + bool current_mode = previous_mode; + if (is_memory_storage()) { + bool is_size_insufficient = is_insufficient(cast_set(size_percentage)); + if (is_size_insufficient) { + current_mode = true; + } else if (current_mode && + size_percentage < config::file_cache_exit_need_evict_cache_in_advance_percent) { + current_mode = false; + } + _need_evict_cache_in_advance.store(current_mode, std::memory_order_relaxed); + _need_evict_cache_in_advance_metrics->set_value(current_mode); + if (previous_mode != current_mode) { + LOG(INFO) << "Memory file cache evict-in-advance mode changed: file_cache=" + << get_base_path() << " enabled=" << current_mode + << " size_percent=" << size_percentage; + } return; } @@ -2124,41 +2220,26 @@ void BlockFileCache::check_need_evict_cache_in_advance() { return; } auto [space_percentage, inode_percentage] = percent; - int size_percentage = static_cast(_cur_cache_size * 100 / _capacity); - auto is_insufficient = [](const int& percentage) { - return percentage >= config::file_cache_enter_need_evict_cache_in_advance_percent; - }; DCHECK_GE(space_percentage, 0); DCHECK_LE(space_percentage, 100); DCHECK_GE(inode_percentage, 0); DCHECK_LE(inode_percentage, 100); - // ATTN: due to that can be changed dynamically, set it to default value if it's invalid - // FIXME: reject with config validator - if (config::file_cache_enter_need_evict_cache_in_advance_percent <= - config::file_cache_exit_need_evict_cache_in_advance_percent) { - LOG_WARNING("config error, set to default value") - .tag("enter", config::file_cache_enter_need_evict_cache_in_advance_percent) - .tag("exit", config::file_cache_exit_need_evict_cache_in_advance_percent); - config::file_cache_enter_need_evict_cache_in_advance_percent = 78; - config::file_cache_exit_need_evict_cache_in_advance_percent = 75; - } - bool previous_mode = _need_evict_cache_in_advance; bool is_space_insufficient = is_insufficient(space_percentage); bool is_inode_insufficient = is_insufficient(inode_percentage); - bool is_size_insufficient = is_insufficient(size_percentage); + bool is_size_insufficient = is_insufficient(cast_set(size_percentage)); if (is_space_insufficient || is_inode_insufficient || is_size_insufficient) { - _need_evict_cache_in_advance = true; - _need_evict_cache_in_advance_metrics->set_value(1); - } else if (_need_evict_cache_in_advance && + current_mode = true; + } else if (current_mode && (space_percentage < config::file_cache_exit_need_evict_cache_in_advance_percent) && (inode_percentage < config::file_cache_exit_need_evict_cache_in_advance_percent) && (size_percentage < config::file_cache_exit_need_evict_cache_in_advance_percent)) { - _need_evict_cache_in_advance = false; - _need_evict_cache_in_advance_metrics->set_value(0); + current_mode = false; } - if (previous_mode != _need_evict_cache_in_advance) { + _need_evict_cache_in_advance.store(current_mode, std::memory_order_relaxed); + _need_evict_cache_in_advance_metrics->set_value(current_mode); + if (previous_mode != current_mode) { // add log for evict cache in advance mode switching - if (_need_evict_cache_in_advance) { + if (current_mode) { LOG(WARNING) << "Entering evict cache in advance mode: " << "file_cache=" << get_base_path() << " space_percent=" << space_percentage @@ -2175,7 +2256,7 @@ void BlockFileCache::check_need_evict_cache_in_advance() { << " size_percent=" << size_percentage << " exit threshold=" << config::file_cache_exit_need_evict_cache_in_advance_percent; } - } else if (_need_evict_cache_in_advance) { + } else if (current_mode) { // print log for evict cache in advance mode running, but less frequently LOG_EVERY_N(WARNING, 10) << "file_cache=" << get_base_path() << " space_percent=" << space_percentage @@ -2197,7 +2278,7 @@ void BlockFileCache::run_background_monitor() { if (config::enable_evict_file_cache_in_advance) { check_need_evict_cache_in_advance(); } else { - _need_evict_cache_in_advance = false; + _need_evict_cache_in_advance.store(false, std::memory_order_relaxed); _need_evict_cache_in_advance_metrics->set_value(0); } @@ -2327,7 +2408,7 @@ void BlockFileCache::run_background_evict_in_advance() { batch = config::file_cache_evict_in_advance_batch_bytes; // Skip if eviction not needed or too many pending recycles - if (!_need_evict_cache_in_advance || + if (!_need_evict_cache_in_advance.load(std::memory_order_relaxed) || _recycle_keys.size_approx() >= config::file_cache_evict_in_advance_recycle_keys_num_threshold) { continue; @@ -2415,7 +2496,8 @@ bool BlockFileCache::try_reserve_during_async_load(size_t size, std::vector to_evict; auto collect_eliminate_fragments = [&](LRUQueue& queue) { for (const auto& [entry_key, entry_offset, entry_size] : queue) { - if (!_disk_resource_limit_mode || removed_size >= size) { + if (!_disk_resource_limit_mode.load(std::memory_order_relaxed) || + removed_size >= size) { break; } auto* cell = get_cell(entry_key, entry_offset, cache_lock); @@ -2448,7 +2530,7 @@ bool BlockFileCache::try_reserve_during_async_load(size_t size, std::string reason = "async load"; remove_file_blocks(to_evict, cache_lock, true, reason); - return !_disk_resource_limit_mode || removed_size >= size; + return !_disk_resource_limit_mode.load(std::memory_order_relaxed) || removed_size >= size; } void BlockFileCache::clear_need_update_lru_blocks() { @@ -2516,6 +2598,9 @@ size_t BlockFileCache::replay_lru_logs_once() { } void BlockFileCache::dump_lru_queues(bool force) { + if (is_memory_storage()) { + return; + } std::unique_lock dump_lock(_dump_lru_queues_mtx); if (config::file_cache_background_lru_dump_tail_record_num > 0 && !ExecEnv::GetInstance()->get_is_upgrading()) { @@ -2552,6 +2637,11 @@ void BlockFileCache::restore_lru_queues_from_disk(std::lock_guard& c std::map BlockFileCache::get_stats() { std::map stats; + size_t size_percentage = 0; + { + SCOPED_CACHE_LOCK(_mutex, this); + size_percentage = calc_size_percentage_unlocked(cache_lock); + } stats["hits_ratio"] = (double)_hit_ratio->get_value(); stats["hits_ratio_5m"] = (double)_hit_ratio_5m->get_value(); stats["hits_ratio_1h"] = (double)_hit_ratio_1h->get_value(); @@ -2594,8 +2684,12 @@ std::map BlockFileCache::get_stats() { (double)_lru_recorder_shadow_queue_element_count_metrics[FileCacheType::DISPOSABLE] ->get_value(); - stats["need_evict_cache_in_advance"] = (double)_need_evict_cache_in_advance; - stats["disk_resource_limit_mode"] = (double)_disk_resource_limit_mode; + stats["need_evict_cache_in_advance"] = + (double)_need_evict_cache_in_advance.load(std::memory_order_relaxed); + stats["disk_resource_limit_mode"] = + (double)_disk_resource_limit_mode.load(std::memory_order_relaxed); + stats["size_percentage"] = (double)size_percentage; + stats["storage_type"] = (double)_storage->get_type(); stats["total_removed_counts"] = (double)_num_removed_blocks->get_value(); stats["total_hit_counts"] = (double)_num_hit_blocks->get_value(); @@ -2647,8 +2741,14 @@ std::map BlockFileCache::get_stats_unsafe() { (double)_lru_recorder_shadow_queue_element_count_metrics[FileCacheType::DISPOSABLE] ->get_value(); - stats["need_evict_cache_in_advance"] = (double)_need_evict_cache_in_advance; - stats["disk_resource_limit_mode"] = (double)_disk_resource_limit_mode; + stats["need_evict_cache_in_advance"] = + (double)_need_evict_cache_in_advance.load(std::memory_order_relaxed); + stats["disk_resource_limit_mode"] = + (double)_disk_resource_limit_mode.load(std::memory_order_relaxed); + stats["size_percentage"] = _capacity == 0 ? 0 + : static_cast(_cur_cache_size) / + static_cast(_capacity) * 100; + stats["storage_type"] = (double)_storage->get_type(); stats["total_removed_counts"] = (double)_num_removed_blocks->get_value(); stats["total_hit_counts"] = (double)_num_hit_blocks->get_value(); diff --git a/be/src/io/cache/block_file_cache.h b/be/src/io/cache/block_file_cache.h index 9a01d07d136a62..491e40b94faa8a 100644 --- a/be/src/io/cache/block_file_cache.h +++ b/be/src/io/cache/block_file_cache.h @@ -467,6 +467,9 @@ class BlockFileCache { size_t get_used_cache_size_unlocked(FileCacheType type, std::lock_guard& cache_lock) const; + bool is_memory_storage() const; + size_t calc_size_percentage_unlocked(std::lock_guard&) const; + void check_disk_resource_limit(); void check_need_evict_cache_in_advance(); @@ -539,8 +542,8 @@ class BlockFileCache { std::thread _cache_background_block_lru_update_thread; std::atomic_bool _async_open_done {false}; // disk space or inode is less than the specified value - bool _disk_resource_limit_mode {false}; - bool _need_evict_cache_in_advance {false}; + std::atomic_bool _disk_resource_limit_mode {false}; + std::atomic_bool _need_evict_cache_in_advance {false}; bool _is_initialized {false}; // strategy diff --git a/be/src/io/cache/block_file_cache_factory.cpp b/be/src/io/cache/block_file_cache_factory.cpp index 0ffc9cea365d6f..2f3503adb7ae0e 100644 --- a/be/src/io/cache/block_file_cache_factory.cpp +++ b/be/src/io/cache/block_file_cache_factory.cpp @@ -36,6 +36,7 @@ #include #include +#include "common/cast_set.h" #include "common/config.h" #include "core/block/block.h" #include "information_schema/schema_scanner_helper.h" @@ -312,8 +313,8 @@ std::vector FileCacheFactory::get_base_paths() { return paths; } -std::string validate_capacity(const std::string& path, int64_t new_capacity, - int64_t& valid_capacity) { +std::string validate_disk_capacity(const std::string& path, int64_t new_capacity, + int64_t& valid_capacity) { struct statfs stat; if (statfs(path.c_str(), &stat) < 0) { auto ret = fmt::format("reset capacity {} statfs error {}. ", path, strerror(errno)); @@ -340,16 +341,35 @@ std::string validate_capacity(const std::string& path, int64_t new_capacity, return ""; } +std::string validate_capacity(BlockFileCache* cache, int64_t new_capacity, + int64_t& valid_capacity) { + if (cache->get_storage()->get_type() == FileCacheStorageType::MEMORY) { + valid_capacity = new_capacity; + if (valid_capacity <= 0) { + return fmt::format("The memory cache {} capacity must be greater than zero. ", + cache->get_base_path()); + } + return ""; + } + return validate_disk_capacity(cache->get_base_path(), new_capacity, valid_capacity); +} + std::string FileCacheFactory::reset_capacity(const std::string& path, int64_t new_capacity) { std::stringstream ss; size_t total_capacity = 0; if (path.empty()) { - for (auto& [p, cache] : _path_to_cache) { + std::vector> reset_targets; + reset_targets.reserve(_path_to_cache.size()); + for (auto& entry : _path_to_cache) { + auto* cache = entry.second; int64_t valid_capacity = 0; - ss << validate_capacity(p, new_capacity, valid_capacity); + ss << validate_capacity(cache, new_capacity, valid_capacity); if (valid_capacity <= 0) { return ss.str(); } + reset_targets.emplace_back(cache, cast_set(valid_capacity)); + } + for (auto& [cache, valid_capacity] : reset_targets) { ss << cache->reset_capacity(valid_capacity); total_capacity += cache->capacity(); } @@ -358,11 +378,11 @@ std::string FileCacheFactory::reset_capacity(const std::string& path, int64_t ne } else { if (auto iter = _path_to_cache.find(path); iter != _path_to_cache.end()) { int64_t valid_capacity = 0; - ss << validate_capacity(path, new_capacity, valid_capacity); + ss << validate_capacity(iter->second, new_capacity, valid_capacity); if (valid_capacity <= 0) { return ss.str(); } - ss << iter->second->reset_capacity(valid_capacity); + ss << iter->second->reset_capacity(cast_set(valid_capacity)); for (auto& [p, cache] : _path_to_cache) { total_capacity += cache->capacity(); diff --git a/be/src/io/cache/lru_queue_recorder.cpp b/be/src/io/cache/lru_queue_recorder.cpp index 314a444e232e9e..7ee7069d969ed0 100644 --- a/be/src/io/cache/lru_queue_recorder.cpp +++ b/be/src/io/cache/lru_queue_recorder.cpp @@ -35,6 +35,9 @@ size_t file_cache_type_index(FileCacheType type) { void LRUQueueRecorder::record_queue_event(FileCacheType type, CacheLRULogType log_type, const UInt128Wrapper hash, const size_t offset, const size_t size) { + if (_mgr->is_memory_storage()) { + return; + } if (config::file_cache_background_lru_dump_tail_record_num <= 0) { return; } diff --git a/be/test/io/cache/block_file_cache_mem_storage_mode_test.cpp b/be/test/io/cache/block_file_cache_mem_storage_mode_test.cpp new file mode 100644 index 00000000000000..67e83669878922 --- /dev/null +++ b/be/test/io/cache/block_file_cache_mem_storage_mode_test.cpp @@ -0,0 +1,465 @@ +// 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 "io/cache/block_file_cache_factory.h" +#include "io/cache/block_file_cache_test_common.h" + +namespace doris::io { + +class BlockFileCacheMemModeTest : public BlockFileCacheTest { +public: + void SetUp() override { + const auto* test_info = ::testing::UnitTest::GetInstance()->current_test_info(); + _cache_path = caches_dir / test_info->name() / ""; + fs::remove_all(_cache_path); + fs::create_directories(_cache_path); + + _origin_enter_need_evict = config::file_cache_enter_need_evict_cache_in_advance_percent; + _origin_exit_need_evict = config::file_cache_exit_need_evict_cache_in_advance_percent; + _origin_enter_disk_limit = config::file_cache_enter_disk_resource_limit_mode_percent; + _origin_exit_disk_limit = config::file_cache_exit_disk_resource_limit_mode_percent; + _origin_enable_evict_in_advance = config::enable_evict_file_cache_in_advance; + _origin_lru_dump_tail_record_num = config::file_cache_background_lru_dump_tail_record_num; + } + + void TearDown() override { + config::file_cache_enter_need_evict_cache_in_advance_percent = _origin_enter_need_evict; + config::file_cache_exit_need_evict_cache_in_advance_percent = _origin_exit_need_evict; + config::file_cache_enter_disk_resource_limit_mode_percent = _origin_enter_disk_limit; + config::file_cache_exit_disk_resource_limit_mode_percent = _origin_exit_disk_limit; + config::enable_evict_file_cache_in_advance = _origin_enable_evict_in_advance; + config::file_cache_background_lru_dump_tail_record_num = _origin_lru_dump_tail_record_num; + + auto* sync_point = SyncPoint::get_instance(); + sync_point->disable_processing(); + sync_point->clear_all_call_backs(); + + fs::remove_all(_cache_path); + } + +protected: + FileCacheSettings make_settings(const std::string& storage = "memory") const { + FileCacheSettings settings; + settings.storage = storage; + settings.capacity = 100_mb; + settings.max_file_block_size = 10_mb; + settings.max_query_cache_size = 100_mb; + settings.query_queue_size = 100_mb; + settings.query_queue_elements = 16; + settings.index_queue_size = 30_mb; + settings.index_queue_elements = 8; + settings.disposable_queue_size = 30_mb; + settings.disposable_queue_elements = 8; + settings.ttl_queue_size = 30_mb; + settings.ttl_queue_elements = 8; + return settings; + } + + void wait_async_open(BlockFileCache& cache) const { + for (int i = 0; i < 100; ++i) { + if (cache.get_async_open_success()) { + return; + } + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + } + FAIL() << "cache async open did not finish"; + } + + bool wait_until(const std::function& predicate, int64_t timeout_ms = 2000) const { + auto deadline = std::chrono::steady_clock::now() + std::chrono::milliseconds(timeout_ms); + while (std::chrono::steady_clock::now() < deadline) { + if (predicate()) { + return true; + } + std::this_thread::sleep_for(std::chrono::milliseconds(5)); + } + return predicate(); + } + + void fill_memory_cache(BlockFileCache& cache, size_t total_size) const { + TUniqueId query_id; + query_id.hi = 1; + query_id.lo = 1; + CacheContext context; + ReadStatistics rstats; + context.stats = &rstats; + context.cache_type = FileCacheType::NORMAL; + context.query_id = query_id; + auto key = BlockFileCache::hash("mem-mode-key"); + + for (size_t offset = 0; offset < total_size; offset += 10_mb) { + auto holder = cache.get_or_set(key, offset, 10_mb, context); + auto blocks = fromHolder(holder); + ASSERT_EQ(blocks.size(), 1); + ASSERT_TRUE(blocks[0]->get_or_set_downloader() == FileBlock::get_caller_id()); + download_into_memory(blocks[0]); + } + } + + void fill_cache(BlockFileCache& cache, const UInt128Wrapper& key, size_t total_size, + bool use_memory_download) const { + TUniqueId query_id; + query_id.hi = 11; + query_id.lo = 11; + CacheContext context; + ReadStatistics rstats; + context.stats = &rstats; + context.cache_type = FileCacheType::NORMAL; + context.query_id = query_id; + + for (size_t offset = 0; offset < total_size; offset += 10_mb) { + auto holder = cache.get_or_set(key, offset, 10_mb, context); + auto blocks = fromHolder(holder); + ASSERT_EQ(blocks.size(), 1); + ASSERT_TRUE(blocks[0]->get_or_set_downloader() == FileBlock::get_caller_id()); + if (use_memory_download) { + download_into_memory(blocks[0]); + } else { + download(blocks[0]); + } + } + } + + fs::path _cache_path; + int64_t _origin_enter_need_evict = 0; + int64_t _origin_exit_need_evict = 0; + int64_t _origin_enter_disk_limit = 0; + int64_t _origin_exit_disk_limit = 0; + bool _origin_enable_evict_in_advance = false; + int64_t _origin_lru_dump_tail_record_num = 0; +}; + +TEST_F(BlockFileCacheMemModeTest, check_need_evict_memory_enter_exit) { + auto settings = make_settings(); + BlockFileCache cache(_cache_path.string(), settings); + + config::file_cache_enter_need_evict_cache_in_advance_percent = 70; + config::file_cache_exit_need_evict_cache_in_advance_percent = 65; + + cache._cur_cache_size = 80_mb; + cache.check_need_evict_cache_in_advance(); + ASSERT_TRUE(cache._need_evict_cache_in_advance); + + cache._cur_cache_size = 60_mb; + cache.check_need_evict_cache_in_advance(); + ASSERT_FALSE(cache._need_evict_cache_in_advance); +} + +TEST_F(BlockFileCacheMemModeTest, zero_capacity_memory_should_keep_monitor_and_stats_available) { + auto settings = make_settings(); + settings.capacity = 0; + BlockFileCache cache(_cache_path.string(), settings); + + cache.check_disk_resource_limit(); + cache.check_need_evict_cache_in_advance(); + EXPECT_FALSE(cache._disk_resource_limit_mode.load(std::memory_order_relaxed)); + EXPECT_FALSE(cache._need_evict_cache_in_advance.load(std::memory_order_relaxed)); + EXPECT_EQ(cache.get_stats().at("size_percentage"), 0); + EXPECT_EQ(cache.get_stats_unsafe().at("size_percentage"), 0); + + ASSERT_TRUE(cache.initialize().ok()); + wait_async_open(cache); + EXPECT_EQ(cache.get_stats().at("size_percentage"), 0); +} + +TEST_F(BlockFileCacheMemModeTest, check_need_evict_memory_config_fallback) { + auto settings = make_settings(); + BlockFileCache cache(_cache_path.string(), settings); + + config::file_cache_enter_need_evict_cache_in_advance_percent = 70; + config::file_cache_exit_need_evict_cache_in_advance_percent = 75; + + cache._cur_cache_size = 80_mb; + cache.check_need_evict_cache_in_advance(); + + ASSERT_EQ(config::file_cache_enter_need_evict_cache_in_advance_percent, 78); + ASSERT_EQ(config::file_cache_exit_need_evict_cache_in_advance_percent, 75); + ASSERT_TRUE(cache._need_evict_cache_in_advance); +} + +TEST_F(BlockFileCacheMemModeTest, check_disk_mode_memory_enter_exit) { + auto settings = make_settings(); + BlockFileCache cache(_cache_path.string(), settings); + + config::file_cache_enter_disk_resource_limit_mode_percent = 90; + config::file_cache_exit_disk_resource_limit_mode_percent = 80; + + cache._cur_cache_size = 95_mb; + cache.check_disk_resource_limit(); + ASSERT_TRUE(cache._disk_resource_limit_mode); + + cache._cur_cache_size = 70_mb; + cache.check_disk_resource_limit(); + ASSERT_FALSE(cache._disk_resource_limit_mode); +} + +TEST_F(BlockFileCacheMemModeTest, check_disk_mode_disk_compat) { + auto settings = make_settings("disk"); + BlockFileCache cache(_cache_path.string(), settings); + + config::file_cache_enter_disk_resource_limit_mode_percent = 90; + config::file_cache_exit_disk_resource_limit_mode_percent = 80; + + auto* sync_point = SyncPoint::get_instance(); + sync_point->set_call_back("BlockFileCache::disk_used_percentage:1", + [](std::vector&& values) { + auto* percent = try_any_cast*>(values.back()); + percent->first = 95; + percent->second = 70; + }); + sync_point->enable_processing(); + cache.check_disk_resource_limit(); + ASSERT_TRUE(cache._disk_resource_limit_mode); + + sync_point->clear_all_call_backs(); + sync_point->set_call_back("BlockFileCache::disk_used_percentage:1", + [](std::vector&& values) { + auto* percent = try_any_cast*>(values.back()); + percent->first = 70; + percent->second = 70; + }); + cache.check_disk_resource_limit(); + ASSERT_FALSE(cache._disk_resource_limit_mode); +} + +TEST_F(BlockFileCacheMemModeTest, reset_capacity_memory_should_work_and_mode_converge) { + auto settings = make_settings(); + BlockFileCache cache(_cache_path.string(), settings); + ASSERT_TRUE(cache.initialize().ok()); + wait_async_open(cache); + + fill_memory_cache(cache, 60_mb); + ASSERT_EQ(cache._cur_cache_size, 60_mb); + + cache._disk_resource_limit_mode = false; + cache._disk_limit_mode_metrics->set_value(0); + + auto message = cache.reset_capacity(30_mb); + ASSERT_FALSE(message.empty()); + ASSERT_LE(cache._cur_cache_size, 30_mb); + ASSERT_FALSE(cache._disk_resource_limit_mode); + + config::file_cache_enter_disk_resource_limit_mode_percent = 90; + config::file_cache_exit_disk_resource_limit_mode_percent = 80; + cache.check_disk_resource_limit(); + ASSERT_TRUE(cache._disk_resource_limit_mode); +} + +TEST_F(BlockFileCacheMemModeTest, reset_capacity_memory_should_not_overcount_unreleasable_blocks) { + auto settings = make_settings(); + BlockFileCache cache(_cache_path.string(), settings); + ASSERT_TRUE(cache.initialize().ok()); + wait_async_open(cache); + + TUniqueId query_id; + query_id.hi = 7; + query_id.lo = 7; + CacheContext context; + ReadStatistics rstats; + context.stats = &rstats; + context.cache_type = FileCacheType::NORMAL; + context.query_id = query_id; + auto key = BlockFileCache::hash("mem-held-key"); + + auto holder = cache.get_or_set(key, 0, 10_mb, context); + auto blocks = fromHolder(holder); + ASSERT_EQ(blocks.size(), 1); + ASSERT_TRUE(blocks[0]->get_or_set_downloader() == FileBlock::get_caller_id()); + download_into_memory(blocks[0]); + + fill_memory_cache(cache, 30_mb); + ASSERT_EQ(cache._cur_cache_size, 40_mb); + + auto message = cache.reset_capacity(15_mb); + ASSERT_FALSE(message.empty()); + ASSERT_LE(cache._cur_cache_size, 15_mb); + ASSERT_EQ(cache._cur_cache_size, 10_mb); + + blocks.clear(); + holder.file_blocks.clear(); +} + +TEST_F(BlockFileCacheMemModeTest, run_background_monitor_memory_mode_flip) { + auto settings = make_settings(); + config::enable_evict_file_cache_in_advance = true; + config::file_cache_enter_need_evict_cache_in_advance_percent = 70; + config::file_cache_exit_need_evict_cache_in_advance_percent = 65; + config::file_cache_enter_disk_resource_limit_mode_percent = 90; + config::file_cache_exit_disk_resource_limit_mode_percent = 80; + + auto* sync_point = SyncPoint::get_instance(); + sync_point->set_call_back("BlockFileCache::set_sleep_time", + [](auto&& args) { *try_any_cast(args[0]) = 10; }); + sync_point->enable_processing(); + + BlockFileCache cache(_cache_path.string(), settings); + ASSERT_TRUE(cache.initialize().ok()); + wait_async_open(cache); + + { + std::lock_guard lock(cache._mutex); + cache._cur_cache_size = 95_mb; + } + ASSERT_TRUE(wait_until([&cache] { + return cache._need_evict_cache_in_advance && cache._disk_resource_limit_mode; + })); + + { + std::lock_guard lock(cache._mutex); + cache._cur_cache_size = 60_mb; + } + ASSERT_TRUE(wait_until([&cache] { + return !cache._need_evict_cache_in_advance && !cache._disk_resource_limit_mode; + })); +} + +TEST_F(BlockFileCacheMemModeTest, stats_memory_contains_storage_type_and_size_percentage) { + auto settings = make_settings(); + BlockFileCache cache(_cache_path.string(), settings); + + cache._cur_cache_size = 25_mb; + cache._cur_cache_size_metrics->set_value(cache._cur_cache_size); + + auto stats = cache.get_stats(); + ASSERT_TRUE(stats.contains("size_percentage")); + ASSERT_TRUE(stats.contains("storage_type")); + ASSERT_EQ(stats["size_percentage"], 25); + ASSERT_EQ(stats["storage_type"], FileCacheStorageType::MEMORY); +} + +TEST_F(BlockFileCacheMemModeTest, memory_storage_should_skip_lru_restore_and_dump) { + config::file_cache_background_lru_dump_tail_record_num = 100; + auto settings = make_settings(); + auto key = BlockFileCache::hash("memory-lru-restore-key"); + auto memory_dump_dir = fs::current_path() / "memory"; + fs::remove_all(memory_dump_dir); + fs::create_directories(memory_dump_dir); + + { + BlockFileCache cache(_cache_path.string(), settings); + ASSERT_TRUE(cache.initialize().ok()); + wait_async_open(cache); + + fill_cache(cache, key, 20_mb, true); + ASSERT_EQ(cache._cur_cache_size, 20_mb); + cache.dump_lru_queues(true); + } + + auto dump_file = memory_dump_dir / "lru_dump_normal.tail"; + ASSERT_FALSE(fs::exists(dump_file)); + + std::ofstream(dump_file, std::ios::binary) << "stale-memory-dump"; + ASSERT_TRUE(fs::exists(dump_file)); + + BlockFileCache cache2(_cache_path.string(), settings); + ASSERT_TRUE(cache2.initialize().ok()); + wait_async_open(cache2); + + ASSERT_EQ(cache2._cur_cache_size, 0); + ASSERT_EQ(cache2._normal_queue.get_elements_num_unsafe(), 0); + ASSERT_TRUE(cache2.get_blocks_by_key(key).empty()); + + fs::remove_all(memory_dump_dir); +} + +TEST_F(BlockFileCacheMemModeTest, memory_storage_should_skip_lru_queue_recorder) { + config::file_cache_background_lru_dump_tail_record_num = 100; + auto settings = make_settings(); + BlockFileCache cache(_cache_path.string(), settings); + ASSERT_TRUE(cache.initialize().ok()); + wait_async_open(cache); + + auto key = BlockFileCache::hash("memory-storage-skip-lru-recorder"); + fill_cache(cache, key, 10_mb, true); + + EXPECT_EQ(cache._lru_recorder->get_lru_queue_update_cnt_from_last_dump(FileCacheType::NORMAL), + 0); +} + +TEST_F(BlockFileCacheMemModeTest, disk_storage_should_keep_lru_restore) { + config::file_cache_background_lru_dump_tail_record_num = 100; + auto settings = make_settings("disk"); + auto key = BlockFileCache::hash("disk-lru-restore-key"); + auto origin_cache_base_path = cache_base_path; + cache_base_path = _cache_path.string(); + + { + BlockFileCache cache(_cache_path.string(), settings); + ASSERT_TRUE(cache.initialize().ok()); + wait_async_open(cache); + + fill_cache(cache, key, 10_mb, false); + cache.dump_lru_queues(true); + } + + BlockFileCache cache2(_cache_path.string(), settings); + ASSERT_TRUE(cache2.initialize().ok()); + wait_async_open(cache2); + + ASSERT_EQ(cache2._cur_cache_size, 10_mb); + ASSERT_EQ(cache2._normal_queue.get_elements_num_unsafe(), 1); + auto blocks = cache2.get_blocks_by_key(key); + ASSERT_EQ(blocks.size(), 1); + ASSERT_TRUE(blocks.contains(0)); + + cache_base_path = origin_cache_base_path; +} + +TEST_F(BlockFileCacheMemModeTest, factory_should_reset_memory_without_disk_validation) { + FileCacheFactory factory; + auto memory_settings = make_settings(); + auto disk_settings = make_settings("disk"); + auto disk_path = (_cache_path / "disk").string(); + + ASSERT_TRUE(factory.create_file_cache("memory", memory_settings).ok()); + ASSERT_TRUE(factory.create_file_cache(disk_path, disk_settings).ok()); + ASSERT_EQ(factory.get_capacity(), 200_mb); + + auto message = factory.reset_capacity("memory", 60_mb); + EXPECT_FALSE(message.empty()); + EXPECT_EQ(factory.get_by_path("memory")->capacity(), 60_mb); + EXPECT_EQ(factory.get_capacity(), 160_mb); + + message = factory.reset_capacity("", 40_mb); + EXPECT_FALSE(message.empty()); + EXPECT_EQ(factory.get_by_path("memory")->capacity(), 40_mb); + EXPECT_EQ(factory.get_by_path(disk_path)->capacity(), 40_mb); + EXPECT_EQ(factory.get_capacity(), 80_mb); +} + +TEST_F(BlockFileCacheMemModeTest, factory_should_prevalidate_all_resets_before_mutation) { + FileCacheFactory factory; + auto settings = make_settings("disk"); + auto disk_path_a = (_cache_path / "disk_a").string(); + auto disk_path_b = (_cache_path / "disk_b").string(); + + ASSERT_TRUE(factory.create_file_cache(disk_path_a, settings).ok()); + ASSERT_TRUE(factory.create_file_cache(disk_path_b, settings).ok()); + auto paths = factory.get_base_paths(); + ASSERT_EQ(paths.size(), 2); + fs::remove_all(paths.back()); + + auto message = factory.reset_capacity("", 60_mb); + EXPECT_NE(message.find("statfs error"), std::string::npos); + EXPECT_EQ(factory.get_by_path(disk_path_a)->capacity(), 100_mb); + EXPECT_EQ(factory.get_by_path(disk_path_b)->capacity(), 100_mb); + EXPECT_EQ(factory.get_capacity(), 200_mb); +} + +} // namespace doris::io diff --git a/be/test/service/http/file_cache_action_test.cpp b/be/test/service/http/file_cache_action_test.cpp index ce4f8b02372728..d98796b8bb2c4e 100644 --- a/be/test/service/http/file_cache_action_test.cpp +++ b/be/test/service/http/file_cache_action_test.cpp @@ -242,4 +242,31 @@ TEST_F(FileCacheActionTest, clear_value_uses_async_remove) { EXPECT_TRUE(_cache->get_blocks_by_key(hash).empty()); } +TEST_F(FileCacheActionTest, reset_memory_capacity_without_memory_directory) { + reset_file_cache_factory(); + std::filesystem::remove_all("memory"); + doris::io::FileCacheSettings settings; + settings.storage = "memory"; + settings.capacity = 1024 * 1024; + settings.max_file_block_size = 64 * 1024; + settings.query_queue_size = 1024 * 1024; + settings.query_queue_elements = 16; + auto* factory = doris::io::FileCacheFactory::instance(); + ASSERT_TRUE(factory->create_file_cache("memory", settings).ok()); + _cache = factory->get_by_path("memory"); + ASSERT_NE(_cache, nullptr); + + HttpRequest req(_evhttp_req); + std::string json_metrics; + req._params["op"] = "reset"; + req._params["path"] = "memory"; + req._params["capacity"] = std::to_string(2 * 1024 * 1024); + + Status status = _action->_handle_header(&req, &json_metrics); + + EXPECT_TRUE(status.ok()); + EXPECT_EQ(_cache->capacity(), 2 * 1024 * 1024); + EXPECT_EQ(factory->get_capacity(), 2 * 1024 * 1024); +} + } // namespace doris