如何在Futures 0.2中合并Stream?替代Futures 0.1的merge方法
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 (likeselect) live in Futures 0.2+. selectandselect_allhandle backpressure correctly, just like the oldmergemethod did.- These methods prioritize immediate emission of available elements, unlike
zipwhich waits for matching pairs across streams.
内容的提问来源于stack exchange,提问作者Basic Coder

