<A, E>(
iterable: AsyncIterable<A>,
onError: (error: unknown) => E
): Stream<A, E>Creates a stream from an AsyncIterable.
Example (Creating a stream from an AsyncIterable)
import { Data, Effect, Stream } from "effect"
class StreamError extends Data.TaggedError("StreamError")<{ readonly cause: unknown }> {}
const iterable = (async function*() {
yield 1
yield 2
yield 3
})()
Effect.runPromise(Effect.gen(function*() {
const stream = Stream.fromAsyncIterable(iterable, (cause) => new StreamError({ cause }))
const values = yield* Stream.runCollect(stream)
yield* Effect.sync(() => console.log(values))
}))
// [ 1, 2, 3 ]export const const fromAsyncIterable: <A, E>(
iterable: AsyncIterable<A>,
onError: (error: unknown) => E
) => Stream<A, E>
Creates a stream from an AsyncIterable.
Example (Creating a stream from an AsyncIterable)
import { Data, Effect, Stream } from "effect"
class StreamError extends Data.TaggedError("StreamError")<{ readonly cause: unknown }> {}
const iterable = (async function*() {
yield 1
yield 2
yield 3
})()
Effect.runPromise(Effect.gen(function*() {
const stream = Stream.fromAsyncIterable(iterable, (cause) => new StreamError({ cause }))
const values = yield* Stream.runCollect(stream)
yield* Effect.sync(() => console.log(values))
}))
// [ 1, 2, 3 ]
fromAsyncIterable = <function (type parameter) A in <A, E>(iterable: AsyncIterable<A>, onError: (error: unknown) => E): Stream<A, E>A, function (type parameter) E in <A, E>(iterable: AsyncIterable<A>, onError: (error: unknown) => E): Stream<A, E>E>(
iterable: AsyncIterable<A>iterable: interface AsyncIterable<T, TReturn = any, TNext = any>AsyncIterable<function (type parameter) A in <A, E>(iterable: AsyncIterable<A>, onError: (error: unknown) => E): Stream<A, E>A>,
onError: (error: unknown) => EonError: (error: unknownerror: unknown) => function (type parameter) E in <A, E>(iterable: AsyncIterable<A>, onError: (error: unknown) => E): Stream<A, E>E
): interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E>(iterable: AsyncIterable<A>, onError: (error: unknown) => E): Stream<A, E>A, function (type parameter) E in <A, E>(iterable: AsyncIterable<A>, onError: (error: unknown) => E): Stream<A, E>E> => const fromChannel: <
Arr extends Arr.NonEmptyReadonlyArray<any>,
E,
R
>(
channel: Channel.Channel<
Arr,
E,
void,
unknown,
unknown,
unknown,
R
>
) => Stream<
Arr extends Arr.NonEmptyReadonlyArray<infer A>
? A
: never,
E,
R
>
Creates a stream from a array-emitting Channel.
Example (Creating a stream from an array-emitting channel)
import { Channel, Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const channel = Channel.succeed([1, 2, 3] as const)
const stream = Stream.fromChannel(channel)
const result = yield* Stream.runCollect(stream)
yield* Console.log(result)
})
// Output: [ 1, 2, 3 ]
fromChannel(import ChannelChannel.const fromAsyncIterableArray: <A, D, E>(
iterable: AsyncIterable<A, D>,
onError: (error: unknown) => E
) => Channel<Arr.NonEmptyReadonlyArray<A>, E, D>
Creates a channel from an AsyncIterable, emitting each yielded value as a
single-element non-empty array.
Details
The iterator's return value becomes the channel's done value. Thrown or
rejected iterator errors are converted with onError. If the channel scope
closes early and the iterator has a return method, that method is called.
fromAsyncIterableArray(iterable: AsyncIterable<A>iterable, onError: (error: unknown) => EonError))