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

Rust技术问询:基于websocket-rs、tokio和SSL的异步WebSocket客户端实现

解决WebSocket客户端异步收发+工作线程通信的问题

一、修复过时示例的mpsc问题

你找到的旧示例使用的是旧版futures 0.1的futures::sync::mpsc,这个API在新版futures(0.3+)中已经迁移到futures_channel::mpsc,同时Tokio也提供了更契合自身生态的异步mpsc实现tokio::sync::mpsc。修复步骤如下:

  1. 更新依赖
    修改Cargo.toml,使用支持异步和TLS的新版依赖:

    [dependencies]
    tokio = { version = "1.0", features = ["full"] }
    websocket = { version = "0.26", features = ["async-tls"] }
    futures = "0.3"
    futures-channel = "0.3"
    
  2. 替换mpsc导入
    将示例中的use futures::sync::mpsc;替换为二选一:

    // 适配旧示例逻辑的futures-channel实现
    use futures_channel::mpsc;
    // 更推荐的Tokio原生实现,和异步运行时更兼容
    use tokio::sync::mpsc;
    
  3. 适配异步语法
    旧示例用的是futures 0.1的链式调用(map/then),新版Rust建议用async/await替换,简化代码逻辑,比如将future.and_then(...)改为直接await调用。

二、推荐的现代实现思路(基于Tokio+异步WebSocket)

如果不想修复旧示例,直接用现代异步生态实现更可靠。核心是用Tokio稳定版的select!宏同时监听WebSocket入站消息和工作线程的消息,无需忙等待,同时原生支持TLS。

实现逻辑

  1. 建立带TLS的异步WebSocket连接
  2. 创建双向mpsc通道:主线程→工作线程(分发耗时任务),工作线程→主线程(发送消息到服务器)
  3. 启动工作线程池,处理耗时任务后将结果发回主线程
  4. 主线程通过tokio::select!同时处理WebSocket入站事件(响应Ping、分发任务)和工作线程的消息(转发到服务器)

完整代码示例

# Cargo.toml
[package]
name = "async-ws-client"
version = "0.1.0"
edition = "2021"

[dependencies]
tokio = { version = "1.35", features = ["full"] }
websocket = { version = "0.26", features = ["async-tls"] }
async-tls = "0.12"
use std::time::Duration;
use tokio::sync::mpsc;
use websocket::async::Client;
use websocket::message::OwnedMessage;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 连接带TLS的WebSocket服务器
    let ws_url = "wss://echo.websocket.org";
    let (mut client, response) = Client::connect(ws_url).await?;
    println!("Connected to server: {:?}", response.status);

    // 创建双向通信通道:
    // tx_to_server: 工作线程→主线程,用于发送消息到服务器
    // tx_to_worker: 主线程→工作线程,用于分发耗时任务
    let (tx_to_server, mut rx_from_worker) = mpsc::channel(100);
    let (tx_to_worker, mut rx_from_main) = mpsc::channel(100);

    // 启动2个工作线程处理耗时任务
    for worker_id in 0..2 {
        let worker_tx = tx_to_server.clone();
        tokio::spawn(async move {
            while let Some(task_content) = rx_from_main.recv().await {
                // 模拟耗时操作(比如计算、IO)
                tokio::time::sleep(Duration::from_secs(1)).await;
                println!("Worker {} completed task: {}", worker_id, task_content);

                // 将处理结果发回主线程,由主线程发送到服务器
                let reply_msg = OwnedMessage::Text(format!(
                    "Worker {} response: {}",
                    worker_id, task_content
                ));
                if let Err(e) = worker_tx.send(reply_msg).await {
                    eprintln!("Worker {} failed to send reply: {}", worker_id, e);
                }
            }
        });
    }

    // 主线程核心循环:同时监听WebSocket和工作线程消息
    loop {
        tokio::select! {
            // 处理WebSocket入站消息
            incoming_msg = client.recv() => {
                let msg = incoming_msg?;
                match msg {
                    OwnedMessage::Ping(ping_data) => {
                        // 强制响应Ping,维持连接
                        client.send(OwnedMessage::Pong(ping_data)).await?;
                        println!("Sent Pong to server");
                    }
                    OwnedMessage::Text(text) => {
                        println!("Received from server: {}", text);
                        // 把耗时任务分发到工作线程
                        if let Err(e) = tx_to_worker.send(text).await {
                            eprintln!("Failed to send task to workers: {}", e);
                        }
                    }
                    OwnedMessage::Close(close_frame) => {
                        println!("Server initiated close: {:?}", close_frame);
                        // 响应关闭帧
                        client.send(OwnedMessage::Close(None)).await?;
                        break;
                    }
                    // 忽略其他类型消息
                    _ => {}
                }
            }
            // 处理工作线程发来的消息,发送到服务器
            worker_msg = rx_from_worker.recv() => {
                let msg = worker_msg.ok_or("Worker communication channel closed")?;
                client.send(msg).await?;
                println!("Sent worker reply to server");
            }
        }
    }

    Ok(())
}

方案优势

  • 原生支持TLS:依赖websocket-rs的async-tls特性,自动处理TLS握手和加密通信
  • 无忙等待:用Tokio稳定的select!宏同时监听多个异步事件,CPU利用率高
  • 线程安全:通过mpsc通道实现主线程和工作线程的安全通信,避免直接共享WebSocket连接
  • 易扩展:可以根据需求调整工作线程数量,或者添加更多异步任务到select!中

内容的提问来源于stack exchange,提问作者Maximilian Keßler

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 22:45:54