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

如何在`fn poll_..`中多次poll Future?AsyncRead解码器实现疑问

实现AsyncRead时Pin Project的重复poll问题与唤醒方式的正确性

核心误区澄清

你提到调用this.reader.poll_read()会消耗this.reader,其实是误解——self.project()得到的this.reader是Pin<&mut R>类型(因为你给reader加了#[pin]标记),调用poll_read只是对这个可变引用的方法调用,并不会消耗它。你觉得没法重复调用,是因为代码里直接用?处理返回值,一旦poll_read返回Pending,函数就提前返回了,而非this.reader被消耗。

手动唤醒的做法是否正确?

你用cx.waker().wake_by_ref()再返回Pending的方式能跑,但绝对不是正确的异步IO实践:

  • 这会导致无意义的忙轮询,浪费CPU资源;
  • 违反了AsyncRead的契约:只有当底层IO事件就绪时才应该唤醒任务,主动唤醒自己相当于强制调度,完全违背了异步IO的设计初衷。

正确实现思路:基于状态机分阶段处理

你的AsyncDecoder必须靠状态机(也就是你定义的DecoderState)来跟踪解码进度,把“读文件头→读块头→读压缩块→循环”拆分成独立阶段,每次poll_read只处理当前阶段的操作:

  1. 每个阶段读取对应数据时,检查poll_read的返回:
    • 若返回Poll::Ready(Ok(())),说明该阶段数据已读完,切换到下一个状态;
    • 若返回Poll::Pending,直接返回Pending即可——底层reader已经注册了当前任务的waker,数据就绪时会自动唤醒,不需要手动操作;
    • 若返回错误,直接向上传递错误。
  2. 必须保证每次返回Ready(Ok(()))时,至少向输出buf写入一个字节(除非是EOF),否则会违反AsyncRead的契约(未写入字节即表示EOF)。

简化实现示例

#[pin_project]
pub struct AsyncDecoder<R> {
    #[pin]
    reader: R,
    state: DecoderState,
    // 存储阶段读取的临时数据,比如文件头、块头
    temp_buf: Vec<u8>,
    // 保存未写完的解压数据,下次poll_read优先写入
    pending_decompressed: Vec<u8>,
}

enum DecoderState {
    ReadingFileHeader,
    ReadingBlockHeader(usize), // 记录需要读取的块头长度
    ReadingCompressedBlock(usize), // 记录需要读取的压缩块长度
    Done,
}

impl<R> AsyncRead for AsyncDecoder<R> where R: AsyncRead {
    fn poll_read(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &mut ReadBuf<'_>,
    ) -> Poll<io::Result<()>> {
        let mut this = self.project();

        // 先处理上次未写完的解压数据
        if !this.pending_decompressed.is_empty() {
            let copy_len = this.pending_decompressed.len().min(buf.remaining());
            buf.put_slice(&this.pending_decompressed[..copy_len]);
            this.pending_decompressed.drain(..copy_len);
            // 写入了数据,直接返回Ready
            if copy_len > 0 {
                return Poll::Ready(Ok(()));
            }
        }

        loop {
            match this.state {
                DecoderState::ReadingFileHeader => {
                    // 准备4字节空间存文件头
                    this.temp_buf.resize(4, 0);
                    let mut read_buf = ReadBuf::new(&mut this.temp_buf);

                    match this.reader.as_mut().poll_read(cx, &mut read_buf)? {
                        Poll::Ready(_) => {
                            // 解析文件头,获取后续块头长度(示例逻辑)
                            let block_header_len = parse_file_header(&this.temp_buf);
                            *this.state = DecoderState::ReadingBlockHeader(block_header_len);
                        }
                        Poll::Pending => return Poll::Pending,
                    }
                }
                DecoderState::ReadingBlockHeader(len) => {
                    this.temp_buf.resize(len, 0);
                    let mut read_buf = ReadBuf::new(&mut this.temp_buf);

                    match this.reader.as_mut().poll_read(cx, &mut read_buf)? {
                        Poll::Ready(_) => {
                            // 解析块头,获取压缩块长度(示例逻辑)
                            let block_len = parse_block_header(&this.temp_buf);
                            *this.state = DecoderState::ReadingCompressedBlock(block_len);
                        }
                        Poll::Pending => return Poll::Pending,
                    }
                }
                DecoderState::ReadingCompressedBlock(len) => {
                    this.temp_buf.resize(len, 0);
                    let mut read_buf = ReadBuf::new(&mut this.temp_buf);

                    match this.reader.as_mut().poll_read(cx, &mut read_buf)? {
                        Poll::Ready(_) => {
                            // 解压数据
                            let decompressed = decompress(&this.temp_buf)?;
                            // 写入输出buf,剩余的存到pending_decompressed
                            let copy_len = decompressed.len().min(buf.remaining());
                            buf.put_slice(&decompressed[..copy_len]);
                            if copy_len < decompressed.len() {
                                this.pending_decompressed.extend_from_slice(&decompressed[copy_len..]);
                            }
                            // 切换回读块头状态,继续循环
                            match get_next_block_header_len() {
                                Ok(next_len) => *this.state = DecoderState::ReadingBlockHeader(next_len),
                                Err(_) => *this.state = DecoderState::Done,
                            }
                            // 写入了数据就返回,没写入的话继续循环处理下一个阶段
                            if copy_len > 0 {
                                return Poll::Ready(Ok(()));
                            }
                        }
                        Poll::Pending => return Poll::Pending,
                    }
                }
                DecoderState::Done => {
                    // 没有更多数据,返回EOF(此时buf未被填充,符合AsyncRead定义)
                    return Poll::Ready(Ok(()));
                }
            }
        }
    }
}

// 示例辅助函数(需自行实现)
fn parse_file_header(buf: &[u8]) -> usize { 4 }
fn parse_block_header(buf: &[u8]) -> usize { 8 }
fn decompress(buf: &[u8]) -> io::Result<Vec<u8>> { Ok(Vec::new()) }
fn get_next_block_header_len() -> io::Result<usize> { Ok(4) }

关键注意事项

  • 永远不要一次性把所有解码步骤写死,必须通过状态机分阶段处理,每次只做当前阶段的IO操作;
  • 底层reader返回Pending时,直接返回即可,手动唤醒是画蛇添足,还会浪费资源;
  • 处理解压数据时,一定要考虑输出buf空间不足的情况,把剩余数据存到结构体的临时字段里,下次poll_read优先处理;
  • 严格遵守AsyncRead的契约:返回Ready(Ok(()))时,要么至少写入一个字节,要么确实是EOF。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 08:45:28