Streams
Stream is the async cousin of Iterator: it produces a sequence of items over time, driven by .await instead of .next(). Use streams for paginated APIs, WebSockets, file chunks, and event feeds.
Search across all documentation pages
Stream is the async cousin of Iterator: it produces a sequence of items over time, driven by .await instead of .next(). Use streams for paginated APIs, WebSockets, file chunks, and event feeds.
Quick-reference recipe card - copy-paste ready.
use tokio_stream::{self as stream, StreamExt};
#[tokio::main]
async fn main() {
let mut s = stream::iter(vec![10, 20, 30]);
while let Some(n) = s.next().await {
println!("{n}");
}
}When to reach for this: Consuming paginated HTTP, database cursors, Kafka/NATS messages, or SSE/WebSocket events.
use tokio::sync::mpsc;
use tokio_stream::{StreamExt, wrappers::ReceiverStream};
async fn produce(tx: mpsc::Sender<i32>) {
for i in 1..=5 {
tx.send(i).await.unwrap();
}
}
#[tokio::main]
async fn main() {
let (tx, rx) = mpsc::channel(8);
tokio::spawn(produce(tx));
let mut stream = ReceiverStream::new(rx);
let sum: i32 = stream.map(|n| n * 2).sum().await;
println!("{sum}");
}What this demonstrates:
Stream via ReceiverStream.StreamExt provides map, filter, take, buffer_unordered.sum().await consumes the stream.Stream::poll_next: Like Future::poll but yields Option<Item> until None.StreamExt: Ergonomic methods from futures / tokio-stream crates.| Adapter | Purpose |
|---|---|
map / filter | Transform items |
take(n) | Limit count |
timeout | Per-item deadline |
buffer_unordered(n) | Concurrent map with limit |
use futures::Stream;
use std::pin::Pin;
use std::task::{Context, Poll};
struct Counter {
n: u32,
max: u32,
}
impl Stream for Counter {
type Item = u32;
fn poll_next(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<u32>> {
if self.n >= self.max {
Poll::Ready(None)
} else {
let v = self.n;
self.n += 1;
Poll::Ready(Some(v))
}
}
}Stream impls mirror Future pinning rules.try_stream! macro builds fallible streams ergonomically.mpsc or stream::iter with rate limits.map closures must not block. Fix: then with async offload.StreamExt from the crate matching your stream type. Fix: use tokio_stream::StreamExt vs futures::StreamExt (enable features).for await in stable Rust - Use while let Some or StreamExt::next. Fix: Loop with .await on next().| Alternative | Use When | Don't Use When |
|---|---|---|
Vec + loop | Small finite batches | Live unbounded events |
| Channels only | Task coordination | You need iterator-style adapters |
| Polling loop + sleep | Simple cron tick | High-throughput event processing |
| Callback API | C FFI bridge | Idiomatic Rust async services |
Iterator is synchronous pull. Stream is async pull - each item may wait for I/O.
tokio_stream::StreamExt for Tokio streams; futures::StreamExt is broader - import one consistently.
Axum can return Body::from_stream wrapping a Stream<Item = Result<Bytes, Error>>.
sqlx::query(...).fetch(&pool) returns a stream of rows - process with try_next().await.
select! on streams, StreamExt::merge, or futures::stream::select for fair merge.
Bounded channels block producers; buffer_unordered limits in-flight async maps.
TryStream with try_next().await or map_err + ? inside try_stream!.
stream::iter(vec![...]) for deterministic inputs; timeout in tests to catch hangs.
Advanced streams use pin_project for pin-projecting struct fields in poll_next.
tokio-tungstenite messages form a stream of frames - map to app events.
Stack versions: This page was written for Rust 1.97.0 (edition 2024), Tokio 1.x, Axum 0.8, serde 1.0, sqlx 0.8, clap 4, and Polars 0.46+.
Reviewed by Chris St. John·Last updated Jul 16, 2026