如何从调用方法中取消Tokio spawn创建的任务线程
解决Tokio任务无法通过abort取消的问题
你的代码里任务无法被abort取消的核心原因是使用了阻塞线程的同步sleep:std::thread::sleep会直接占用Tokio工作线程,让任务无法回到Tokio调度器,导致abort发出的取消信号无法被处理。要解决这个问题,需要做以下关键修改:
核心修改点
- 把所有
std::thread::sleep替换为Tokio提供的异步sleep:tokio::time::sleep(Duration::from_millis(...)).await,这样任务在等待时会主动让出线程,Tokio能及时处理取消请求。 - 利用Tokio异步sleep的取消特性:当任务被
abort时,异步sleep会返回Err(Elapsed),我们可以通过这个信号终止任务循环。
修改后的完整代码
use std::sync::mpsc::channel; use std::sync::mpsc::{Sender, Receiver, TryRecvError}; use tokio::time::{Duration, sleep}; // 改用Tokio的异步sleep use rand::distributions::{Uniform, Distribution}; #[tokio::main] async fn main() { execution_cycle().await; } async fn execution_cycle() { let (tx_first, rx_first) = channel::<Message>(); let (tx_second, rx_second) = channel::<Message>(); let handle_sender_first = tokio::spawn(sender_thread(tx_first)); let handle_sender_second = tokio::spawn(sender_thread(tx_second)); let handle_receiver = tokio::spawn(receiver_thread(rx_first, rx_second)); let mut thread_rng = rand::thread_rng(); let rng_generator = Uniform::from(1..10); loop { let cancel_from_cycle = rng_generator.sample(&mut thread_rng); if cancel_from_cycle == 9 { println!("Aborting from the execution cycle."); // 取消所有任务 handle_receiver.abort(); handle_sender_first.abort(); handle_sender_second.abort(); break; } // 避免循环占用过多资源,添加短等待 sleep(Duration::from_millis(500)).await; } // 等待任务彻底结束,确认取消状态 let _ = handle_sender_first.await; let _ = handle_sender_second.await; let _ = handle_receiver.await; println!("handle_sender_first finished: {}", handle_sender_first.is_finished()); println!("handle_sender_second finished: {}", handle_sender_second.is_finished()); println!("handle_receiver finished: {}", handle_receiver.is_finished()); } async fn sender_thread(tx: Sender<Message>) { let mut thread_rng = rand::thread_rng(); let rng_generator = Uniform::from(1..20); let mut random_id = rng_generator.sample(&mut thread_rng); while random_id != 9 { let msg = Message { id: random_id, text: "hello".to_owned() }; println!("Sending message {}.", msg.id); random_id = rng_generator.sample(&mut thread_rng); println!("Generated id {}.", random_id); if let Err(error) = tx.send(msg) { println!("Sending error {:?}", error); break; } // 使用Tokio异步sleep,响应取消信号 if let Err(_) = sleep(Duration::from_millis(2000)).await { println!("Sender task aborted."); break; } } } async fn receiver_thread(rx_first: Receiver<Message>, rx_second: Receiver<Message>) { let mut channel_open_first = true; let mut channel_open_second = true; let mut thread_rng = rand::thread_rng(); let rng_generator = Uniform::from(1..15); let mut random_event = rng_generator.sample(&mut thread_rng); while channel_open_first && channel_open_second && random_event != 9 { channel_open_first = receiver_inner(&rx_first); channel_open_second = receiver_inner(&rx_second); random_event = rng_generator.sample(&mut thread_rng); println!("Generated event {}.", random_event); // 使用Tokio异步sleep,响应取消信号 if let Err(_) = sleep(Duration::from_millis(800)).await { println!("Receiver task aborted."); break; } } } fn receiver_inner(rx: &Receiver<Message>) -> bool { let value = rx.try_recv(); match value { Ok(msg) => { println!("Message {} received: {}", msg.id, msg.text); }, Err(error) => { if error != TryRecvError::Empty { println!("{}", error); return false; } } } true } struct Message { id: usize, text: String, }
额外优化建议:用CancellationToken实现可控取消
如果需要更优雅的取消逻辑(比如让任务在取消前执行清理工作),可以使用Tokio的CancellationToken:
- 确保Tokio依赖包含
sync特性:tokio = { version = "1.x", features = ["full"] } - 在
execution_cycle中创建CancellationToken,并克隆传递给所有任务 - 任务在循环中主动检查取消状态,或通过
select!同时等待业务逻辑和取消信号
示例修改片段:
// 在execution_cycle中 use tokio::sync::CancellationToken; let token = CancellationToken::new(); let cloned_token1 = token.clone(); let cloned_token2 = token.clone(); let cloned_token3 = token.clone(); let handle_sender_first = tokio::spawn(async move { sender_thread(tx_first, cloned_token1).await }); // 其他任务同理传递克隆的token // 触发取消时 token.cancel(); // 在sender_thread中 async fn sender_thread(tx: Sender<Message>, token: CancellationToken) { let mut thread_rng = rand::thread_rng(); let rng_generator = Uniform::from(1..20); let mut random_id = rng_generator.sample(&mut thread_rng); while random_id != 9 && !token.is_cancelled() { // ...原有发送逻辑 // 同时等待sleep和取消信号 tokio::select! { _ = sleep(Duration::from_millis(2000)) => {}, _ = token.cancelled() => { println!("Sender task cancelled, cleaning up..."); break; } } } }
内容的提问来源于stack exchange,提问作者Ron Lukas
相关产品推荐
相关产品推荐

