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

如何为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 20:34:55