Stream bridges
fromStream, fxStream y .toStream() — cruza libremente entre el Stream de Dart y el FxAsyncIterable de FxDart.
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 9 | FxDart | mientras estás ocupado |
|---|---|---|
iterateEach | fromStream / .pull() | FIFO sin pérdidas — pausa la fuente, encola cada valor |
iterateLatest | fromStreamLatest / .pullLatest() | descarta lo superado — quédate solo con el no leído más nuevo |
iterateBuffered | fromStreamChunked / .pullChunked() | lote — emite las llegadas como una lista |
iterateNext | fromStreamNext / .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.fx | stream.fxEvents | |
|---|---|---|
| devuelve | FxAsync<T> | FxEvents<T> |
| equivale a | fxStream(stream) | fxEvents(stream) |
| modelo | pull — datos bajo demanda | push — eventos en el tiempo |
| quién marca el ritmo | el consumidor, un next() cada vez | el stream; los operadores reajustan el tiempo |
| operadores típicos | map, filter, concurrent, toList | debounce, throttle, switchMap, combineLatest |
| contrapresión | sí — nada se tira hasta pedirlo | no — 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.
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