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
相关产品推荐
相关产品推荐

