如何在`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只处理当前阶段的操作:
- 每个阶段读取对应数据时,检查
poll_read的返回:- 若返回
Poll::Ready(Ok(())),说明该阶段数据已读完,切换到下一个状态; - 若返回
Poll::Pending,直接返回Pending即可——底层reader已经注册了当前任务的waker,数据就绪时会自动唤醒,不需要手动操作; - 若返回错误,直接向上传递错误。
- 若返回
- 必须保证每次返回
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
相关产品推荐
相关产品推荐

