Stream 브리지

fromStream, fxStream, .toStream() — Dart의 Stream과 FxDart의 FxAsyncIterable 사이를 자유롭게 오갑니다.

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

강의

fromStream은 단일 구독이든 브로드캐스트든 어떤 Stream이라도 FxAsyncIterable로 바꿔 줍니다. 덕분에 소켓이나 파일, 위젯의 이벤트 스트림 등 Dart가 Stream으로 건네주는 모든 데이터 위에서 FxDart의 연산자 전체(map, filter, concurrent, …)를 쓸 수 있습니다. fxStream(stream)도 같은 일을 하지만, 날것의 FxAsyncIterable 대신 체이닝 가능한 FxAsync를 바로 반환합니다 — fxfxAsync에 대응하는 비동기 버전인 셈입니다.

반대 방향으로는 .toStream()FxAsyncIterable(또는 FxAsync 체인)을 끝까지 돌리면서 그 값들을 평범한 Stream으로 다시 내보냅니다 — 다른 API(예컨대 StreamBuilder)가 이를 요구할 때 유용합니다. 한 가지 주의할 점은 toStream()이 언제나 순차적으로 값을 끌어당기며, 상류의 concurrent(n)을 무시한다는 것입니다 — 병렬 처리가 실제로 일어나길 원한다면 concurrentconcurrentPool을 체인에 먼저 적용한 다음 .toStream()을 호출하세요. 스트림 변환 자체는 병렬성을 더해 주지 않습니다.

푸시에서 풀로 건너오는 것은 연산 하나가 아닙니다 — 넷입니다. 소비자가 바쁜 동안에도 스트림은 계속 내보낼 수 있기 때문입니다. RxJS 9는 넷을 iterateEach, iterateLatest, iterateBuffered, iterateNext로 부릅니다. FxDart는 이를 fromStream*(날것 iterable)과 FxEvents.pull*(체인)에 대응시킵니다:

RxJS 9FxDart바쁜 동안
iterateEachfromStream / .pull()무손실 FIFO — 소스를 일시정지하고 모든 값을 큐에 담음
iterateLatestfromStreamLatest / .pullLatest()낡은 값 버리기 — 읽히지 않은 최신 값만 유지
iterateBufferedfromStreamChunked / .pullChunked()배치 — 도착분을 리스트 하나로 내보냄
iterateNextfromStreamNext / .pullNext()수요 게이트로 버리기 — 기다리는 pull이 없을 때 온 것은 무시

fromStream이 기본값이자 데모 1이 쓰는 것이고, 파일이나 소켓은 바이트를 잃으면 안 되기 때문입니다. 최신만 쓰는 UI에는 latest, 묶음으로 일하고 싶으면 chunked, 낡은 이벤트가 공백보다 나쁠 때는 next를 고르세요.

데모 1 · fromStream과 fxStream

둘 다 Stream.fromIterable을 감싸므로, 기존 스트림을 FxDart 연산자에 그대로 흘려보낼 수 있습니다.

데모 2 · 유한한 주기 스트림으로 왕복하기

Stream.periodic은 스스로 끝나지 않으므로 .take(n)으로 데모를 유한하게 만듭니다. 후반부는 반대 방향을 보여 줍니다 — FxAsync 체인을 구성한 다음 .toStream()으로 다시 평범한 Stream으로 내보내는 것입니다.

데모 3 · 스트림을 당기는 네 가지 방법

이미 pull이 기다리는 동안 1, 2, 3의 동기 버스트가 도착합니다. fromStream은 모든 값을 지키고, fromStreamLatest는 최신만 지키고, fromStreamChunked는 리스트 하나로 내보내고, fromStreamNext는 기다리던 pull을 만난 값만 지킵니다. 이벤트 체인 표기는 .pull(), .pullLatest(), .pullChunked(), .pullNext()입니다.

하나의 Stream, 두 개의 체인

Stream은 FxDart의 두 세계 모두에 속하는 유일한 소스라서 getter가 둘 붙습니다. 서로의 변종이 아니라 모델이 다릅니다.

stream.fxstream.fxEvents
결과FxAsync<T>FxEvents<T>
함수 표기fxStream(stream)fxEvents(stream)
모델pull — 수요에 따라 흐르는 데이터push — 시간 위에 놓인 이벤트
속도를 정하는 쪽소비자 — next() 한 번에 하나스트림 — 연산자는 타이밍을 다듬을 뿐
주로 쓰는 연산자map, filter, concurrent, toListdebounce, throttle, switchMap, combineLatest
배압(backpressure)있음 — 요청하기 전에는 당기지 않음없음 — 스트림은 나올 때 나옴

고르는 기준은 이렇습니다. 질문이 “한 번에 몇 개씩?”이면 pull 체인입니다. concurrent(n)은 소비자가 수요를 쥐고 있을 때만 의미가 있기 때문입니다. 질문이 “얼마나 자주, 그리고 어느 것이 이기나?”라면 이벤트 체인입니다. 어느 쪽에서 시작하든 건너올 수 있습니다 — .pull() / .pullLatest() / .pullChunked() / .pullNext()FxEvents를 pull 체인으로, .toStream()FxAsync를 다시 Stream으로 바꿉니다.

자세한 수업: push 체인은 fxEvents, 체인 모델과 getter 표기는 fx, pull 쪽에서만 가능한 일은 concurrent에 있습니다.

직접 해 보기

연습: 이 스트림에서 10 이상인 값만 남겨 보세요.

관련 항목: toAsync — 스트림 대신 평범한 Iterable을 끌어올리기 · 비동기 변형 — *Async 명명 규칙 · concurrent — 실제 병렬 처리를 원하면 toStream() 이전에 적용 · concurrentPool — 완료 순서 방식의 변형 · fxEvents — 푸시 체인; .pull() / .pullLatest() / .pullChunked() / .pullNext()로 되돌아옴