wandres.dev
ASYNC AVANZADO · streams, select, cancelación

Streams: iteradores asincronos

Un Future produce un valor que llega mas tarde; un Stream produce una SECUENCIA de valores que llegan con el tiempo. Es el Iterator del mundo async: donde next devuelve Option, poll_next devuelve Poll de Option. StreamExt hereda el vocabulario map, filter, then y collect, y se consume con while let Some(x) = stream.next().await.

⏱ 18 min

Un Future representa un valor que aún no ha llegado. Pero muchísimo trabajo asíncrono no produce uno, sino muchos valores repartidos en el tiempo: los mensajes que van entrando por un canal, las líneas que un cliente escribe en un socket, los latidos de un temporizador, las filas que una consulta devuelve por lotes. Para eso existe el Stream: es a Future lo que Iterator es a un valor suelto. Donde el iterador tiene next —“dame el siguiente, o None si se acabó”—, el stream tiene poll_next —“dame el siguiente cuando esté listo, o avísame que aún no hay”—. Y así como implementar Iterator te regalaba setenta adaptadores, implementar o usar un Stream te da, vía StreamExt, el mismo vocabulario —map, filter, then, collect— pero asíncrono. Aprender streams es reconocer que ya sabes iterar; solo hay que añadirle la dimensión del tiempo.

🎯 Al terminar esta lección sabrás
  • Entender Stream como el análogo asíncrono de Iterator: poll_next frente a next.
  • Consumir un stream con el bucle while let Some(x) = stream.next().await.
  • Heredar los adaptadores de StreamExt (map, filter, then, collect) y saber que son perezosos.
  • Conocer de dónde salen los streams y cómo procesarlos con concurrencia acotada.

Stream: el Iterator asíncrono

Recuerda el contrato de Iterator: un tipo asociado Item y un método next que devuelve Option<Item>Some mientras haya, None al agotarse—. Stream calca esa forma y le añade lo único que le faltaba para vivir en async: la posibilidad de decir “todavía no”:

pub trait Stream {
    type Item;
    fn poll_next(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
    ) -> Poll<Option<Self::Item>>;
}

Léelo comparándolo con next. Iterator::next devuelve Option<Item> de inmediato; Stream::poll_next devuelve Poll<Option<Item>>, y ahí está toda la diferencia. Son tres desenlaces, no dos:

  • Poll::Ready(Some(v)) → hay un valor listo, v; el stream sigue vivo.
  • Poll::Ready(None) → el stream terminó; no vendrán más valores.
  • Poll::Pending → aún no hay valor; el runtime volverá a sondear cuando lo haya.

Ese Pending es la costura async: es donde el stream cede, igual que un future cede en .await. poll_next es maquinaria de bajo nivel —con Pin y Context, herencia del nivel de Pin—; casi nunca lo llamarás a mano, igual que rara vez llamas a poll de un future. Para eso está la capa ergonómica.

ℹ️
Stream aun no vive en la biblioteca estandar

A día de hoy, en Rust 2024, Stream no está estabilizado en std: vive en el crate futures (como futures::Stream) y lo reexporta tokio_stream. La versión de std se llama AsyncIterator y sigue siendo nightly, y todavía no hay sintaxis for await en estable. Por eso, en la práctica, importas el trait desde futures o tokio_stream y consumes con un while let. El diseño está asentado; solo falta que cruce la última puerta hacia std.

Consumir un stream: while let Some

La forma de recorrer un stream es el eco async del bucle for sobre un iterador. El método .next() —que aporta la extensión StreamExt— devuelve un future de Option<Item>; lo esperas con .await, y mientras salga Some, procesas:

use tokio_stream::StreamExt;

async fn sumar(mut s: impl Stream<Item = u32> + Unpin) -> u32 {
    let mut total = 0;
    while let Some(n) = s.next().await { // cede en cada espera
        total += n; // corre cuando hay valor; None termina el bucle
    }
    total
}

