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

Tokio 0.1出错时如何将Future重新加入Stream以持续轮询

在Tokio 0.1中实现无限重试的多监听轮询循环

要实现无论任何异常都持续发起新请求的无限轮询,核心是让Stream在出错时不终止,而是生成新的请求Future继续轮询。以下是基于Tokio 0.1和Futures 0.1的具体实现方案:

核心思路

利用stream::unfold或stream::repeat构建无限Stream,在每个请求的then回调中捕获成功/失败状态,无论结果如何都返回下一个待轮询的Future,彻底避免因错误导致Stream退出。

代码实现

1. 定义模拟请求Future

先封装一个可能失败的请求逻辑(替换为你实际的HTTP请求,比如基于Hyper的实现):

use futures::{future, stream, Future, Stream};
use tokio::runtime::Runtime;
use rand; // 仅用于模拟随机失败

// 模拟可能失败的HTTP请求
fn make_request() -> impl Future<Item = String, Error = ()> {
    future::result(
        if rand::random() {
            Ok("请求成功: 获取到数据".to_string())
        } else {
            Err(()) // 模拟网络错误、超时等异常
        }
    )
}

2. 构建无限重试Stream

使用stream::unfold创建持续生成请求的Stream,出错时自动重试:

fn infinite_poll_stream() -> impl Stream<Item = String, Error = !> {
    stream::unfold((), |_| {
        make_request()
            .then(|result| {
                // 处理请求结果:成功则打印,失败则输出重试提示
                match result {
                    Ok(response) => println!("{}", response),
                    Err(_) => println!("请求失败,立即重试..."),
                }
                // 无论成功/失败,都返回下一个要轮询的任务
                Ok(Some(((), ())))
            })
    })
    // 将错误类型转为Never(!),因为我们不会让Stream因错误终止
    .map_err(|_| unreachable!())
}

3. 添加重试延迟(可选)

如果需要在出错后延迟重试(避免频繁请求),可以结合tokio::timer::Delay:

use tokio::timer::Delay;
use std::time::{Duration, Instant};

fn make_request_with_delay() -> impl Future<Item = (), Error = ()> {
    make_request()
        .then(|result| {
            match result {
                Ok(res) => {
                    println!("{}", res);
                    future::ok(())
                }
                Err(_) => {
                    println!("请求失败,1秒后重试...");
                    // 延迟1秒后继续
                    Delay::new(Instant::now() + Duration::from_secs(1))
                        .map_err(|_| ()) // 忽略Delay自身的错误
                }
            }
        })
}

fn infinite_poll_with_delay() -> impl Stream<Item = (), Error = !> {
    stream::unfold((), |_| {
        make_request_with_delay()
            .then(|_| Ok(Some(((), ()))))
    })
    .map_err(|_| unreachable!())
}

4. 运行Runtime

在Tokio Runtime中启动无限轮询:

fn main() {
    let mut runtime = Runtime::new().expect("Failed to create runtime");
    // 运行无限Stream,for_each忽略每个元素仅持续轮询
    runtime.block_on(infinite_poll_stream().for_each(|_| Ok(()))).unwrap();
}

关键细节

  • stream::unfold的作用:它通过初始状态和闭包持续生成新的Future,只要闭包返回Some((元素, 新状态)),Stream就会继续轮询。这里我们用空状态(),只关注持续生成请求。
  • then回调的异常捕获:无论请求成功还是失败,then都会返回一个新的Future,确保Stream不会因错误终止。
  • Never类型!:通过map_err(|_| unreachable!())将错误转为Never,因为我们的逻辑已经处理了所有错误场景,不会让错误传递到Stream外层导致终止。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 20:37:35