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

如何获取两个async函数结果?StreamCache实现与测试问题

解答

1. 获取异步函数结果与优化实现

问题根源

join!宏会等待所有传入的Future全部完成才返回,但api.subscribe()返回的是一个无限流(测试代码中最终是future::pending().await,永远不会结束),导致join!永久阻塞,无法继续执行。

正确实现方式

不需要等待流结束,而是要分别启动后台任务处理两个逻辑:

  • 处理fetch():获取全量数据更新缓存,可定期重复调用保证数据一致性
  • 处理subscribe():持续消费流中的增量数据,实时更新缓存

修改后的new方法:

use tokio::time::Duration;

impl StreamCache {
    pub async fn new(api: impl Api + Clone) -> Self {
        let instance = Self {
            results: Arc::new(Mutex::new(HashMap::new())),
        };

        // 处理subscribe实时流
        let results_clone = instance.results.clone();
        let api_clone = api.clone();
        tokio::spawn(async move {
            let mut stream = api_clone.subscribe().await;
            while let Some(result) = stream.next().await {
                match result {
                    Ok((city, temp)) => {
                        let mut results = results_clone.lock().expect("poisoned mutex");
                        results.insert(city, temp);
                    }
                    Err(e) => eprintln!("Subscribe error: {}", e),
                }
            }
        });

        // 处理初始fetch与定期全量更新
        let results_clone_2 = instance.results.clone();
        tokio::spawn(async move {
            // 初始全量拉取
            if let Ok(full_data) = api.fetch().await {
                let mut results = results_clone_2.lock().expect("poisoned mutex");
                results.extend(full_data);
            }

            // 可选:定期刷新全量数据,比如每5分钟一次
            loop {
                tokio::time::sleep(Duration::from_minutes(5)).await;
                if let Ok(full_data) = api.fetch().await {
                    let mut results = results_clone_2.lock().expect("poisoned mutex");
                    results.extend(full_data);
                }
            }
        });

        instance
    }

    // ... 其他方法保持不变
}

注意:需要为Api trait添加Clone约束,确保能在多个后台任务中复用实例;同时依赖tokio runtime来启动后台任务,避免阻塞主线程。

2. 让单元测试生效

原测试问题

原测试中StreamCache::new因join!阻塞无法返回,且即使修复后,缓存更新是异步的,需要等待数据写入后再断言。

修改后的测试代码

#[tokio::test]
async fn works() {
    let cache = StreamCache::new(TestApi::default()).await;

    // 等待异步任务完成缓存更新
    tokio::time::sleep(Duration::from_millis(100)).await;

    assert_eq!(cache.get("Berlin"), Some(29));
    assert_eq!(cache.get("London"), Some(27));
    // subscribe的Paris值32会覆盖fetch的31,最终取最新值
    assert_eq!(cache.get("Paris"), Some(32));
}

同时需确保Cargo.toml中包含正确的依赖:

[dependencies]
async-trait = "0.1"
futures = "0.3"
tokio = { version = "1.0", features = ["full"] }
maplit = "1.0"

3. 实现缓存的后台更新

后台更新的核心逻辑已包含在上述new方法中,拆解为三个部分:

  1. 初始全量更新:启动时调用api.fetch(),将全量数据写入缓存
  2. 实时增量更新:通过api.subscribe()的流,接收单个城市的温度变化并实时更新缓存
  3. 定期全量刷新(可选):定期调用api.fetch()覆盖缓存,避免因流丢失数据导致的不一致

若需要单独触发后台更新,可完善update_in_background方法:

pub fn update_in_background(&self, api: impl Api + Clone) {
    let results_clone = self.results.clone();
    tokio::spawn(async move {
        loop {
            match api.fetch().await {
                Ok(full_data) => {
                    let mut results = results_clone.lock().expect("poisoned mutex");
                    results.extend(full_data);
                }
                Err(e) => eprintln!("Background fetch error: {}", e),
            }
            tokio::time::sleep(Duration::from_minutes(5)).await;
        }
    });
}

该方法可在new中调用,或根据业务需求在外部触发。


内容的提问来源于stack exchange,提问作者R Sun

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 16:49:52