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

使用Tungstenite crate时read_message接收消息几小时后卡住求助

Tungstenite WebSocket 长时间运行后卡在read_message的排查与解决

问题核心

使用tungstenite crate建立WebSocket连接后,初期运行正常,但数小时后程序会卡在socket.read_message()调用处,无法继续执行后续逻辑。切换rust-tls、native-tls特性后问题依旧。

排查方向与解决方案

1. TCP静默断开未被检测

WebSocket基于TCP协议,网络中间设备(路由器、防火墙)通常会对长时间无数据的连接主动断开,但TCP默认机制不会立即通知两端。此时tungstenite的阻塞式read_message()会一直等待数据,不会抛出错误。

解决:启用心跳机制
定期发送Ping帧,要求服务器返回Pong帧,若指定时间内未收到Pong,则判定连接失效并重建。示例代码:

use tungstenite::{Message, WebSocket};
use std::time::{Duration, Instant};
use std::io::ErrorKind;
use chrono::Local;

let mut socket = /* 初始化WebSocket连接 */;
let ping_interval = Duration::from_secs(30); // 每30秒发送一次Ping
let mut last_ping_time = Instant::now();

loop {
    println!("before read_message");
    // 设置非阻塞读取,避免一直卡住
    if let Err(e) = socket.set_nonblocking(true) {
        println!("设置非阻塞失败: {}", e);
        break;
    }

    match socket.read_message() {
        Ok(msg) => {
            println!("after read_message");
            // 收到Pong帧,更新心跳时间
            if msg.is_pong() {
                last_ping_time = Instant::now();
            }
            // 处理业务消息逻辑
            let now = Local::now().timestamp_nanos();
        }
        Err(e) if e.kind() == ErrorKind::WouldBlock => {
            // 无消息可读,检查是否需要发送Ping
            if Instant::now().duration_since(last_ping_time) >= ping_interval {
                if let Err(e) = socket.write_message(Message::Ping(vec![])) {
                    println!("发送Ping失败: {}", e);
                    break;
                }
                last_ping_time = Instant::now();
            }
            // 短暂休眠避免空转消耗CPU
            std::thread::sleep(Duration::from_millis(100));
        }
        Err(e) => {
            println!("读取消息错误: {}", e);
            break;
        }
    }
}

2. 阻塞式读取无超时限制

默认的read_message()是无超时的阻塞调用,若连接异常但未触发TCP错误,会无限期卡住。

解决:添加读取超时
推荐使用异步版本的tokio-tungstenite,配合tokio::time::timeout为读取操作设置超时时间:

use tokio_tungstenite::WebSocketStream;
use tokio::time::{timeout, Duration};
use futures_util::StreamExt;
use chrono::Local;

#[tokio::main]
async fn main() {
    let mut socket = /* 初始化异步WebSocket连接 */;
    let read_timeout = Duration::from_secs(60); // 60秒超时

    loop {
        println!("before read_message");
        match timeout(read_timeout, socket.next()).await {
            Ok(Some(Ok(msg))) => {
                println!("after read_message");
                // 处理消息逻辑
                let now = Local::now().timestamp_nanos();
            }
            Ok(None) => {
                println!("连接已关闭");
                break;
            }
            Ok(Some(Err(e))) => {
                println!("读取消息错误: {}", e);
                break;
            }
            Err(_) => {
                println!("读取超时,连接已失效");
                // 断开并重建连接逻辑
                break;
            }
        }
    }
}

3. 服务器未发送关闭帧直接断开

部分服务器在长时间无交互后会直接断开TCP连接,不发送WebSocket关闭帧,导致客户端无法感知连接状态。

解决:结合心跳+超时双重检测
通过心跳确认连接活性,同时用超时限制读取操作,确保能及时发现这类静默断开,主动重建连接。

4. 资源泄漏间接导致阻塞

长时间运行后若存在内存、文件句柄等资源泄漏,可能间接引发IO操作阻塞。可以使用cargo flamegraph或perf工具分析程序性能,排查是否存在异常的资源占用或线程阻塞情况。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 23:31:46