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

Rust调用带tokio::main的异步websocket函数后后续代码不执行如何解决

问题原因与解决方案

问题根因

  • 滥用#[tokio::main]宏:该宏仅用于程序入口的main函数,作用是启动Tokio运行时并阻塞当前线程,直到标记的异步函数执行完毕。你给ws::create加了该宏,调用时会直接阻塞主线程,而create内部的listener.accept().await是无限循环,永远不会返回,自然不会执行后续代码。
  • 没有对多个无限运行的异步任务做并行调度:WebSocket监听服务和你的数据流处理都是不会终止的异步逻辑,串行调用必然会卡死在第一个任务上。

修复步骤

第一步:调整ws模块的实现

去掉create函数的#[tokio::main]标记,改为普通异步函数,同时新增跨任务通信的广播通道,方便你把数据流内容推给所有WebSocket客户端:

// 其余原有依赖、PeerMap等定义保持不变
use tokio::sync::broadcast;

// 封装WebSocket消息发送器,对外暴露发送方法
#[derive(Clone)]
pub struct WsMessageSender(broadcast::Sender<Vec<u8>>);

impl WsMessageSender {
    pub fn send(&self, payload: Vec<u8>) -> Result<(), broadcast::SendError<Vec<u8>>> {
        self.0.send(payload)?;
        Ok(())
    }
}

pub async fn create() -> Result<(WsMessageSender, impl std::future::Future<Output = Result<(), IoError>>), IoError> {
    let addr = "0.0.0.0:8080".to_string();
    // 初始化广播通道,缓存大小可根据业务调整
    let (broadcast_tx, _) = broadcast::channel(128);
    let sender = WsMessageSender(broadcast_tx.clone());

    let state = PeerMap::new(Mutex::new(HashMap::new()));
    let try_socket = TcpListener::bind(&addr).await;
    let listener = try_socket.expect("Failed to bind");
    println!("Listening on: {}", addr);

    // 把WebSocket服务逻辑封装为Future返回,不直接执行
    let serve_future = async move {
        while let Ok((stream, addr)) = listener.accept().await {
            // 给每个连接生成一个广播接收端
            let client_rx = broadcast_tx.subscribe();
            tokio::spawn(handle_connection(state.clone(), stream, addr, client_rx));
        }
        Ok(())
    };

    Ok((sender, serve_future))
}

// 调整handle_connection,新增广播消息接收逻辑
async fn handle_connection(
    state: PeerMap, 
    stream: TcpStream, 
    addr: SocketAddr,
    mut broadcast_rx: broadcast::Receiver<Vec<u8>>
) -> Result<(), IoError> {
    // 原有WebSocket握手、客户端消息处理逻辑保持不变
    // 新增select分支处理广播消息,推送给当前连接的客户端
    loop {
        tokio::select! {
            // 原有处理客户端上行消息的逻辑,保留即可
            msg = ws_stream.next() => {
                // 原有客户端消息处理代码
            }
            // 处理全局广播消息,下行发给客户端
            Ok(payload) = broadcast_rx.recv() => {
                ws_stream.send(Message::Binary(payload)).await?;
            }
        }
    }
}

第二步:调整main函数实现

把#[tokio::main]移到main函数上,并行调度两个异步任务:

mod ws;

#[tokio::main]
async fn main() -> Result<(), IoError> {
    // 初始化WebSocket服务,拿到消息发送器和服务运行Future
    let (ws_sender, ws_serve_fut) = ws::create().await?;

    // 你的数据流处理逻辑
    let data_process_fut = async move {
        // 初始化你自己的无限数据流
        let mut your_stream = init_your_stream().await;
        while let Some(stream_payload) = your_stream.next().await {
            // 直接调用send方法把数据推给所有WebSocket客户端
            ws_sender.send(stream_payload).unwrap();
        }
        Ok(())
    };

    // 并行等待两个任务执行,任意一个报错都会终止程序
    tokio::try_join!(ws_serve_fut, data_process_fut)?;
    Ok(())
}

额外说明

如果你的数据流不需要等WebSocket服务启动完成才能初始化,也可以把两个任务都用tokio::spawn提交到运行时,再用join!等待所有任务结束,逻辑效果是一致的。

内容的提问来源于stack exchange,提问作者François Richard

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 13:09:03