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

为何接受Mio的TcpStream后首次写入会出现“operation would block”错误?

解析你的代码中"operation would block"错误的原因

咱们先搞清楚这个错误到底是什么意思——它其实是非阻塞IO模式下的正常临时状态,不是真的操作失败。结合你给出的代码,我来一步步拆解原因:

1. 非阻塞Socket的特性

首先,能触发这个错误,说明你的listener或者stream肯定是设置成了非阻塞模式(虽然代码里没显式写,但用poll做IO多路复用的场景,几乎都会把Socket设为非阻塞)。

在非阻塞模式下,write这类IO操作不会像阻塞模式那样一直等待系统完成操作,而是会立刻返回:如果当前操作能直接完成(比如发送缓冲区还有空间),就返回写入的字节数;如果暂时没法完成(比如发送缓冲区满了,或者客户端还没来得及接收数据),就返回WouldBlock错误,告诉你“现在没法写,你过会儿再试”。

2. 你的代码为什么会触发这个错误

你在accept新连接后,直接调用了stream.write("Hello, world!\n".as_bytes()),完全没考虑非阻塞Socket的这个特性:

  • 如果客户端的接收窗口很小,或者网络有延迟,导致Socket的发送缓冲区被占满,write就会立刻返回WouldBlock。
  • 而你的unwrap_or_continue!宏把这个临时错误当成了需要终止当前连接处理的“失败”,直接输出报错信息并执行continue,这就导致你看到了这个错误提示。

3. 怎么修复这个问题

核心思路是:不要把WouldBlock当成错误,而是要缓存待发送的数据,等Socket可写时再重试。这里给你一个修改后的思路示例:

use std::collections::HashMap;
use std::io::{self, Write};
use std::os::unix::io::AsRawFd;
use std::time::Duration;
use mio::{Events, Poll, PollOpt, Ready, Token};

// 定义唯一的Token常量
const LISTENER_TOKEN: Token = Token(0);
// 给新连接分配Token的计数器(简单示例,实际可以用更安全的方式)
static mut NEXT_TOKEN: usize = 1;

fn next_token() -> Token {
    unsafe {
        let token = NEXT_TOKEN;
        NEXT_TOKEN += 1;
        Token(token)
    }
}

fn main() -> io::Result<()> {
    let mut poll = Poll::new()?;
    let mut events = Events::with_capacity(1024);
    // 假设这里已经初始化了listener并注册到poll...

    // 缓存待发送的数据:Token -> 剩余待发送字节
    let mut pending_writes: HashMap<Token, Vec<u8>> = HashMap::new();

    loop {
        poll.poll(&mut events, Some(Duration::from_secs(0)))?;

        for event in events.iter() {
            match event.token() {
                LISTENER_TOKEN => {
                    let (mut stream, addr) = match listener.accept() {
                        Ok(s) => s,
                        Err(e) => {
                            eprintln!("Failed to accept incoming connection: {}", e);
                            continue;
                        }
                    };
                    // 确保Socket是非阻塞模式
                    stream.set_nonblocking(true)?;

                    let message = "Hello, world!\n".as_bytes().to_vec();
                    let token = next_token();

                    // 尝试第一次写入
                    match stream.write(&message) {
                        Ok(n) if n == message.len() => {
                            // 全部写完,无需后续处理
                            println!("Sent message to {}", addr);
                        }
                        Ok(n) => {
                            // 只写了一部分,缓存剩余数据并注册写事件
                            let remaining = message[n..].to_vec();
                            pending_writes.insert(token, remaining);
                            poll.register(
                                &stream,
                                token,
                                Ready::writable(),
                                PollOpt::level(),
                            )?;
                        }
                        Err(e) if e.kind() == io::ErrorKind::WouldBlock => {
                            // 完全没法写,缓存全部数据并注册写事件
                            pending_writes.insert(token, message);
                            poll.register(
                                &stream,
                                token,
                                Ready::writable(),
                                PollOpt::level(),
                            )?;
                        }
                        Err(e) => {
                            eprintln!("Failed to write to {}: {}", addr, e);
                        }
                    }
                }
                token if pending_writes.contains_key(&token) => {
                    // 获取对应的Socket和待发送数据
                    let mut stream = /* 这里需要你根据Token获取对应的Stream,比如用HashMap存储Token和Stream的映射 */;
                    let data = pending_writes.get_mut(&token).unwrap();

                    match stream.write(data) {
                        Ok(n) => {
                            if n == data.len() {
                                // 全部写完,清理缓存并注销事件
                                pending_writes.remove(&token);
                                poll.reregister(&stream, token, Ready::empty(), PollOpt::empty())?;
                                println!("Finished sending data to connection");
                            } else {
                                // 移除已发送的部分,保留剩余数据
                                *data = data[n..].to_vec();
                            }
                        }
                        Err(e) if e.kind() == io::ErrorKind::WouldBlock => {
                            // 还是没法写,继续监听写事件即可
                        }
                        Err(e) => {
                            eprintln!("Failed to write to connection: {}", e);
                            pending_writes.remove(&token);
                            poll.reregister(&stream, token, Ready::empty(), PollOpt::empty())?;
                        }
                    }
                }
                _ => unreachable!(),
            }
        }
    }
}

简单说,就是把没法立刻发送的数据存起来,然后让poll监听这个Socket的可写事件,等系统通知我们“现在可以写了”,再继续发送剩余的数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:31:31