futures_util::stream::Peekable未按预期Peek,next()无法获取已Peek项
问题原因分析
你遇到的问题核心在于每次调用peek_first_message时,都新建了一个独立的Peekable包装器,而futures的Peekable工作机制和标准库Iterator的Peekable有个关键差异:
当你调用Peekable::peek()时,它会从原Stream中读取一个项并缓存到自己内部。但这个Peekable是函数内的局部变量,函数执行完毕后它就会被销毁,缓存的项也随之丢弃——而原Stream的读取指针已经前进了,不会回头。所以下一次调用peek_first_message时,新的Peekable会从原Stream的当前位置(也就是上一次被消费后的位置)读取下一个项,自然每次输出不同的消息。
简单说:你每次临时创建的Peekable把peek到的项“私吞”了,用完就扔,原Stream再也拿不到这个项。
修复方案
要解决这个问题,你需要在整个流程中复用同一个Peekable实例,而不是每次函数调用都新建。调整代码如下:
use futures::StreamExt; use serde::Deserialize; use tokio_tungstenite::{connect_async, tungstenite::Message}; use url::Url; #[derive(Debug, Deserialize)] struct WsMessage { // 你的字段定义 } async fn peek_first_message( read: Pin<&mut impl futures::Stream<Item = Result<Message, tungstenite::Error>>>, ) -> anyhow::Result<()> { let x = read .peek() .await .ok_or(anyhow::anyhow!("websocket closed before first message"))? .as_ref(); let ws: WsMessage = match x { Ok(v) => serde_json::from_slice(v.to_text()?.as_bytes())?, Err(e) => return Err(anyhow::anyhow!("failed to read message: {}", e)), }; println!("{ws:?}"); Ok(()) } #[tokio::main] async fn main() -> anyhow::Result<()> { let url = Url::parse("wss://127.0.0.1:12345").unwrap(); let (ws_stream, _) = connect_async(url).await.expect("Failed to connect"); let (_, read) = ws_stream.split(); // 提前创建Peekable,整个流程复用它 let mut read = read.peekable(); tokio::pin!(read); peek_first_message(read.as_mut()).await?; peek_first_message(read.as_mut()).await?; peek_first_message(read.as_mut()).await?; // 如果之后要正常消费消息,直接用同一个read即可,peek过的项不会丢失 while let Some(msg) = read.next().await { // 处理消息逻辑 } Ok(()) }
关键说明
- 把
Peekable的创建移到main函数中,全程复用同一个实例,这样peek缓存的项会一直保存在Peekable内部,调用next()时会优先返回缓存的项,再继续读取原Stream。 - 标准库
Iterator的Peekable之所以不会有这个问题,是因为你通常会在整个迭代流程中持有同一个Peekable实例,而不是每次迭代都新建。你之前的误解源于没有意识到Stream的Peekable同样需要保持实例的持续性。
内容的提问来源于stack exchange,提问作者Sergio Gliesh
相关产品推荐
相关产品推荐

