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

如何在Rust Tokio中实现可取消Future?负载测试全局取消方案

实现Tokio负载测试器的全局取消逻辑

你要的全局取消逻辑在Tokio生态里有标准实现方案,最常用的是tokio::sync::CancellationToken——这是Tokio官方专门为任务取消场景设计的工具,完全适配你这种需要同时取消数千个客户端的场景。

核心实现思路

  1. 在主函数中创建全局的CancellationToken实例,给每个通过tokio::spawn()启动的客户端任务克隆一份令牌副本。
  2. 主函数监听退出信号(比如Ctrl+C),触发时调用令牌的cancel()方法,向所有持有令牌副本的任务广播取消信号。
  3. 客户端任务中,用tokio::select!同时等待业务逻辑(stream.message())和取消信号(token.cancelled()),一旦取消信号触发,立即执行套接字断开逻辑并退出循环。

完整代码示例

use tokio::sync::CancellationToken;
use tokio::signal::ctrl_c;
// 替换为你的实际Stream类型定义
use your_crate::Stream;

async fn connect_to_server() -> Result<Stream, Box<dyn std::error::Error>> {
    // 替换为你的实际连接逻辑
    Ok(Stream::connect("127.0.0.1:8080").await?)
}

#[tokio::main]
async fn main() {
    // 创建全局取消令牌
    let cancel_token = CancellationToken::new();
    let signal_token = cancel_token.clone();

    // 启动信号监听任务,捕获Ctrl+C后触发全局取消
    tokio::spawn(async move {
        if ctrl_c().await.is_ok() {
            println!("收到退出信号,开始全局取消所有客户端...");
            signal_token.cancel();
        }
    });

    // 启动1000个客户端任务
    for client_id in 0..1000 {
        let token = cancel_token.clone();
        tokio::spawn(async move {
            println!("启动客户端 {}", client_id);
            let mut stream = match connect_to_server().await {
                Ok(s) => s,
                Err(e) => {
                    eprintln!("客户端 {} 连接失败: {}", client_id, e);
                    return;
                }
            };

            loop {
                tokio::select! {
                    // 等待服务器消息
                    msg_result = stream.message() => {
                        match msg_result {
                            Ok(msg) => {
                                // 替换为你的消息处理逻辑
                                println!("客户端 {} 收到消息: {:?}", client_id, msg);
                            }
                            Err(e) => {
                                eprintln!("客户端 {} 消息接收失败: {}", client_id, e);
                                break;
                            }
                        }
                    }
                    // 等待取消信号
                    _ = token.cancelled() => {
                        println!("客户端 {} 收到取消信号,断开连接", client_id);
                        // 替换为你的实际断开逻辑
                        if let Err(e) = stream.disconnect().await {
                            eprintln!("客户端 {} 断开失败: {}", client_id, e);
                        }
                        break;
                    }
                }
            }
            println!("客户端 {} 退出", client_id);
        });
    }

    // 等待所有客户端任务完成后退出主进程
    cancel_token.cancelled().await;
    println!("所有客户端已退出,程序结束");
}

为什么选CancellationToken?

  • 轻量高效:克隆令牌的开销极低,完全适配数千个客户端的大规模场景。
  • 语义清晰:专门为取消场景设计,cancelled()返回的Future会在令牌被取消时立即完成,无需额外处理消息传递细节。
  • 灵活扩展:支持创建子令牌实现层级取消,满足更复杂的分组取消需求。

替代方案:Broadcast通道

如果不想用CancellationToken,也可以用tokio::sync::broadcast通道实现类似逻辑:主函数创建广播发送者,每个客户端任务订阅接收者,信号触发时发送取消消息,客户端任务通过select!监听消息和业务逻辑。但这种方案需要额外处理消息收发的样板代码,优先级低于CancellationToken。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 23:46:30