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

Tokio/Futures 0.1中长轮询阻塞时如何调度其他Future?

问题描述

项目必须使用Tokio/Futures 0.1实现多Consul服务监听,要求当某个Future被长轮询请求阻塞时,其他同类型Future仍能正常轮询。当前实现存在两个核心问题:

  1. reqwest::get是阻塞调用,会卡住整个Runtime,导致其他Future无法被调度
  2. poll并非async函数,无法直接使用await处理异步请求

给定的Cargo.toml配置:

[dependencies]
tokio = "0.1"
futures= "0.1"

测试用例要求每5秒循环输出:

this is a test name1
this is a test name2
this is a test name1
this is a test name2

解决方案

1. 调整依赖,引入异步HTTP客户端

替换阻塞的reqwest为Tokio 0.1兼容的异步版本,同时启用Tokio的定时器特性:

[dependencies]
tokio = { version = "0.1", features = ["timer"] }
futures = "0.1"
reqwest = { version = "0.10", features = ["tokio01"] }

2. 重构Watcher结构,用状态机管理异步任务

不能在poll中执行阻塞操作,需将HTTP请求和延迟任务作为内部异步Future持有,通过状态机跟踪任务进度:

use std::time::{Duration, Instant};
use futures::{Async, Future, Poll};
use reqwest::Client;
use tokio::timer::Delay;

// 状态机定义,跟踪当前异步任务阶段
enum WatcherState {
    WaitingDelay(Delay),          // 等待间隔时间结束
    WaitingRequest(Box<dyn Future<Item = (), Error = ()>>), // 等待HTTP请求完成
    ReadyToLoop,                  // 准备进入下一轮循环
}

pub struct JnsWatcher {
    name: String,
    state: WatcherState,
    client: Client,
    consul_url: String,
    interval: Duration,
}

impl JnsWatcher {
    pub fn new(name: String, consul_url: String, interval: Duration) -> Self {
        JnsWatcher {
            name,
            state: WatcherState::ReadyToLoop,
            client: Client::new(),
            consul_url,
            interval,
        }
    }

    // 创建Consul长轮询请求的异步Future
    fn create_consul_request(&self) -> impl Future<Item = (), Error = ()> {
        self.client
            .get(&self.consul_url)
            .send()
            .map_err(|e| eprintln!("{}请求失败: {}", self.name, e))
            .and_then(|_| Ok(()))
    }
}

impl Future for JnsWatcher {
    type Item = ();
    type Error = ();

    fn poll(&mut self) -> Poll<Self::Item, Self::Error> {
        loop {
            self.state = match self.state {
                WatcherState::ReadyToLoop => {
                    // 初始化延迟任务,等待指定间隔后发起请求
                    let delay = Delay::new(Instant::now() + self.interval);
                    WatcherState::WaitingDelay(delay)
                }
                WatcherState::WaitingDelay(ref mut delay) => {
                    match delay.poll() {
                        Ok(Async::Ready(())) => {
                            // 延迟结束,创建Consul请求Future
                            let req_future = self.create_consul_request();
                            WatcherState::WaitingRequest(Box::new(req_future))
                        }
                        Ok(Async::NotReady) => return Ok(Async::NotReady),
                        Err(e) => {
                            eprintln!("{}延迟任务失败: {}", self.name, e);
                            // 错误处理:重新创建延迟任务
                            WatcherState::WaitingDelay(Delay::new(Instant::now() + self.interval))
                        }
                    }
                }
                WatcherState::WaitingRequest(ref mut req_future) => {
                    match req_future.poll() {
                        Ok(Async::Ready(())) => {
                            // 请求完成,输出内容
                            println!("this is a test {}", self.name);
                            // 准备进入下一轮循环
                            WatcherState::ReadyToLoop
                        }
                        Ok(Async::NotReady) => return Ok(Async::NotReady),
                        Err(_) => {
                            // 请求失败,直接进入下一轮循环
                            WatcherState::ReadyToLoop
                        }
                    }
                }
            };
        }
    }
}

3. 调整测试代码,用Tokio Runtime运行多Watcher

使用Tokio Runtime调度多个异步Watcher,确保任务并发执行:

use futures::future::join_all;
use tokio::runtime::Runtime;

#[test]
fn test_pool() {
    // 创建Tokio Runtime
    let mut rt = Runtime::new().expect("Runtime创建失败");

    // 初始化两个Watcher,间隔5秒
    let w1 = JnsWatcher::new(
        "name1".to_string(),
        "http://consul.example.com/v1/health/service/service1?wait=5s".to_string(),
        Duration::from_secs(5),
    );
    let w2 = JnsWatcher::new(
        "name2".to_string(),
        "http://consul.example.com/v1/health/service/service2?wait=5s".to_string(),
        Duration::from_secs(5),
    );

    // 运行所有Watcher(无限循环,测试时可手动终止)
    rt.block_on_all(join_all(vec![w1, w2])).unwrap();

    println!("test end");
}

核心原理说明

  • 非阻塞IO:所有网络请求使用Tokio兼容的异步客户端,避免阻塞Runtime线程
  • 状态机调度:通过状态机跟踪每个Watcher的任务阶段,每次poll仅处理当前任务的进度,不会阻塞其他Future
  • 并发执行:Tokio Runtime会自动调度多个异步任务,即使某个Watcher处于长轮询等待状态,其他Watcher仍能正常被轮询执行

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 13:37:56