Stream bridges

fromStream, fxStream y .toStream() — cruza libremente entre el Stream de Dart y el FxAsyncIterable de FxDart.

FxAsyncIterable<T> fromStream<T>(Stream<T> stream) // lossless FIFO + pause FxAsyncIterable<T> fromStreamLatest<T>(Stream<T> stream) // drop superseded FxAsyncIterable<List<T>> fromStreamChunked<T>(Stream<T> stream) // batches FxAsyncIterable<T> fromStreamNext<T>(Stream<T> stream) // demand-gated drop FxAsync<T> fxStream<T>(Stream<T> stream) // chain, wraps fromStream Stream<T> FxAsyncIterableToStream.toStream<T>() // extension on FxAsyncIterable Stream<T> FxAsync.toStream() // chain FxAsync<T> FxEvents.pull() // = fromStream FxAsync<T> FxEvents.pullLatest() // = fromStreamLatest FxAsync<List<T>> FxEvents.pullChunked() // = fromStreamChunked FxAsync<T> FxEvents.pullNext() // = fromStreamNext

Lección

fromStream convierte cualquier Stream —de suscripción única o broadcast— en un FxAsyncIterable, para que puedas aplicar todo el conjunto de operadores de FxDart (map, filter, concurrent, …) sobre datos que llegan desde un socket, un fichero, el stream de eventos de un widget o de cualquier otro sitio donde Dart te dé un Stream. fxStream(stream) hace lo mismo, pero devuelve directamente un FxAsync encadenable en lugar de un FxAsyncIterable a secas — el equivalente asíncrono de fx y fxAsync.

En el sentido contrario, .toStream() lleva un FxAsyncIterable (o una cadena FxAsync) hasta el final y reemite sus valores como un Stream normal — útil cuando otra API (un StreamBuilder, por ejemplo) espera uno. Una advertencia: toStream() siempre tira de los valores secuencialmente, ignorando cualquier concurrent(n) que haya aguas arriba — aplica concurrent/concurrentPool a la cadena antes de llamar a .toStream() si quieres que el paralelismo ocurra de verdad; la conversión a stream por sí sola no lo añade.

Cruzar de push a pull no es una operación: son cuatro, porque un stream puede seguir emitiendo mientras el consumidor está ocupado. RxJS 9 nombra las cuatro iterateEach, iterateLatest, iterateBuffered e iterateNext. FxDart las mapea a fromStream* (iterable crudo) y FxEvents.pull* (cadena):

RxJS 9FxDartmientras estás ocupado
iterateEachfromStream / .pull()FIFO sin pérdidas — pausa la fuente, encola cada valor
iterateLatestfromStreamLatest / .pullLatest()descarta lo superado — quédate solo con el no leído más nuevo
iterateBufferedfromStreamChunked / .pullChunked()lote — emite las llegadas como una lista
iterateNextfromStreamNext / .pullNext()descarte por demanda — ignora lo que llegó sin un pull esperando

fromStream es el valor por defecto y el que usa la Demo 1, porque un fichero o un socket no deberían perder bytes. Usa latest cuando una UI solo le importa la lectura actual, chunked cuando el consumidor quiere trabajo en lotes, y next cuando los eventos viejos son peores que los huecos.

Demo 1 · fromStream y fxStream

Ambos envuelven un Stream.fromIterable para que puedas pasar un stream existente por los operadores de FxDart:

Demo 2 · Ida y vuelta, con un stream periódico finito

Stream.periodic nunca termina por sí solo, así que .take(n) mantiene la demo finita. La segunda mitad muestra la dirección inversa — construir una cadena FxAsync y devolverla hacia fuera como un Stream normal con .toStream():

Demo 3 · Cuatro formas de tirar de un stream

Un ráfaga síncrona de 1, 2, 3 llega mientras un pull ya está esperando. fromStream conserva cada valor; fromStreamLatest conserva solo el más nuevo; fromStreamChunked los emite como una lista; fromStreamNext conserva solo el valor que encontró el pull en espera. Las grafías de la cadena de eventos son .pull(), .pullLatest(), .pullChunked(), .pullNext().

Dos cadenas, un Stream

Un Stream es la única fuente que pertenece a las dos mitades de FxDart, así que lleva dos getters. No son variantes uno del otro: son modelos distintos.

stream.fxstream.fxEvents
devuelveFxAsync<T>FxEvents<T>
equivale afxStream(stream)fxEvents(stream)
modelopull — datos bajo demandapush — eventos en el tiempo
quién marca el ritmoel consumidor, un next() cada vezel stream; los operadores reajustan el tiempo
operadores típicosmap, filter, concurrent, toListdebounce, throttle, switchMap, combineLatest
contrapresiónsí — nada se tira hasta pedirlono — el stream emite cuando emite

La regla práctica: si la pregunta es «¿cuántos a la vez?» quieres la cadena pull, porque concurrent(n) solo significa algo cuando el consumidor controla la demanda. Si la pregunta es «¿con qué frecuencia, y cuál gana?» quieres la cadena de eventos. Empieza por cualquiera y cruza — .pull() / .pullLatest() / .pullChunked() / .pullNext() convierten un FxEvents en la cadena pull, .toStream() convierte un FxAsync de nuevo en un Stream.

Lecciones completas: fxEvents para la cadena push, fx para el modelo de cadena y las formas con getter, y concurrent para lo que solo el lado pull puede hacer.

Pruébalo tú

Ejercicio: quédate solo con los valores >= 10 de este stream.

Relacionado: toAsync — eleva un Iterable normal en su lugar · variantes asíncronas — la convención de nombres *Async · concurrent — aplícalo antes de toStream() para tener paralelismo real · concurrentPool — variante por orden de finalización · fxEvents — la cadena push; .pull() / .pullLatest() / .pullChunked() / .pullNext() cruzan de vuelta