Tokio Channel无法接收流式Futures返回数据的问题求助
问题分析与解决
你的代码存在两个核心问题,导致Receiver无法正常读取通道数据:
1. 无限流阻塞后续代码执行
你用futures::stream::iter(1..)创建了无限迭代器,尽管jsonplaceholder的/posts/{id}只有1-100的有效数据,但buffer_unordered(100)会持续发起请求,即使遇到404错误后仍会循环,导致resps.for_each(...).await永远无法执行完毕,后面的rx.recv()循环根本没机会运行。
2. 发送端未完全销毁导致通道无法关闭
即便流能正常结束,main函数中保留的原始tx发送端会让Tokio的MPSC通道一直处于开放状态,只要有一个发送端存在,rx.recv()就会持续等待新数据,不会返回None,接收循环无法终止。
修复后的代码
use futures::stream::StreamExt; use serde::{Deserialize, Serialize}; use tokio::sync::mpsc; #[derive(Deserialize, Serialize, Debug)] #[serde(rename_all = "camelCase")] pub struct Post { pub user_id: u16, pub id: u16, pub title: String, pub body: String, } #[tokio::main] async fn main() { // 创建请求客户端 let client = reqwest::Client::new(); // 改为有限流:仅请求1-100的posts(匹配jsonplaceholder的有效范围) let resps = futures::stream::iter(1..=100) .map(|i| { let client = client.clone(); async move { let url = format!("https://jsonplaceholder.typicode.com/posts/{i}"); client.get(url).send().await } }) .buffer_unordered(100); let (tx, mut rx) = mpsc::channel::<Post>(50); // 异步生成任务处理请求与发送,避免阻塞接收逻辑 tokio::spawn(async move { resps .for_each(|r| async { match r { Ok(resp) => { if resp.status().is_success() { // 反序列化成功后发送到通道 if let Ok(post) = resp.json::<Post>().await { let _ = tx.send(post).await; } } else { println!("请求失败,状态码: {}", resp.status()); } } Err(e) => { println!("请求错误: {e}!"); } } }) .await; // 任务结束后tx自动被drop,通道关闭 }); // 接收通道数据并处理 while let Some(post) = rx.recv().await { println!("{:?}", post); } }
关键修复点
- 有限流替代无限流:使用
1..=100匹配API的有效数据范围,确保流能正常结束。 - 并行处理请求与接收:用
tokio::spawn将请求发送逻辑放到独立任务中,让请求和数据接收并行执行,避免前者阻塞后者。 - 自动关闭通道:将
tx移动到异步任务中,任务结束后tx自动销毁,通道关闭,rx.recv()会返回None,接收循环正常终止。 - 增加状态码校验:避免对错误响应做反序列化,减少不必要的报错。
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

