Rust基于tokio DelayQueue实现缓存时如何正确调用poll_purge方法?
问题根因
poll_purge方法要求传入的&mut Context<'_>是Rust异步运行时轮询Future时由运行时传入的上下文参数,普通业务代码无法手动构造或获取该参数,也不应该直接调用poll_开头的底层轮询方法——这类方法是为实现Future/Stream trait设计的底层API,不是给上层业务逻辑直接调用的。
推荐实现方案
你不需要手动维护poll_purge的轮询逻辑,DelayQueue本身已经实现了Stream trait,可以直接通过异步流的方式拉取过期key,配合tokio的后台任务自动完成清理,全程不需要手动接触Context参数。
另外你原有代码的insert方法存在逻辑缺陷:重复插入同一个key时旧的过期记录不会被移除,会导致同一个key被多次加入过期队列,触发重复/错误删除,下面的代码已经修复该问题,同时补全了跨任务共享缓存所需的线程安全包装:
# Cargo.toml 依赖保持不变 [dependencies] futures = "0.3" tokio-util = { version = "0.7.3", features = ["full"] } tokio = { version = "1.19.2", features = ["full"] }
use futures::StreamExt; use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; use tokio::sync::Mutex; use tokio_util::time::{delay_queue, DelayQueue}; type CacheKey = String; type Value = String; struct Cache { entries: HashMap<CacheKey, (Value, delay_queue::Key)>, expirations: DelayQueue<CacheKey>, } const TTL_SECS: u64 = 10; impl Cache { fn new() -> Cache { Cache { entries: HashMap::new(), expirations: DelayQueue::new(), } } fn insert(&mut self, key: CacheKey, value: Value) { // 重复插入时先移除旧的过期记录 if let Some((_, old_delay_key)) = self.entries.remove(&key) { self.expirations.remove(&old_delay_key); } let delay_key = self .expirations .insert(key.clone(), Duration::from_secs(TTL_SECS)); self.entries.insert(key, (value, delay_key)); } fn get(&self, key: &CacheKey) -> Option<&Value> { self.entries.get(key).map(|(v, _)| v) } fn remove(&mut self, key: &CacheKey) { if let Some((_, delay_key)) = self.entries.remove(key) { self.expirations.remove(&delay_key); } } } #[tokio::main] async fn main() { // 跨异步任务共享缓存需要用Arc<Mutex>做线程安全包装 let cache = Arc::new(Mutex::new(Cache::new())); // 启动后台常驻清理任务 let purge_cache = cache.clone(); tokio::spawn(async move { let mut cache_guard = purge_cache.lock().await; let mut expirations = std::mem::replace(&mut cache_guard.expirations, DelayQueue::new()); drop(cache_guard); loop { match expirations.next().await { Some(Ok(expired_entry)) => { let mut cache_guard = purge_cache.lock().await; // 删除过期的缓存条目 cache_guard.entries.remove(expired_entry.get_ref()); // 把运行过程中新插入的过期任务合并到当前流 let new_expirations = std::mem::replace(&mut cache_guard.expirations, DelayQueue::new()); expirations.extend(new_expirations); drop(cache_guard); } // 队列异常时直接退出清理任务 _ => break, } } }); // 正常业务操作 let mut cache_guard = cache.lock().await; cache_guard.insert("k1".to_string(), "v1".to_string()); cache_guard.insert("k2".to_string(), "v2".to_string()); cache_guard.insert("k3".to_string(), "v3".to_string()); cache_guard.insert("k4".to_string(), "v4".to_string()); cache_guard.insert("k5".to_string(), "v5".to_string()); drop(cache_guard); // 模拟等待缓存过期 tokio::time::sleep(Duration::from_secs(12)).await; let cache_guard = cache.lock().await; assert!(cache_guard.get(&"k1".to_string()).is_none()); println!("过期缓存已自动清理完成"); }
注意事项
- 所有
poll_前缀的方法都是异步生态的底层API,上层业务开发优先使用对应的async方法、Stream/Future组合子,不要直接调用poll方法避免出错。 - 如果缓存的读写并发量很高,可以把
tokio::sync::Mutex换成parking_lot::Mutex进一步降低锁开销,逻辑不需要改动。
内容的提问来源于stack exchange,提问作者Shoobi
相关产品推荐
相关产品推荐

