如何为futures::stream::iter()设置步长10?编译报错求助
解决方法
核心问题是:futures::stream::Iter是Stream类型,而step_by是标准库Iterator trait的方法,不能直接在Stream上调用。你需要先给底层的迭代器设置步长,再转换成Stream。
修改方案
将stream::iter(1..).step_by(10)改为stream::iter((1..).step_by(10))——先让整数范围迭代器按步长10生成元素,再把这个处理后的迭代器转换成Stream。
修改后的完整代码
use futures::stream::{self as stream, StreamExt, TryStreamExt}; 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() { // create iterator that will stream async responses let client = reqwest::Client::new(); let (tx, mut rx) = mpsc::channel::<Post>(2); tokio::spawn(async move { let _ = stream::iter((1..).step_by(10)) // make the request .then(|i| { let client = &client; let url = format!("https://jsonplaceholder.typicode.com/posts/{i}"); client.get(url).send() }) // deserialize .and_then(|resp| resp.json()) .try_for_each_concurrent(2, |r| async { let tx_cloned = tx.clone(); let _ = tx_cloned.send(r).await; Ok(()) }) .await; }); // consume responses from our channel to do future things with results... while let Some(r) = rx.recv().await { println!("{:?}", r); } }
补充说明
如果后续需要在Stream层面实现类似步长过滤的逻辑(比如基于Stream元素的动态步长),可以结合enumerate和filter:
stream::iter(1..) .enumerate() .filter(|(idx, _)| idx % 10 == 0) .map(|(_, val)| val)
但对于固定步长的场景,直接在迭代器阶段处理效率更高。
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

