如何在Rust多任务间共享并重置tokio::time::Sleep异步定时器
多任务共享异步定时器阻塞问题分析与解决
我正在开发一个能在多个异步任务间共享的Rust异步定时器,实现代码如下:
use std::{ future::{self, Future}, pin::Pin, sync::{Arc, Mutex}, time::Duration, }; use tokio::time::{self, Instant, Sleep}; struct Foo(Mutex<Pin<Box<Sleep>>>); impl Foo { fn new(sleep: Sleep) -> Self { Self(Mutex::new(Box::pin(sleep))) } async fn sleep(&self) { future::poll_fn(|cx| self.0.lock().unwrap().as_mut().poll(cx)).await } fn reset(&self, deadline: Instant) { self.0.lock().unwrap().as_mut().reset(deadline); } } async fn task1(foo: Arc<Foo>) { println!("starting task 1 ..."); let start = Instant::now(); foo.sleep().await; let time = start.elapsed().as_millis(); println!("task 1 complete in {time} millis "); } async fn task2(foo: Arc<Foo>) { println!("starting task 2 ..."); let start = Instant::now(); foo.sleep().await; let time = start.elapsed().as_millis(); println!("task 2 complete in {time} millis "); } #[tokio::main] pub async fn main() { let sleep = time::sleep(Duration::from_secs(3)); let foo = Arc::new(Foo::new(sleep)); let task1 = tokio::spawn(task1(foo.clone())); let task2 = tokio::spawn(task2(foo)); tokio::join!(task1, task2); }
运行后输出:
starting task 2 ... starting task 1 ... task 1 complete in 3005 millis // stuck here
问题:只有一个任务完成,另一个任务陷入阻塞。是不是因为第二个被轮询的future的waker覆盖了第一个?我知道FutureExt::shared方法,但它会获取future的所有权,而我需要在其他任务等待定时器时仍能重置定时器的能力。
问题原因分析
你的猜测完全正确,核心问题就是**Sleep future的waker被覆盖**。
Tokio的Sleep内部仅维护一个waker实例,当多个任务通过互斥锁轮询同一个Sleep时,后一个任务的waker会直接覆盖前一个的。当定时器到期时,只有最后一次设置的waker会被触发,其他任务的waker根本没机会收到唤醒信号,导致这些任务永久阻塞。
另外你的实现还有个隐藏问题:Mutex的持有时间。poll_fn里调用lock().unwrap()会持有互斥锁直到poll调用完成,虽然Sleep::poll本身不会阻塞,但多个任务同时抢锁会造成不必要的等待。不过这里导致任务阻塞的主要原因还是waker被覆盖。
解决方案:维护多任务waker列表
要实现多任务共享且支持重置的定时器,需要自己维护所有等待任务的waker,而不是依赖Sleep单一的waker。改进后的实现如下:
use std::{ future::{self, Future}, pin::Pin, sync::{Arc, Mutex}, task::{Context, Waker}, time::Duration, }; use tokio::time::{self, Instant, Sleep}; struct SharedTimer { inner: Mutex<SharedTimerInner>, } struct SharedTimerInner { sleep: Pin<Box<Sleep>>, waiters: Vec<Waker>, } impl SharedTimer { fn new(duration: Duration) -> Self { let sleep = Box::pin(time::sleep(duration)); Self { inner: Mutex::new(SharedTimerInner { sleep, waiters: Vec::new(), }), } } async fn sleep(&self) { future::poll_fn(|cx| { let mut inner = self.inner.lock().unwrap(); match inner.sleep.as_mut().poll(cx) { std::task::Poll::Ready(()) => { // 定时器到期,唤醒所有等待任务 for waker in inner.waiters.drain(..) { waker.wake(); } std::task::Poll::Ready(()) } std::task::Poll::Pending => { // 避免重复添加同一任务的waker if !inner.waiters.iter().any(|w| w.will_wake(cx.waker())) { inner.waiters.push(cx.waker().clone()); } std::task::Poll::Pending } } }).await } fn reset(&self, deadline: Instant) { let mut inner = self.inner.lock().unwrap(); inner.sleep.as_mut().reset(deadline); // 重置后唤醒所有等待任务,让它们重新轮询定时器状态 for waker in inner.waiters.drain(..) { waker.wake(); } } } async fn task1(timer: Arc<SharedTimer>) { println!("starting task 1 ..."); let start = Instant::now(); timer.sleep().await; let time = start.elapsed().as_millis(); println!("task 1 complete in {time} millis "); } async fn task2(timer: Arc<SharedTimer>) { println!("starting task 2 ..."); let start = Instant::now(); timer.sleep().await; let time = start.elapsed().as_millis(); println!("task 2 complete in {time} millis "); } #[tokio::main] pub async fn main() { let timer = Arc::new(SharedTimer::new(Duration::from_secs(3))); let task1 = tokio::spawn(task1(timer.clone())); let task2 = tokio::spawn(task2(timer)); tokio::join!(task1, task2); }
关键改进点
- 维护waker列表:在内部结构体中新增
waiters向量,保存所有等待任务的waker,彻底避免覆盖问题。 - 到期唤醒所有任务:当定时器到期时,遍历唤醒所有等待的任务,确保每个任务都能收到完成信号。
- 重置主动唤醒:调用
reset后立即唤醒所有等待任务,让它们重新检查定时器状态,保证重置操作能立即生效。 - 去重waker:添加waker前检查是否已存在,避免同一任务的waker被重复添加,减少不必要的唤醒操作。
修改后两个任务都会在定时器到期后正常完成,同时保留了随时重置定时器的能力。
内容的提问来源于stack exchange,提问作者EAimTY
相关产品推荐
相关产品推荐

