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

如何让Tokio Stream实现并发运行?

问题分析与解决:Tokio Stream 并发执行失效的原因及修复

原因分析

  • StreamExt::then 是串行处理流元素:它会等待当前元素对应的异步任务执行完成后,才会拉取流中的下一个元素并启动其任务。你的代码中每个 sleep(1秒) 任务是依次执行的,因此每隔1秒打印一个值,完全没有并发效果。
  • chunks_timeout 仅负责结果收集:它的作用是把上游流的结果按指定数量或超时时间打包成块,但不会改变上游流的执行逻辑——上游还是串行处理元素,所以最终结果虽然按块收集,执行过程依然是串行的。

解决方案

要实现流元素的并发处理,需要使用 Tokio Stream 提供的并发处理方法,比如 buffer_unordered 或 for_each_concurrent。

方案1:使用 buffer_unordered

buffer_unordered 允许同时处理最多 N 个流元素的异步任务,任务完成后会将结果放入输出流,实现指定数量的并发执行:

use tokio;
use tokio_stream::{self as stream, StreamExt};
use std::time::Duration;

#[tokio::main]
async fn main() {
    let stream = stream::iter(0..10)
        // 将每个元素映射为异步任务(使用map而非then)
        .map(|i| async move {
            tokio::time::sleep(Duration::from_secs(1)).await;
            println!("i={:?}", i);
        })
        // 同时并发运行最多3个任务
        .buffer_unordered(3);

    let _result: Vec<_> = stream.collect().await;
}

这段代码会同时启动前3个任务,1秒后这3个任务同时完成并打印;随后自动启动接下来的3个任务,再1秒后打印,以此类推,符合你期望的“每隔1秒打印一组3个数字”的效果。

方案2:使用 for_each_concurrent

如果不需要收集任务结果,只是要并发执行任务,for_each_concurrent 更简洁,直接指定并发数并处理每个元素:

use tokio;
use tokio_stream::{self as stream, StreamExt};
use std::time::Duration;

#[tokio::main]
async fn main() {
    stream::iter(0..10)
        // 最多同时处理3个任务
        .for_each_concurrent(3, |i| async move {
            tokio::time::sleep(Duration::from_secs(1)).await;
            println!("i={:?}", i);
        })
        .await;
}

补充说明

你提到用 join_all 可以正常并发,是因为 join_all 会一次性将所有异步任务提交到 Tokio runtime 并发执行,而 then 是严格串行处理流元素,两者的执行模型完全不同。

内容的提问来源于stack exchange,提问作者Mathieu Dutour Sikiric

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 16:42:22