Hyperlinkv0.8.0-beta.28

Stream

Stream.fromQueueconsteffect/Stream.ts:1293
<A, E>(queue: Queue.Dequeue<A, E>): Stream<A, Exclude<E, Cause.Done>>

Creates a stream that pulls values from a Queue.Dequeue.

Details

The stream emits non-empty batches of queued values and ends when the queue fails with Cause.Done; other queue failures are propagated.

Example (Creating a stream from a queue of values)

import { Console, Effect, Queue, Stream } from "effect"

const program = Effect.gen(function*() {
  const queue = yield* Queue.unbounded<number>()
  yield* Queue.offer(queue, 1)
  yield* Queue.offer(queue, 2)
  yield* Queue.offer(queue, 3)
  yield* Queue.shutdown(queue)

  const stream = Stream.fromQueue(queue)
  const values = yield* Stream.runCollect(stream)
  yield* Console.log(values)
})

Effect.runPromise(program)
// Output: [ 1, 2, 3 ]
constructors
Source effect/Stream.ts:12932 lines
export const fromQueue = <A, E>(queue: Queue.Dequeue<A, E>): Stream<A, Exclude<E, Cause.Done>> =>
  fromChannel(Channel.fromQueueArray(queue))
Referenced by 2 symbols