This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 461c6d2331b branch-4.1: [fix](filecache) convert blocks back to TTL
when an expired TTL is extended #67971 (#68092)
461c6d2331b is described below
commit 461c6d2331bf57c246b33a5a77a059b7b1901be7
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Thu Sep 17 21:56:26 2026 +0800
branch-4.1: [fix](filecache) convert blocks back to TTL when an expired TTL
is extended #67971 (#68092)
Cherry-picked from #67971
Co-authored-by: Xin Liao <[email protected]>
---
be/src/io/cache/block_file_cache.h | 4 +
be/src/io/cache/block_file_cache_ttl_mgr.cpp | 184 ++++++++----
be/src/io/cache/block_file_cache_ttl_mgr.h | 45 ++-
be/test/io/cache/block_file_cache_ttl_mgr_test.cpp | 308 ++++++++++++++++++++-
4 files changed, 478 insertions(+), 63 deletions(-)
diff --git a/be/src/io/cache/block_file_cache.h
b/be/src/io/cache/block_file_cache.h
index 191dec86d98..21dae1c0452 100644
--- a/be/src/io/cache/block_file_cache.h
+++ b/be/src/io/cache/block_file_cache.h
@@ -309,6 +309,10 @@ public:
[[nodiscard]] bool get_async_open_success() const { return
_async_open_done; }
+ // The manager that keeps cached blocks in the cache type their tablet's
TTL asks for.
+ // Exposed so that tests can drive it deterministically rather than race
its threads.
+ BlockFileCacheTtlMgr* get_ttl_mgr() { return _ttl_mgr.get(); }
+
BlockFileCache& operator=(const BlockFileCache&) = delete;
BlockFileCache(const BlockFileCache&) = delete;
diff --git a/be/src/io/cache/block_file_cache_ttl_mgr.cpp
b/be/src/io/cache/block_file_cache_ttl_mgr.cpp
index b394591d42d..924526709b8 100644
--- a/be/src/io/cache/block_file_cache_ttl_mgr.cpp
+++ b/be/src/io/cache/block_file_cache_ttl_mgr.cpp
@@ -41,6 +41,8 @@ BlockFileCacheTtlMgr::BlockFileCacheTtlMgr(BlockFileCache*
mgr, CacheBlockMetaSt
: _mgr(mgr), _meta_store(meta_store), _stop_background(false) {
_tablet_id_set_size_metrics = std::make_shared<bvar::Status<size_t>>(
_mgr->get_base_path().c_str(),
"file_cache_ttl_mgr_tablet_id_set_size", 0);
+ _ttl_info_map_size_metrics = std::make_shared<bvar::Status<size_t>>(
+ _mgr->get_base_path().c_str(),
"file_cache_ttl_mgr_ttl_info_map_size", 0);
resume();
}
@@ -85,6 +87,17 @@ void BlockFileCacheTtlMgr::register_tablet_id(int64_t
tablet_id) {
_tablet_id_queue.enqueue(tablet_id);
}
+void BlockFileCacheTtlMgr::update_ttl_info_map_size_metrics() {
+ if (_ttl_info_map_size_metrics) {
+ _ttl_info_map_size_metrics->set_value(_ttl_info_map.size());
+ }
+}
+
+size_t BlockFileCacheTtlMgr::tracked_tablet_num() {
+ std::lock_guard<std::mutex> lock(_ttl_info_mutex);
+ return _ttl_info_map.size();
+}
+
void BlockFileCacheTtlMgr::run_background_tablet_id_flush() {
Thread::set_self_name("ttl_mgr_flush");
@@ -168,6 +181,83 @@ FileBlocks
BlockFileCacheTtlMgr::get_file_blocks_from_tablet_id(int64_t tablet_i
return result;
}
+void BlockFileCacheTtlMgr::reconcile_tablet_blocks(int64_t tablet_id) {
+ // Serialize all conversions of this tablet. Whichever caller takes this
lock last re-reads
+ // the state below and has the final say, so the update and expiration
threads cannot fight
+ // over the same blocks and strand them in the loser's cache type.
+ std::lock_guard<std::mutex>
transition_lock(transition_lock_for(tablet_id));
+
+ bool want_ttl = false;
+ bool blocks_promoted = false;
+ {
+ // Deliberately re-read the map rather than trust what the caller saw:
the expiration
+ // thread picks its candidates up to a full gc interval before getting
here.
+ std::lock_guard<std::mutex> lock(_ttl_info_mutex);
+ auto it = _ttl_info_map.find(tablet_id);
+ if (it != _ttl_info_map.end()) {
+ if (it->second.ttl == 0 && !it->second.blocks_promoted) {
+ // No TTL, and none of its blocks were put in the TTL queue by
us: nothing left
+ // to track. Dropped here rather than after a conversion, so
that a tablet which
+ // settles without needing one still stops being walked on
every round.
+ _ttl_info_map.erase(it);
+ update_ttl_info_map_size_metrics();
+ return;
+ }
+ want_ttl = it->second.is_ttl_active(UnixSeconds());
+ blocks_promoted = it->second.blocks_promoted;
+ }
+ }
+
+ // The two directions are not symmetric.
+ //
+ // Promotion is edge triggered: once the blocks are in the TTL queue
nothing takes them out
+ // behind our back, so rescanning a tablet whose TTL is still running is
pure waste. This is
+ // also what makes a TTL rewritten to another still-valid value free,
which matters where
+ // the property is rewritten on a schedule.
+ //
+ // Demotion is level triggered: blocks can still land in the TTL queue
after a tablet was
+ // demoted, and what is recorded here is per tablet, so it cannot tell
whether any have.
+ // Rescanning is the only way to collect them.
+ if (want_ttl && blocks_promoted) {
+ return;
+ }
+
+ // Scan and convert outside _ttl_info_mutex: this walks the meta store and
takes the cache
+ // lock once per block, which is far too long to hold a mutex the other
thread needs.
+ const auto target_type = want_ttl ? FileCacheType::TTL :
FileCacheType::NORMAL;
+ FileBlocks blocks = get_file_blocks_from_tablet_id(tablet_id);
+ size_t converted = 0;
+ bool all_converted = true;
+ for (auto& block : blocks) {
+ if (block->cache_type() == target_type) {
+ continue;
+ }
+ auto st = block->change_cache_type(target_type);
+ if (st.ok()) {
+ ++converted;
+ } else {
+ all_converted = false;
+ LOG(WARNING) << "Failed to convert block to " <<
cache_type_to_string(target_type)
+ << " cache_type, tablet_id=" << tablet_id << ", err="
<< st;
+ }
+ }
+ if (converted > 0) {
+ LOG(INFO) << "converted cached blocks to " <<
cache_type_to_string(target_type)
+ << ", tablet_id=" << tablet_id << ", block_num=" << converted
+ << ", scanned=" << blocks.size();
+ }
+
+ {
+ std::lock_guard<std::mutex> lock(_ttl_info_mutex);
+ auto it = _ttl_info_map.find(tablet_id);
+ if (it != _ttl_info_map.end() && all_converted) {
+ // Left alone when a block failed to convert, so the mismatch
stays visible to the
+ // next round and gets retried instead of being recorded as done.
+ it->second.blocks_promoted = want_ttl;
+ }
+ }
+}
+
void BlockFileCacheTtlMgr::run_backgroud_update_ttl_info_map() {
Thread::set_self_name("ttl_mgr_update");
@@ -203,7 +293,9 @@ void
BlockFileCacheTtlMgr::run_backgroud_update_ttl_info_map() {
}
{
std::lock_guard<std::mutex> lock(_ttl_info_mutex);
- _ttl_info_map.erase(tablet_id);
+ if (_ttl_info_map.erase(tablet_id) > 0) {
+ update_ttl_info_map_size_metrics();
+ }
}
} else {
LOG(WARNING) << "Failed to get tablet meta for
tablet_id: " << tablet_id
@@ -220,47 +312,36 @@ void
BlockFileCacheTtlMgr::run_backgroud_update_ttl_info_map() {
}
}
- // Update TTL info map
- bool need_convert_from_ttl = false;
+ // Record the TTL this tablet currently has, then let
reconcile_tablet_blocks()
+ // decide whether that moves its blocks between the TTL and
normal queues.
+ bool tracked = false;
{
std::lock_guard<std::mutex> lock(_ttl_info_mutex);
+ auto it = _ttl_info_map.find(tablet_id);
if (ttl > 0) {
- auto old_info_it = _ttl_info_map.find(tablet_id);
- bool was_zero_ttl = (old_info_it ==
_ttl_info_map.end() ||
- old_info_it->second.ttl == 0);
- _ttl_info_map[tablet_id] = TtlInfo {ttl, tablet_ctime};
-
- // If TTL changed from 0 to non-zero, convert blocks
to TTL type
- if (was_zero_ttl) {
- FileBlocks blocks =
get_file_blocks_from_tablet_id(tablet_id);
- for (auto& block : blocks) {
- if (block->cache_type() != FileCacheType::TTL)
{
- auto change_status =
-
block->change_cache_type(FileCacheType::TTL);
- if (!change_status.ok()) {
- LOG(WARNING) << "Failed to convert
block to TTL cache_type";
- }
- }
- }
+ if (it == _ttl_info_map.end()) {
+ _ttl_info_map.emplace(tablet_id, TtlInfo {ttl,
tablet_ctime});
+ update_ttl_info_map_size_metrics();
+ } else {
+ // Keep blocks_promoted: it records what we did to
the blocks, not
+ // what the tablet meta says.
+ it->second.ttl = ttl;
+ it->second.tablet_ctime = tablet_ctime;
}
- } else {
- // Periodically reconcile blocks restored from
persisted TTL metadata,
- // because _ttl_info_map is rebuilt only in memory
after restart.
- need_convert_from_ttl =
- _ttl_info_map.erase(tablet_id) > 0 ||
need_full_reconcile;
+ tracked = true;
+ } else if (it != _ttl_info_map.end()) {
+ // Hold on to the entry until the blocks are actually
demoted; it
+ // drops itself once the tablet has settled back to
NORMAL.
+ it->second.ttl = 0;
+ tracked = true;
}
}
- if (need_convert_from_ttl) {
- FileBlocks blocks =
get_file_blocks_from_tablet_id(tablet_id);
- for (auto& block : blocks) {
- if (block->cache_type() == FileCacheType::TTL) {
- auto st =
block->change_cache_type(FileCacheType::NORMAL);
- if (!st.ok()) {
- LOG(WARNING) << "Failed to convert block back
to NORMAL cache_type";
- }
- }
- }
+ // An untracked tablet reconciles to NORMAL, which is how TTL
blocks restored
+ // from persisted metadata are cleaned up after a restart.
Gated on the periodic
+ // round so that ordinary non-TTL tablets are not walked every
time.
+ if (tracked || need_full_reconcile) {
+ reconcile_tablet_blocks(tablet_id);
}
}
@@ -279,29 +360,26 @@ void
BlockFileCacheTtlMgr::run_backgroud_expiration_check() {
while (!_stop_background.load(std::memory_order_acquire)) {
try {
- std::map<int64_t, TtlInfo> ttl_info_copy;
-
- // Copy TTL info for processing
+ // Collect tablets whose TTL has run out.
+ std::vector<int64_t> expired_tablet_ids;
{
std::lock_guard<std::mutex> lock(_ttl_info_mutex);
- ttl_info_copy = _ttl_info_map;
+ uint64_t current_time = UnixSeconds();
+ for (const auto& [tablet_id, ttl_info] : _ttl_info_map) {
+ if (ttl_info.ttl > 0 &&
!ttl_info.is_ttl_active(current_time)) {
+ expired_tablet_ids.push_back(tablet_id);
+ }
+ }
}
- uint64_t current_time = UnixSeconds();
-
- for (const auto& [tablet_id, ttl_info] : ttl_info_copy) {
- if (ttl_info.tablet_ctime + ttl_info.ttl < current_time) {
- // Tablet has expired, convert TTL blocks back to NORMAL
type
- FileBlocks blocks =
get_file_blocks_from_tablet_id(tablet_id);
- for (auto& block : blocks) {
- if (block->cache_type() == FileCacheType::TTL) {
- auto st =
block->change_cache_type(FileCacheType::NORMAL);
- if (!st.ok()) {
- LOG(WARNING) << "Failed to convert block back
to NORMAL cache_type";
- }
- }
- }
+ // Only a candidate list: reconcile_tablet_blocks() re-reads the
tablet's state
+ // under the per-tablet lock, so a TTL extended in between is
never demoted on the
+ // strength of what was observed here.
+ for (int64_t tablet_id : expired_tablet_ids) {
+ if (_stop_background.load(std::memory_order_acquire)) {
+ break;
}
+ reconcile_tablet_blocks(tablet_id);
}
std::this_thread::sleep_for(
diff --git a/be/src/io/cache/block_file_cache_ttl_mgr.h
b/be/src/io/cache/block_file_cache_ttl_mgr.h
index 8cd677446f1..59e2233f284 100644
--- a/be/src/io/cache/block_file_cache_ttl_mgr.h
+++ b/be/src/io/cache/block_file_cache_ttl_mgr.h
@@ -22,7 +22,9 @@
#include <bvar/bvar.h>
#include <concurrentqueue.h>
+#include <array>
#include <atomic>
+#include <limits>
#include <map>
#include <memory>
#include <mutex>
@@ -38,8 +40,25 @@ class BlockFileCache;
class CacheBlockMetaStore;
struct TtlInfo {
- uint64_t ttl;
- uint64_t tablet_ctime;
+ uint64_t ttl = 0;
+ uint64_t tablet_ctime = 0;
+ // True once this manager has put the tablet's blocks into the TTL queue
and nothing has
+ // taken them out since. Note it records what we did, not what the blocks
are: it starts
+ // false for a tablet seen for the first time after a restart, whose
blocks on disk may
+ // well be TTL already, because this process has no record of putting them
there and must
+ // scan to find out.
+ bool blocks_promoted = false;
+
+ // Whether this tablet's blocks belong in the TTL queue right now.
+ bool is_ttl_active(uint64_t now) const {
+ if (ttl == 0 || tablet_ctime == 0) {
+ return false;
+ }
+ if (tablet_ctime > std::numeric_limits<uint64_t>::max() - ttl) {
+ return false;
+ }
+ return tablet_ctime + ttl >= now;
+ }
};
class BlockFileCacheTtlMgr {
@@ -48,6 +67,9 @@ public:
~BlockFileCacheTtlMgr();
void register_tablet_id(int64_t tablet_id);
+ // Number of tablets whose TTL state is currently tracked. Mirrors the
+ // file_cache_ttl_mgr_ttl_info_map_size bvar; entry leaks are otherwise
invisible.
+ size_t tracked_tablet_num();
void stop();
void resume();
@@ -61,6 +83,19 @@ public:
private:
FileBlocks get_file_blocks_from_tablet_id(int64_t tablet_id);
+ // Drive this tablet's cached blocks to the cache type its current TTL
state asks for.
+ // Both background threads funnel through here and it is serialized per
tablet, so they can
+ // never scan the same tablet concurrently and leave the blocks in
whatever type the scan
+ // that happened to finish last wrote.
+ void reconcile_tablet_blocks(int64_t tablet_id);
+
+ // Caller must hold _ttl_info_mutex.
+ void update_ttl_info_map_size_metrics();
+
+ std::mutex& transition_lock_for(int64_t tablet_id) {
+ return _transition_locks[static_cast<uint64_t>(tablet_id) %
kTransitionLockStripes];
+ }
+
private:
// Tablet ids waiting to be deduplicated + set of unique ids known to have
cached data
moodycamel::ConcurrentQueue<int64_t> _tablet_id_queue;
@@ -79,7 +114,13 @@ private:
std::mutex _ttl_info_mutex;
+ // Striped locks serializing block conversions per tablet. Lock order is
always
+ // _transition_locks[i] -> _ttl_info_mutex; _ttl_info_mutex is never held
across a block scan.
+ static constexpr size_t kTransitionLockStripes = 64;
+ std::array<std::mutex, kTransitionLockStripes> _transition_locks;
+
std::shared_ptr<bvar::Status<size_t>> _tablet_id_set_size_metrics;
+ std::shared_ptr<bvar::Status<size_t>> _ttl_info_map_size_metrics;
};
} // namespace doris::io
diff --git a/be/test/io/cache/block_file_cache_ttl_mgr_test.cpp
b/be/test/io/cache/block_file_cache_ttl_mgr_test.cpp
index 56f6d3d43d6..dc38ba48a09 100644
--- a/be/test/io/cache/block_file_cache_ttl_mgr_test.cpp
+++ b/be/test/io/cache/block_file_cache_ttl_mgr_test.cpp
@@ -223,6 +223,14 @@ protected:
ASSERT_TRUE(_cache->initialize());
ASSERT_TRUE(wait_for_condition([this]() { return
_cache->get_async_open_success(); },
std::chrono::seconds(5)));
+
+ // initialize() starts a TTL manager of its own against this cache.
Every case below
+ // drives one it owns, and two of them converting the same blocks
makes both the
+ // conversions and the scan counts nondeterministic -- a block can be
demoted before
+ // the case has finished setting up the state it means to exercise.
+ if (auto* cache_owned_ttl_mgr = _cache->get_ttl_mgr()) {
+ cache_owned_ttl_mgr->stop();
+ }
}
void TearDown() override {
@@ -368,16 +376,18 @@ TEST_F(BlockFileCacheTtlMgrTest,
NonTtlTabletWithoutPriorTtlInfoSkipsBlockScan)
auto block = create_block(kTabletId, "non-ttl-tablet", 0, 1024, &hash);
persist_block_meta(kTabletId, hash, block->range().left,
block->range().size());
- std::atomic<int64_t> block_scan_count {0};
+ // Held by value in the callback: the background threads outlive this
stack frame, and
+ // neither the guard nor disable_processing() synchronizes with a callback
in flight.
+ auto block_scan_count = std::make_shared<std::atomic<int64_t>>(0);
auto* sync_point = SyncPoint::get_instance();
sync_point->clear_all_call_backs();
sync_point->clear_trace();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"BlockFileCacheTtlMgr::get_file_blocks_from_tablet_id",
- [&block_scan_count](std::vector<std::any>&& args) {
+ [block_scan_count](std::vector<std::any>&& args) {
if (doris::try_any_cast<int64_t>(args[0]) == kTabletId) {
- block_scan_count.fetch_add(1, std::memory_order_relaxed);
+ block_scan_count->fetch_add(1, std::memory_order_relaxed);
}
},
&guard);
@@ -390,11 +400,13 @@ TEST_F(BlockFileCacheTtlMgrTest,
NonTtlTabletWithoutPriorTtlInfoSkipsBlockScan)
[this]() { return fake_engine()->get_tablet_meta_call_count() >=
2; },
std::chrono::seconds(5));
std::this_thread::sleep_for(std::chrono::milliseconds(100));
+ // Join the background threads before the callback and its captures go
away.
+ _ttl_mgr.reset();
sync_point->disable_processing();
sync_point->clear_trace();
EXPECT_TRUE(update_thread_observed);
- EXPECT_EQ(0, block_scan_count.load(std::memory_order_relaxed));
+ EXPECT_EQ(0, block_scan_count->load(std::memory_order_relaxed));
EXPECT_EQ(FileCacheType::NORMAL, block->cache_type());
}
@@ -411,16 +423,18 @@ TEST_F(BlockFileCacheTtlMgrTest,
PeriodicReconcileDemotesTtlBlockWithoutPriorTtl
FileCacheType::TTL, expiration_time);
ASSERT_EQ(FileCacheType::TTL, block->cache_type());
- std::atomic<int64_t> block_scan_count {0};
+ // Held by value in the callback: the background threads outlive this
stack frame, and
+ // neither the guard nor disable_processing() synchronizes with a callback
in flight.
+ auto block_scan_count = std::make_shared<std::atomic<int64_t>>(0);
auto* sync_point = SyncPoint::get_instance();
sync_point->clear_all_call_backs();
sync_point->clear_trace();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"BlockFileCacheTtlMgr::get_file_blocks_from_tablet_id",
- [&block_scan_count](std::vector<std::any>&& args) {
+ [block_scan_count](std::vector<std::any>&& args) {
if (doris::try_any_cast<int64_t>(args[0]) == kTabletId) {
- block_scan_count.fetch_add(1, std::memory_order_relaxed);
+ block_scan_count->fetch_add(1, std::memory_order_relaxed);
}
},
&guard);
@@ -432,11 +446,13 @@ TEST_F(BlockFileCacheTtlMgrTest,
PeriodicReconcileDemotesTtlBlockWithoutPriorTtl
bool demoted =
wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::NORMAL; },
std::chrono::seconds(5));
+ // Join the background threads before the callback and its captures go
away.
+ _ttl_mgr.reset();
sync_point->disable_processing();
sync_point->clear_trace();
EXPECT_TRUE(demoted);
- EXPECT_GE(block_scan_count.load(std::memory_order_relaxed), 1);
+ EXPECT_GE(block_scan_count->load(std::memory_order_relaxed), 1);
}
TEST_F(BlockFileCacheTtlMgrTest, TabletTtlRemovedMovesBlocksBackToNormal) {
@@ -461,4 +477,280 @@ TEST_F(BlockFileCacheTtlMgrTest,
TabletTtlRemovedMovesBlocksBackToNormal) {
std::chrono::seconds(5)));
}
+TEST_F(BlockFileCacheTtlMgrTest, ExpiredTtlExtendedMovesBlocksBackToTtl) {
+ constexpr int64_t kTabletId = 6006;
+ auto tablet = std::make_shared<FakeTablet>(UnixSeconds(), 120);
+ fake_engine()->add_tablet(kTabletId, tablet);
+
+ UInt128Wrapper hash;
+ auto block = create_block(kTabletId, "ttl-extend-after-expire", 0, 1024,
&hash);
+ persist_block_meta(kTabletId, hash, block->range().left,
block->range().size());
+
+ _ttl_mgr = std::make_unique<BlockFileCacheTtlMgr>(_cache.get(),
_meta_store.get());
+ _ttl_mgr->register_tablet_id(kTabletId);
+
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::TTL; },
+ std::chrono::seconds(5)));
+
+ // Let the TTL expire. The block goes back to NORMAL while the manager
keeps a non-zero TTL
+ // recorded for the tablet, which is the state that used to wedge the
promotion path.
+ tablet->set_creation_time(UnixSeconds() - 120);
+ tablet->set_ttl_seconds(1);
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::NORMAL; },
+ std::chrono::seconds(5)));
+
+ // Extending an already expired TTL to one that has not expired has to
bring the block back.
+ tablet->set_ttl_seconds(30758400);
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::TTL; },
+ std::chrono::seconds(5)));
+}
+
+TEST_F(BlockFileCacheTtlMgrTest,
ExtendedTtlThatIsStillExpiredKeepsBlocksNormal) {
+ constexpr int64_t kTabletId = 7007;
+ const int64_t creation_time = UnixSeconds() - 7200;
+ auto tablet = std::make_shared<FakeTablet>(creation_time, 60);
+ fake_engine()->add_tablet(kTabletId, tablet);
+
+ UInt128Wrapper hash;
+ auto block = create_block(kTabletId, "ttl-extend-still-expired", 0, 1024,
&hash);
+ persist_block_meta(kTabletId, hash, block->range().left,
block->range().size());
+
+ _ttl_mgr = std::make_unique<BlockFileCacheTtlMgr>(_cache.get(),
_meta_store.get());
+ _ttl_mgr->register_tablet_id(kTabletId);
+
+ int64_t call_count = fake_engine()->get_tablet_meta_call_count();
+ ASSERT_TRUE(wait_for_condition(
+ [&]() { return fake_engine()->get_tablet_meta_call_count() >=
call_count + 2; },
+ std::chrono::seconds(5)));
+ ASSERT_EQ(FileCacheType::NORMAL, block->cache_type());
+
+ // A longer TTL that is still in the past must not promote anything.
+ tablet->set_ttl_seconds(120);
+ call_count = fake_engine()->get_tablet_meta_call_count();
+ ASSERT_TRUE(wait_for_condition(
+ [&]() { return fake_engine()->get_tablet_meta_call_count() >=
call_count + 3; },
+ std::chrono::seconds(5)));
+ EXPECT_EQ(FileCacheType::NORMAL, block->cache_type());
+}
+
+TEST_F(BlockFileCacheTtlMgrTest,
RewritingTtlToAnotherValidValueDoesNotRescanBlocks) {
+ constexpr int64_t kTabletId = 8008;
+ auto tablet = std::make_shared<FakeTablet>(UnixSeconds(), 3600);
+ fake_engine()->add_tablet(kTabletId, tablet);
+
+ UInt128Wrapper hash;
+ auto block = create_block(kTabletId, "ttl-rewrite-valid", 0, 1024, &hash);
+ persist_block_meta(kTabletId, hash, block->range().left,
block->range().size());
+
+ _ttl_mgr = std::make_unique<BlockFileCacheTtlMgr>(_cache.get(),
_meta_store.get());
+ _ttl_mgr->register_tablet_id(kTabletId);
+
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::TTL; },
+ std::chrono::seconds(5)));
+
+ // The promotion is recorded only after every block has been converted, so
the manager has
+ // not necessarily finished with the tablet at the moment the type flips.
Let a couple of
+ // rounds pass before counting, or the tail of the promotion is charged to
the rewrites.
+ int64_t settled_after = fake_engine()->get_tablet_meta_call_count();
+ ASSERT_TRUE(wait_for_condition(
+ [&]() { return fake_engine()->get_tablet_meta_call_count() >=
settled_after + 2; },
+ std::chrono::seconds(5)));
+
+ // Held by value in the callback: the background threads outlive this
stack frame, and
+ // neither the guard nor disable_processing() synchronizes with a callback
in flight.
+ auto block_scan_count = std::make_shared<std::atomic<int64_t>>(0);
+ auto* sync_point = SyncPoint::get_instance();
+ sync_point->clear_all_call_backs();
+ sync_point->clear_trace();
+ SyncPoint::CallbackGuard guard;
+ sync_point->set_call_back(
+ "BlockFileCacheTtlMgr::get_file_blocks_from_tablet_id",
+ [block_scan_count](std::vector<std::any>&& args) {
+ if (doris::try_any_cast<int64_t>(args[0]) == kTabletId) {
+ block_scan_count->fetch_add(1, std::memory_order_relaxed);
+ }
+ },
+ &guard);
+ sync_point->enable_processing();
+
+ // Automated jobs rewrite this property regularly. As long as the tablet
stays in the same
+ // state, none of those rewrites may trigger another walk of the meta
store.
+ for (int64_t ttl : {7200, 1800, 5400}) {
+ tablet->set_ttl_seconds(ttl);
+ int64_t call_count = fake_engine()->get_tablet_meta_call_count();
+ ASSERT_TRUE(wait_for_condition(
+ [&]() { return fake_engine()->get_tablet_meta_call_count() >=
call_count + 2; },
+ std::chrono::seconds(5)));
+ }
+
+ // Join the background threads before the callback and its captures go
away.
+ _ttl_mgr.reset();
+ sync_point->disable_processing();
+ sync_point->clear_trace();
+
+ EXPECT_EQ(0, block_scan_count->load(std::memory_order_relaxed));
+ EXPECT_EQ(FileCacheType::TTL, block->cache_type());
+}
+
+TEST_F(BlockFileCacheTtlMgrTest, TtlExtensionWinsOverConcurrentExpirationScan)
{
+ constexpr int64_t kTabletId = 9009;
+ auto tablet = std::make_shared<FakeTablet>(UnixSeconds(), 3600);
+ fake_engine()->add_tablet(kTabletId, tablet);
+
+ UInt128Wrapper hash;
+ auto block = create_block(kTabletId, "ttl-extend-during-demote", 0, 1024,
&hash);
+ persist_block_meta(kTabletId, hash, block->range().left,
block->range().size());
+
+ _ttl_mgr = std::make_unique<BlockFileCacheTtlMgr>(_cache.get(),
_meta_store.get());
+ _ttl_mgr->register_tablet_id(kTabletId);
+
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::TTL; },
+ std::chrono::seconds(5)));
+
+ // Stall the demotion scan midway so the TTL can be extended underneath
it, reproducing the
+ // window where the expiration check acts on a view of the tablet that is
already stale.
+ // Held by value in the callback rather than captured by reference: the
background threads
+ // outlive this stack frame, and neither the guard nor
disable_processing() synchronizes
+ // with a callback already in flight. A callback that stalls makes that
window wide.
+ auto demote_scan_entered = std::make_shared<std::atomic<bool>>(false);
+ auto release_demote_scan = std::make_shared<std::atomic<bool>>(false);
+ auto* sync_point = SyncPoint::get_instance();
+ sync_point->clear_all_call_backs();
+ sync_point->clear_trace();
+ SyncPoint::CallbackGuard guard;
+ sync_point->set_call_back(
+ "BlockFileCacheTtlMgr::get_file_blocks_from_tablet_id",
+ [demote_scan_entered, release_demote_scan](std::vector<std::any>&&
args) {
+ if (doris::try_any_cast<int64_t>(args[0]) != kTabletId) {
+ return;
+ }
+ if (demote_scan_entered->exchange(true,
std::memory_order_acq_rel)) {
+ return;
+ }
+ while (!release_demote_scan->load(std::memory_order_acquire)) {
+ std::this_thread::sleep_for(std::chrono::milliseconds(5));
+ }
+ },
+ &guard);
+ sync_point->enable_processing();
+
+ tablet->set_creation_time(UnixSeconds() - 3600);
+ tablet->set_ttl_seconds(1);
+
+ bool scan_stalled = wait_for_condition(
+ [&]() { return
demote_scan_entered->load(std::memory_order_acquire); },
+ std::chrono::seconds(10));
+
+ // Extend the TTL while the demotion is still in flight.
+ tablet->set_creation_time(UnixSeconds());
+ tablet->set_ttl_seconds(30758400);
+ std::this_thread::sleep_for(std::chrono::milliseconds(200));
+ // Must come before the join below, or stop() waits on a thread parked in
the callback.
+ release_demote_scan->store(true, std::memory_order_release);
+
+ bool ends_as_ttl = wait_for_condition(
+ [&]() { return block->cache_type() == FileCacheType::TTL; },
std::chrono::seconds(10));
+ _ttl_mgr.reset();
+ sync_point->disable_processing();
+ sync_point->clear_trace();
+
+ ASSERT_TRUE(scan_stalled);
+ EXPECT_TRUE(ends_as_ttl);
+}
+
+// _ttl_info_map is rebuilt in memory only, while the cache type of each block
survives on disk.
+// A tablet whose TTL expired while this BE was down therefore comes back with
TTL blocks and no
+// recorded state, and nothing later in the tablet's life re-examines them.
+TEST_F(BlockFileCacheTtlMgrTest,
ExpiredTabletDemotesTtlBlocksRestoredFromDisk) {
+ constexpr int64_t kTabletId = 10010;
+ const int64_t creation_time = UnixSeconds() - 7200;
+ auto tablet = std::make_shared<FakeTablet>(creation_time, 60);
+ fake_engine()->add_tablet(kTabletId, tablet);
+
+ const uint64_t expiration_time = static_cast<uint64_t>(creation_time) + 60;
+ UInt128Wrapper hash;
+ auto block = create_block(kTabletId, "ttl-restored-expired", 0, 1024,
&hash, FileCacheType::TTL,
+ expiration_time);
+ // Asserted before the block is published to the meta store: a scan can
only reach a block
+ // that is listed there, so until then no manager can convert it out from
under us.
+ ASSERT_EQ(FileCacheType::TTL, block->cache_type());
+ persist_block_meta(kTabletId, hash, block->range().left,
block->range().size(),
+ FileCacheType::TTL, expiration_time);
+
+ _ttl_mgr = std::make_unique<BlockFileCacheTtlMgr>(_cache.get(),
_meta_store.get());
+ _ttl_mgr->register_tablet_id(kTabletId);
+
+ EXPECT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::NORMAL; },
+ std::chrono::seconds(5)));
+}
+
+// Blocks can still land in the TTL queue after a tablet's existing ones were
demoted, and the
+// recorded state is per tablet, so it cannot tell that they have. They have
to be collected too.
+TEST_F(BlockFileCacheTtlMgrTest,
ExpiredTabletDemotesTtlBlocksCachedAfterDemotion) {
+ constexpr int64_t kTabletId = 11011;
+ auto tablet = std::make_shared<FakeTablet>(UnixSeconds(), 120);
+ fake_engine()->add_tablet(kTabletId, tablet);
+
+ UInt128Wrapper hash;
+ auto block = create_block(kTabletId, "ttl-expire-then-cache", 0, 1024,
&hash);
+ persist_block_meta(kTabletId, hash, block->range().left,
block->range().size());
+
+ _ttl_mgr = std::make_unique<BlockFileCacheTtlMgr>(_cache.get(),
_meta_store.get());
+ _ttl_mgr->register_tablet_id(kTabletId);
+
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::TTL; },
+ std::chrono::seconds(5)));
+
+ tablet->set_creation_time(UnixSeconds() - 120);
+ tablet->set_ttl_seconds(1);
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::NORMAL; },
+ std::chrono::seconds(5)));
+
+ // A fresh block lands in the TTL queue after the demotion.
+ const uint64_t expiration_time = UnixSeconds() + 3600;
+ UInt128Wrapper late_hash;
+ auto late_block = create_block(kTabletId, "ttl-expire-then-cache-late", 0,
1024, &late_hash,
+ FileCacheType::TTL, expiration_time);
+ // Same ordering as above, and here it matters: the manager is running by
this point and
+ // sweeps this tablet every round, so publishing first would race the
assertion.
+ ASSERT_EQ(FileCacheType::TTL, late_block->cache_type());
+ persist_block_meta(kTabletId, late_hash, late_block->range().left,
late_block->range().size(),
+ FileCacheType::TTL, expiration_time);
+
+ EXPECT_TRUE(
+ wait_for_condition([&]() { return late_block->cache_type() ==
FileCacheType::NORMAL; },
+ std::chrono::seconds(5)));
+}
+
+// Dropping the TTL of a tablet that had already expired leaves nothing to
convert, but the
+// tablet still has to stop being tracked -- otherwise it holds a map entry
for the life of the
+// process and never again qualifies for the periodic reconcile.
+TEST_F(BlockFileCacheTtlMgrTest, TtlClearedAfterExpiryStopsTrackingTablet) {
+ constexpr int64_t kTabletId = 12012;
+ auto tablet = std::make_shared<FakeTablet>(UnixSeconds(), 120);
+ fake_engine()->add_tablet(kTabletId, tablet);
+
+ UInt128Wrapper hash;
+ auto block = create_block(kTabletId, "ttl-cleared-after-expiry", 0, 1024,
&hash);
+ persist_block_meta(kTabletId, hash, block->range().left,
block->range().size());
+
+ _ttl_mgr = std::make_unique<BlockFileCacheTtlMgr>(_cache.get(),
_meta_store.get());
+ _ttl_mgr->register_tablet_id(kTabletId);
+
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::TTL; },
+ std::chrono::seconds(5)));
+ ASSERT_TRUE(wait_for_condition([&]() { return
_ttl_mgr->tracked_tablet_num() == 1; },
+ std::chrono::seconds(5)));
+
+ tablet->set_creation_time(UnixSeconds() - 120);
+ tablet->set_ttl_seconds(1);
+ ASSERT_TRUE(wait_for_condition([&]() { return block->cache_type() ==
FileCacheType::NORMAL; },
+ std::chrono::seconds(5)));
+
+ tablet->set_ttl_seconds(0);
+ EXPECT_TRUE(wait_for_condition([&]() { return
_ttl_mgr->tracked_tablet_num() == 0; },
+ std::chrono::seconds(5)));
+ EXPECT_EQ(FileCacheType::NORMAL, block->cache_type());
+}
+
} // namespace doris::io
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]