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

如何从调用方法中取消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:

  1. 确保Tokio依赖包含sync特性:tokio = { version = "1.x", features = ["full"] }
  2. 在execution_cycle中创建CancellationToken,并克隆传递给所有任务
  3. 任务在循环中主动检查取消状态,或通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:10:27