Tokio/Futures 0.1中长轮询阻塞时如何调度其他Future?
问题描述
项目必须使用Tokio/Futures 0.1实现多Consul服务监听,要求当某个Future被长轮询请求阻塞时,其他同类型Future仍能正常轮询。当前实现存在两个核心问题:
reqwest::get是阻塞调用,会卡住整个Runtime,导致其他Future无法被调度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
相关产品推荐
相关产品推荐

