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

BufferedReader行迭代器无报错挂起:Redis模拟TCP服务器异常

问题根源

你的代码逻辑错误在于等待TCP连接关闭才会处理命令,但redis-cli使用长连接,发送完PING命令后不会主动断开连接,导致程序一直卡在读取循环里,直到你中断redis-cli(此时连接关闭)才会执行后续的命令处理逻辑。

而用nc测试时,echo发送完数据后会关闭标准输入,nc随之关闭连接,程序读取到EOF后跳出循环,才能正常返回PONG。

解决方案

需要按照Redis的RESP协议规范,解析出完整的一条命令后就停止读取,而不是一直等待连接关闭。以PING命令的RESP格式*1\r\n$4\r\nPING\r\n为例,解析流程应该是:

  • 读取第一行*1,知道这是包含1个元素的数组
  • 读取第二行$4,知道下一个元素的长度是4
  • 读取第三行PING,拿到完整的命令参数
  • 此时已经获取到完整的PING命令,立即跳出读取循环,处理并返回响应

修正后的核心代码(基于你补充的异步版本)

修正handle_stream函数

async fn handle_stream(mut _stream: TcpStream) {
    let mut sock = BufReader::new(_stream);
    let mut line = String::new();

    // 1. 读取数组长度行,格式如 *1
    line.clear();
    if sock.read_line(&mut line).await.unwrap() == 0 {
        return;
    }
    let array_len = line.trim().strip_prefix('*').and_then(|s| s.parse::<usize>().ok()).unwrap_or(0);
    if array_len == 0 {
        write_string_to_stream(String::from("-ERR invalid array format\r\n"), &mut sock).await;
        return;
    }

    let mut command = RedisCommand::UNKNOWN;
    let mut out = String::new();

    // 2. 遍历数组中的每个元素
    for i in 0..array_len {
        // 读取元素长度行,格式如 $4
        line.clear();
        if sock.read_line(&mut line).await.unwrap() == 0 {
            return;
        }
        let elem_len = line.trim().strip_prefix('$').and_then(|s| s.parse::<usize>().ok()).unwrap_or(0);
        if elem_len == 0 {
            write_string_to_stream(String::from("-ERR invalid element format\r\n"), &mut sock).await;
            return;
        }

        // 读取元素内容(包含末尾的\r\n)
        let mut elem_buf = vec![0; elem_len + 2];
        if sock.read_exact(&mut elem_buf).await.is_err() {
            write_string_to_stream(String::from("-ERR failed to read element\r\n"), &mut sock).await;
            return;
        }
        let elem_str = String::from_utf8_lossy(&elem_buf[0..elem_len]).trim().to_string();

        // 第一个元素是命令
        if i == 0 {
            if let Ok(cmd) = RedisCommand::from_str(&elem_str) {
                command = cmd;
            } else {
                write_string_to_stream(String::from("-ERR unknown command\r\n"), &mut sock).await;
                return;
            }
        }
        // ECHO命令的第二个元素是要返回的内容
        if i == 1 && command == RedisCommand::ECHO {
            out = elem_str;
        }
    }

    // 3. 处理命令并返回符合RESP格式的响应
    match command {
        RedisCommand::UNKNOWN => write_string_to_stream(String::from("-ERR invalid command\r\n"), &mut sock).await,
        RedisCommand::ECHO => write_string_to_stream(format!("+{}\r\n", out), &mut sock).await,
        RedisCommand::PING => write_string_to_stream(String::from("+PONG\r\n"), &mut sock).await,
    }
}

修正write_string_to_stream函数

Redis客户端期望响应严格符合RESP协议,同时需要强制刷新缓冲区确保数据立即发送:

async fn write_string_to_stream(out: String, sock: &mut BufReader<TcpStream>) {
    let mut writer = BufWriter::new(sock.get_mut());
    if let Err(_) = writer.write_all(out.as_bytes()).await {
        writer.shutdown().await.ok();
        return;
    }
    // 强制刷新缓冲区,避免数据滞留
    let _ = writer.flush().await;
}
额外注意事项
  • RESP协议的行结束符是\r\n,不是单纯的\n,解析和响应时都要严格遵循
  • 异步代码必须使用Tokio提供的异步IO方法,禁止混用标准库的同步IO,否则会阻塞Tokio运行时
  • Redis默认使用长连接,命令处理完成后无需立即关闭连接,客户端可继续发送后续命令

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 19:44:52