You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 15:27:15