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

如何在Futures 0.2中合并Stream?替代Futures 0.1的merge方法

Merging Streams in Futures 0.2 (and Replacing 0.1's merge)

Hey there! Let's break down how to handle stream merging in Futures 0.2 (the base for the modern futures crate 1.x series) since the old merge method from 0.1 isn't part of the core API anymore.

First: The Equivalent of 0.1's merge in 0.2

In Futures 0.1, Stream::merge let you combine two streams so elements from either were emitted the moment they became available. In Futures 0.2, this exact behavior is replaced by the select method from the StreamExt trait (you'll need to import this trait to access it).

Crucially, unlike zip (which waits for elements from both streams to pair them up), select will immediately yield elements from whichever stream produces them first—exactly what you want instead of zip.

Here's a straightforward example:

use futures::{stream, Stream, StreamExt};
use tokio; // Using tokio as our async runtime

#[tokio::main]
async fn main() {
    // Create two sample streams with different data types
    let stream_strings = stream::iter(vec!["hello", "world", "!"]);
    let stream_numbers = stream::iter(vec![1, 2, 3]);

    // Merge the two streams with select
    let merged_stream = stream_strings.select(stream_numbers);

    // Process each element as it arrives
    merged_stream.for_each(|item| async move {
        println!("Received: {:?}", item);
    }).await;
}

When you run this, you'll see elements from both streams interleaved (the exact order depends on which stream is ready first, but it won't wait for pairs like zip does).

Merging More Than Two Streams

If you need to combine three or more streams, use stream::select_all instead. This function takes an iterator of boxed streams (to ensure they all have the same type) and merges all of them into a single stream.

Example:

use futures::{stream, Stream, StreamExt};
use tokio;

#[tokio::main]
async fn main() {
    let stream1 = stream::iter(vec![10, 20, 30]);
    let stream2 = stream::iter(vec!["foo", "bar", "baz"]);
    let stream3 = stream::iter(vec![true, false]);

    // Box each stream to unify their types
    let streams = vec![
        stream1.boxed(),
        stream2.boxed(),
        stream3.boxed(),
    ];

    // Merge all three streams into one
    let merged_stream = stream::select_all(streams);

    merged_stream.for_each(|item| async move {
        println!("Got: {:?}", item);
    }).await;
}

Quick Key Notes

  • Don't skip importing StreamExt—it's where all handy stream extension methods (like select) live in Futures 0.2+.
  • select and select_all handle backpressure correctly, just like the old merge method did.
  • These methods prioritize immediate emission of available elements, unlike zip which waits for matching pairs across streams.

内容的提问来源于stack exchange,提问作者Basic Coder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:25:52