Lee la línea clave despacio: s.next().await es “pídeme el siguiente y cede mientras no llegue”. Cuando el stream produce Some(n), el cuerpo del bucle corre; cuando produce None, el while let termina, exactamente como un for se acaba al agotarse el iterador. La diferencia es que entre un valor y el siguiente puede pasar tiempo real, y en ese tiempo la tarea no ocupa el hilo: lo cede para que otras avancen.

Un detalle que el compilador te recordará: .next() exige que el stream sea Unpin, o que esté pineado. Si tienes un stream que no lo es, lo fijas en el sitio con tokio::pin! antes del bucle, igual que hacías con futures reutilizados en select!.

let flujo = crear_stream();
tokio::pin!(flujo); // ahora next() lo acepta
while let Some(x) = flujo.next().await {
    manejar(x);
}

StreamExt: el vocabulario heredado

Igual que Iterator gana decenas de métodos por defecto sobre next, un Stream gana todo un vocabulario a través del trait de extensión StreamExt. Los adaptadores son los mismos que ya conoces, con dos añadidos para lo asíncrono, y todos son perezosos: construyen un stream nuevo que no hace nada hasta que lo consumes.

use tokio_stream::StreamExt;

let resumen: Vec<String> = eventos()
    .filter(|e| e.nivel >= Nivel::Warn) // conserva solo los graves
    .map(|e| e.mensaje)                 // transforma cada elemento
    .take(100)                          // corta a los primeros 100
    .collect()                          // consume: recorre y acumula
    .await;                             // collect sobre un stream es async

Dos adaptadores no tienen gemelo en Iterator porque solo tienen sentido con el tiempo: then, que es como map pero su función es async —devuelve un future por cada elemento y lo espera—, y collect, que aquí es un consumidor asíncrono y por eso se .await. La regla mental es directa: map para transformar con una función síncrona, then cuando cada elemento dispara una espera.

// then: por cada id del stream, lanza una peticion async y espera su respuesta.
let cuerpos = ids()
    .then(|id| async move { descargar(id).await })
    .collect::<Vec<_>>()
    .await;

De dónde salen los streams y cómo acotarlos

Rara vez implementarás poll_next a mano; casi siempre envuelves algo que ya produce valores en el tiempo. Las fuentes habituales: un mpsc::Receiver convertido en stream con tokio_stream::wrappers::ReceiverStream, un temporizador periódico con IntervalStream, o un Vec vuelto stream con tokio_stream::iter. Y aquí aparece la conexión con la lección 1: un stream de futures puede procesarse con concurrencia acotada, algo que join_all no te daba.

use futures::stream::StreamExt;

// Procesa hasta 8 descargas A LA VEZ, no todas de golpe ni de una en una.
let resultados: Vec<_> = futures::stream::iter(urls)
    .map(|u| async move { descargar(u).await })
    .buffer_unordered(8) // techo de concurrencia: 8 en vuelo
    .collect()
    .await;

buffer_unordered(n) mantiene hasta n futures corriendo simultáneamente y emite cada resultado en cuanto está —sin respetar el orden—; buffered(n) hace lo mismo pero conserva el orden de entrada. Es el punto medio exacto entre la latencia de procesar uno a uno y el riesgo de join_all de disparar diez mil peticiones a la vez. Además, como el stream es pull —el consumidor tira con poll_next—, obtienes contrapresión natural: nadie produce más rápido de lo que consumes.

🔮

Future: un valor

Una espera, un resultado. Output es el tipo prometido. Se conduce con .await.

🔁

Iterator: muchos, sync

Una secuencia disponible ya. next da Option. Se recorre con for.

🌊

Stream: muchos, async

Una secuencia repartida en el tiempo. poll_next da Poll de Option. Se recorre con while let Some.

🎚️

buffer_unordered

Procesa un stream de futures con un techo de concurrencia. Ni de uno en uno ni todos a la vez: n en vuelo.

flowchart LR
N[Consumidor hace next await] --> PN[poll_next del Stream]
PN --> R{Estado}
R -->|Ready Some v| V[Entrega v corre el cuerpo del while]
R -->|Pending| Y[Cede al runtime espera el wake]
R -->|Ready None| F[El stream termino sale del bucle]
Y --> PN
V --> N
style V fill:#a6e3a1,color:#11111b
style Y fill:#fab387,color:#11111b
style F fill:#89b4fa,color:#11111b
El cuadrante que completa: async multiplica la iteracion por el tiempo

Hay una simetría profunda escondida en estos cuatro nombres, y verla reordena todo el modelo mental de la concurrencia en Rust. Cruza dos ejes: cuántos valores produce algo —uno o muchos— y cuándo están disponibles —ya o con el tiempo—. Un valor, ya: un T corriente. Muchos, ya: un Iterator, la secuencia que puedes recorrer de un tirón. Uno, con el tiempo: un Future, la promesa de un resultado que aún no llegó. Y la casilla que faltaba, muchos y con el tiempo: el Stream. No es una abstracción nueva y ajena que haya que aprender desde cero, es la composición de dos ideas que ya dominabas —iterar y esperar—, y por eso hereda de ambas. De Iterator toma la forma del contrato: un tipo Item, un método que da el siguiente, la pereza que hace que un adaptador no compute hasta que alguien consume, el ecosistema entero de map y filter construido sobre un único método esencial. De Future toma la costura del tiempo: el Poll::Pending que dice “todavía no”, la cesión al runtime, la reanudación exacta donde se paró. La consecuencia práctica es que casi todo lo que produce datos poco a poco —una conexión de red, un archivo que se lee por trozos, una suscripción a eventos, una cola de trabajos— habla el mismo idioma, compone con los mismos adaptadores y se recorre con el mismo while let. Y la consecuencia de diseño es que la contrapresión sale gratis: como el stream es pull —el consumidor tira del siguiente con poll_next y no antes—, el productor nunca corre más rápido que quien consume, sin que tengas que orquestar ningún control de flujo. Un modelo push, donde los valores te llegan encima los pidas o no, tendría que inventar colas, límites y descartes para no ahogarte; el pull asíncrono resuelve eso por construcción. Iterar era recorrer el espacio; con streams, recorres el tiempo con las mismas herramientas.

📝
Lo esencial

Un Stream es el Iterator de async: poll_next devuelve Poll<Option<Item>>, con tres desenlaces —Ready(Some), Ready(None), Pending—. Se consume con while let Some(x) = stream.next().await, donde .next() (de StreamExt) es un future de Option. StreamExt hereda map, filter, collect y añade then (map async) y un collect que se .await; todo perezoso. Stream aún vive en futures/tokio_stream, no en std. Los streams se obtienen envolviendo fuentes (canales, temporizadores) y se procesan con concurrencia acotada vía buffer_unordered, con contrapresión natural por ser pull.

⚔️ Recorre el tiempo con streams
  1. Convierte un Vec<u32> en stream con tokio_stream::iter y súmalo con un bucle while let Some; compáralo mentalmente con hacerlo con un for.
  2. Encadena filter, map y take sobre un stream y consúmelo con collect().await; comprueba que nada corre hasta el collect (pereza).
  3. Usa ReceiverStream para envolver un mpsc::Receiver y consúmelo con while let; envía valores desde otra tarea y observa cómo el consumidor cede entre uno y otro.
  4. Reescribe un map seguido de .await manual como un solo then, y explica por qué then existe y map no basta cuando cada elemento dispara una espera.
  5. Procesa una lista de URLs con buffer_unordered(4) y compáralo con join_all (todas a la vez) y con un bucle secuencial (una a una); razona qué problema resuelve el techo de concurrencia.