Skip to content

Streams — Effectful Sequences Over Time

Pull-based, chunked, backpressured Effect Streams — construction from every source, transformation with concurrency, consuming with Sinks, encoding channels (Ndjson & SchemaBinary), resource safety, and a paginated enrichment pipeline.

An Effect<A,E,R> produces one A. A Stream<A,E,R> produces many — lazily, on demand, with typed errors, services, interruptibility, and resource lifetimes baked in. Unlike a Promise-derived async iterable or a push Observable, a Stream is pull-based, chunked, and backpressured by default.

Model Who drives Buffering Backpressure Errors Resource scope
Promise<T> / async iterable producer pushes none ad-hoc (consumer must buffer) untyped throws manual try/finally
Observable (Rx) producer pushes unbounded subscription buffer operators untyped subscription cleanup
Array<T> fully eager whole array in memory none (materialized) not applicable none
Queue<A,E> channel (both) bounded policy offer suspends when full (see ch.15) typed E + Done queue lifecycle
Stream<A,E,R> consumer pulls chunked pull waits for consumer typed E Scope per pull/operator
Stream + Sink<A,...> consumer + sink chunked fold sink drives pulls typed Effect.scoped

This chapter is the Stream system: mental model, construction from every source the runtime supports, transformation with concurrency, consumption via the run* family and Sinks, chunk semantics, encoding channels (Ndjson, SchemaBinary), resource safety, pubsub interop, error handling, and a realistic paginated pipeline.

Mental model — pull, chunked, backpressured

Section titled “Mental model — pull, chunked, backpressured”

A Stream’s inner representation is a Channel (packages/effect/src/Stream.ts:93). Operationally:

  1. Pull. Each step requests the next Chunk<A> from upstream. Effect work between chunks is a suspension point — other fibers run, interruptions are observed, TestClock can move time.
  2. Chunked. Upstream emits NonEmptyArray<A> chunks; downstream maps run per element but the runtime batches transport. You choose logical batching (Chunk size), internals amortize effect overhead.
  3. Backpressured. A slow consumer simply does not pull — upstream work is not scheduled. No buffer(∞) hiding an OOM.
  4. Scoped. Each pull runs in a scope; callbacks, listeners, and platform handles are added as finalizers and tore down in LIFO order when the stream ends, fails, or is interrupted — same discipline as fibers (ch.13) and resources (ch.12).
Pull: consumer drives; upstream suspends until asked
Rendering diagram…
src/stream-pure.ts
import { Effect, Stream } from "effect"
const empty: Stream.Stream<never> = Stream.empty
const one: Stream.Stream<number> = Stream.succeed(42)
const few: Stream.Stream<number> = Stream.make(1, 2, 3)
const arr: Stream.Stream<number> = Stream.fromIterable([1, 2, 3])
const arrChunked: Stream.Stream<number> =
Stream.fromIterable([1, 2, 3], { chunkSize: 2 }) // logical chunk size hint
const ranged: Stream.Stream<number> = Stream.range(1, 100) // inclusive
const iterated: Stream.Stream<number> = Stream.iterate(1, (n) => n * 2).pipe(Stream.take(10))
const repeated: Stream.Stream<string> = Stream.make("hi").pipe(Stream.repeat(Schedule.recurs(4)))
src/stream-effects.ts
import { Effect, Schedule, Stream } from "effect"
declare const fetchUser: Effect.Effect<User, FetchError>
type User = { id: string }
class FetchError extends Error {}
// one value from an effect
const oneUser: Stream.Stream<User, FetchError> = Stream.fromEffect(fetchUser)
// one fallback value into an effect
const fromOption = Effect.succeed(42).pipe(Stream.fromEffect)
// repeating an effect
const polls: Stream.Stream<string, never> =
Stream.fromEffect(Effect.succeed("ping")).pipe(
Stream.repeat(Schedule.forever) // Stream repeats effect on its element schedule
)
// repeating with a dedicated schedule on an effect
const healthTicks: Stream.Stream<string, never> =
Stream.fromEffectSchedule(
Effect.succeed("ping"),
Schedule.spaced("30 seconds")
).pipe(Stream.take(3)) // => ["ping", "ping", "ping"] over 60s (virtual under TestClock)
// from pull effect — manual chunk-per-pull control
const fromPull: Stream.Stream<number> = Stream.fromPull(
Effect.succeed(() => Effect.succeed([[1, 2, 3]] as const)) // pull fn
)
// unwrap a lazy stream effect
const lazy: Stream.Stream<number, string> =
Stream.unwrap(Effect.succeed(Stream.make(1, 2, 3)))
// scoped resource-producing stream (resource released when stream ends)
const scopedStream: Stream.Stream<number> =
Stream.scoped(Effect.acquireRelease(
Effect.succeed(Stream.make(1, 2, 3)),
() => Effect.log("cleanup")
))
src/stream-paginate.ts
import { Effect, Option, Stream } from "effect"
interface Page { items: Array<Item>; nextCursor: string | null }
interface Item { id: string }
declare const fetchPage: (cursor: string | null) => Effect.Effect<Page, FetchError>
class FetchError extends Error {}
// paginate: seed S, each step → Effect< [chunk A, Option<next S>] >
const allItems: Stream.Stream<Item, FetchError> =
Stream.paginate(null as string | null, (cursor) =>
Effect.gen(function* () {
const page = yield* fetchPage(cursor)
const next: Option.Option<string | null> =
page.nextCursor === null ? Option.none() : Option.some(page.nextCursor)
return [page.items, next] as const
})
).pipe(Stream.flatMap((items) => Stream.fromIterable(items)))
// convenience when each API fetch returns exactly one "batch" that is already
// the desired element type:
const pagesAsValues: Stream.Stream<Array<Item>, FetchError> =
Stream.paginate(null as string | null, Effect.fn(function* (cursor) {
const page = yield* fetchPage(cursor)
const next = page.nextCursor === null ? Option.none() : Option.some(page.nextCursor)
return [[page.items] as Array<Array<Item>>, next] // chunk shape matches batch
}))
// when the page API itself is synchronous (no Effect), `unfold` covers it
const unfolded = Stream.unfold(0, (n) =>
n >= 5 ? Option.none() : Option.some([n, n + 1] as const)
)

From callbacks, listeners, and async iterables

Section titled “From callbacks, listeners, and async iterables”
src/stream-callback.ts
import { Effect, Queue, Stream } from "effect"
declare const element: HTMLElement
class LetterError extends Error { constructor(readonly cause: unknown) { super("letter") } }
// DOM / EventTarget
const clicks: Stream.Stream<PointerEvent> =
Stream.fromEventListener<PointerEvent>(element, "click")
// removal is scoped — listener detached when stream ends/interrupted
// signatures in packages/effect/src/Stream.ts:1405:
// fromEventListener<A>(target: EventTarget, event: string, options?: AddEventListenerOptions)
// manual callback source — acquireRelease semantics, Queue.offerUnsafe inside
const fromCallback: Stream.Stream<number> =
Stream.callback<number>((queue) =>
Effect.gen(function* () {
const handler = (n: number) => Queue.offerUnsafe(queue, n)
// setup
bus.on("tick", handler)
// teardown when stream ends or is interrupted
yield* Effect.addFinalizer(() => Effect.sync(() => bus.off("tick", handler)))
})
)
// async iterable (Node ReadableStream branches, third-party SDK paginators, etc.)
const fromAI: Stream.Stream<string, LetterError> =
Stream.fromAsyncIterable(
asyncIterable(),
(cause) => new LetterError({ cause })
)
async function* asyncIterable(): AsyncIterable<string> {
yield "a"; yield "b"
}
declare const bus: { on: (e: string, h: (n: number) => void) => void; off: (e: string, h: (n: number) => void) => void }

From queues, pubsubs, and platform handles

Section titled “From queues, pubsubs, and platform handles”
src/stream-queues-platform.ts
import { Effect, PubSub, Queue, Stream } from "effect"
import { NodeStream } from "@effect/platform-node"
import { Readable } from "node:stream"
import type { Scope } from "effect"
declare const letterQueue: Queue.Queue<string, LetterError>
class LetterError extends Error {}
const fromQ: Stream.Stream<string, LetterError | Cause.Done> =
Stream.fromQueue(letterQueue)
// Produced as channel — Done excluded? see Stream vs Queue note below
// correctly: Stream.fromQueue(q) : Stream<A, Exclude<E, Cause.Done>>
// Queue's Done is converted to stream termination, not a stream error.
const fromPs: Stream.Stream<string> =
Stream.fromPubSub(yield* PubSub.bounded<string>(1024))
// scoped subscription internally — ends when pubsub shuts down
// subscription variant (already subscribed)
declare const sub: PubSub.Subscription<string>
const fromSub: Stream.Stream<string> = Stream.fromSubscription(sub)
// Node Readable → Stream<A, UnknownError>
const nodeSource: Stream.Stream<Uint8Array, import("effect/Cause").UnknownError> =
NodeStream.fromReadable<Uint8Array, import("effect/Cause").UnknownError>({
evaluate: () => Readable.from(Buffer.from("hello")),
onError: (cause) => new (class extends Error { _tag = "UnknownError"; cause = cause })() as any,
closeOnDone: true // destroy Readable when stream ends (default true)
})
// signature (packages/platform/node-shared/src/NodeStream.ts):
// fromReadable<A = Uint8Array, E = UnknownError>(options: {
// evaluate: () => Readable | NodeJS.ReadableStream,
// onError?: (error: unknown) => E,
// chunkSize?: number, bufferSize?: number, closeOnDone?: boolean
// }): Stream<A, E>
// Scheduling-shaped sources
const ticks: Stream.Stream<void> = Stream.tick("1 seconds") // Stream<void> at interval
const scheduled: Stream.Stream<number> =
Stream.fromSchedule(Schedule.spaced("500 millis")) // Stream<Output> of schedule
import { Schedule } from "effect"
Constructor Signature sketch Resource note
Stream.empty Stream<never> —
Stream.succeed(a) / make(...a) / fromIterable(iter) value → stream sync
Stream.fromEffect(eff) Effect<A,E,R> → Stream<A,E,R> —
Stream.fromEffectSchedule(eff, schedule) Effect + Schedule → Stream schedule controls inter-element delay
Stream.fromEffectRepeat(eff) unbounded repeat chunked
Stream.paginate(s, f: S → Eff<[A[], Option<S>]>) paginated cursor tail-recursive
Stream.unfold(s, f) sync paginate —
Stream.fromAsyncIterable(iter, onError) AsyncIterable → Stream interruptible
Stream.fromEventListener(target, event) EventTarget → Stream scoped add/remove
Stream.callback(f: Queue → Eff<void,_,Scope>) callback → stream acquireRelease in f; offerUnsafe inside
Stream.fromQueue(queue) Queue → Stream ends on queue Done
Stream.fromPubSub(pubsub) PubSub → Stream scoped subscription
Stream.tick(interval) Duration.Input → Stream<void>
Stream.fromSchedule(schedule) Schedule → Stream<Output>
NodeStream.fromReadable({ evaluate, onError, closeOnDone }) Readable → Stream<A,E> scoped destroy
src/stream-maps.ts
import { Effect, Schedule, Stream } from "effect"
class EnrichError extends Error {}
const source = Stream.make(1, 2, 3, 4, 5, 6, 7, 8)
const mapped: Stream.Stream<string> = source.pipe(Stream.map((n) => `n=${n}`))
const filtered: Stream.Stream<number> = source.pipe(Stream.filter((n) => n % 2 === 0))
const filteredMap: Stream.Stream<string> = source.pipe(
Stream.filterMap((n) => n % 2 === 0 ? ({ _tag: "Some", value: `even:${n}` } as any) : ({ _tag: "None" } as any))
)
// effectful maps — with concurrency when appropriate
const mappedEffect: Stream.Stream<string, EnrichError> =
source.pipe(
Stream.mapEffect((n) => Effect.succeed(`enriched:${n}`), { concurrency: 4 })
)
// flattening — each element to a stream (bounded parallelism)
const flat: Stream.Stream<number> = source.pipe(
Stream.flatMap((n) => Stream.make(n, n * 2), { concurrency: 2 })
)
const switchMapped: Stream.Stream<number> = source.pipe(
Stream.switchMap((n) => Stream.make(n)) // latest wins, previous interrupted
)
// chunk-level (avoid per-element Effect overhead when possible)
const mapChunks: Stream.Stream<number> = source.pipe(
Stream.mapChunks((chunk) => chunk.map((n) => n * 2))
)
Operator Type Concurrency option
Stream.map(f) A → B —
Stream.mapEffect(f, { concurrency? }) A → Eff<B,E,R> number | "unbounded"
Stream.filter(pred) / filterMap / filterEffect predicate filterEffect supports concurrency
Stream.mapArray(f) / mapArrayEffect Array<A> → Array<B> (per-chunk) mapArrayEffect supports concurrency
Stream.flatMap(f, { concurrency? }) A → Stream<B> bounded merge
Stream.switchMap(f) A → Stream<B> last input’s stream wins
Stream.tap(f) / tapBoth side-effect —
Stream.scan(init, f) / scanEffect fold emitting intermediates —
Stream.grouped(n) / groupedWithin(n, duration) window into nested stream —
src/stream-take.ts
import { Stream } from "effect"
const source = Stream.range(1, 1_000)
const first10: Stream.Stream<number> = source.pipe(Stream.take(10))
const allBut3: Stream.Stream<number> = source.pipe(Stream.drop(3))
const until: Stream.Stream<number> = source.pipe(Stream.takeUntil((n) => n === 42))
const whil: Stream.Stream<number> = source.pipe(Stream.takeWhile((n) => n < 100))
const effectWhile: Stream.Stream<number> = source.pipe(
Stream.takeWhileEffect((n) => Effect.succeed(n < 100))
)
const droppedWhile: Stream.Stream<number> = source.pipe(Stream.dropWhile((n) => n < 5))
// by schedule / duration / predicate effect
const timed: Stream.Stream<number> = source.pipe(Stream.timeout("5 seconds"))
const grouped: Stream.Stream<Stream.Stream<number>> = source.pipe(Stream.grouped(10))
const groupedTimed = source.pipe(Stream.groupedWithin(100, "500 millis"))
import { Effect } from "effect"

take/drop are pull-aware: upstream not yet produced beyond the boundary is never executed — time- and resource-safe.

src/stream-concurrency.ts
import { Effect, Stream } from "effect"
const ids = Stream.range(1, 50)
const enriched = ids.pipe(
Stream.mapEffect(
(id) => Effect.gen(function* () {
yield* Effect.sleep("20 millis")
return { id, name: `user:${id}` }
}),
{ concurrency: 8 }
)
)
// flatMap parallelism — encloses Stream producing effectful inner streams
const withInner: Stream.Stream<string> = Stream.range(1, 20).pipe(
Stream.flatMap(
(n) => Stream.fromEffect(Effect.succeed(`row:${n}`)),
{ concurrency: 4 }
)
)
// ordering note: concurrency-merged results are emitted as inner streams complete;
// input order is NOT preserved under unbounded concurrency. Bound with ordering
// requirements by post-sorting or by sequential mapping outside the concurrency window.

A Stream does nothing until run. Different collectors turn its emission into an Effect.

src/stream-consume.ts
import { Effect, Option, Sink, Stream } from "effect"
const source: Stream.Stream<number, string> = Stream.make(1, 2, 3)
const collected: Effect.Effect<Array<number>, string> = Stream.runCollect(source)
// => Effect<[1,2,3]>
const drained: Effect.Effect<void, string> = Stream.runDrain(source)
// discard all, run for effects only
const first: Effect.Effect<Option.Option<number>, string> = Stream.runHead(source)
// => Some(1)
const last: Effect.Effect<Option.Option<number>, string> = Stream.runLast(source)
const counted: Effect.Effect<number, string> = Stream.runCount(source)
const summed: Effect.Effect<number, string> = Stream.runSum(source as Stream.Stream<number, string>)
const folded: Effect.Effect<number, string> =
Stream.runFold(source, 0, (acc, n) => acc + n)
const forEach: Effect.Effect<void, string> =
Stream.runForEach(source, (n) => Effect.log(`saw ${n}`))
const forEachWhile: Effect.Effect<void, string> =
Stream.runForEachWhile(source, (n) => Effect.succeed(n < 2)) // stop after predicate fails
const str: Effect.Effect<string, string> =
Stream.mkString(source.pipe(Stream.map(String)))
// sinking — Stream.run(sink) for custom reductions and leftover-aware folds
const sinkSum: Sink.Sink<number, number> = Sink.sum // Sink<Out, In, Leftover, Err, Env>
const viaSink: Effect.Effect<number, string> = Stream.run(source, sinkSum)
const sinkHead: Sink.Sink<Option.Option<number>, number> = Sink.head()
const headViaSink: Effect.Effect<Option.Option<number>, string> = Stream.run(source, sinkHead)
// conversions
const asAsyncIterable: Effect.Effect<AsyncIterable<number>, never, never> =
Stream.toAsyncIterableEffect(source)
const pull: Effect.Effect<() => Effect.Effect<import("effect/Chunk").Chunk<number>>> =
Stream.toPull(source) // manual pull loop — advanced, scopes handled via the returned effect's scope
Consumer Result Effect
Stream.runCollect Effect<Array<A>, E, R>
Stream.runDrain Effect<void, E, R>
Stream.runHead / runLast Effect<Option<A>, E, R>
Stream.runCount / runSum Effect<number, E, R>
Stream.runFold(init, f) / runFoldEffect fold → Effect<S,E,R>
Stream.runForEach(f) / runForEachWhile iterate with Effect
Stream.runForEachArray(f) chunk-batched iteration
Stream.run(sink) run through Sink<Out, In, L, E, R>
Stream.mkString / mkArrayBuffer / mkUint8Array collect bytes/text
Stream.toReadableStream / toReadableStreamEffect expose as Web ReadableStream
Stream.toAsyncIterableEffect expose as JS async iterable

Sink is a Channel that consumes chunks and produces a value (and optional leftovers). Built-ins in packages/effect/src/Sink.ts:

Sink Shape
Sink.succeed(a) / fail(e) / never constant
Sink.head() / last() first / last element
Sink.take(n) Array<In> of n
Sink.sum number (reduce(0, +))
Sink.count number
Sink.drain void
Sink.reduce(init, f) generic fold
Sink.every(pred) / some(pred) boolean fold
Sink.collectArray / collectAll gather

Encoding channels — Ndjson & SchemaBinary pipelines

Section titled “Encoding channels — Ndjson & SchemaBinary pipelines”

Streams carry typed values; bytes on the wire need codecs. Effect provides Channel-based codecs in effect/unstable/encoding that plug into streams with Stream.pipeThroughChannel (or Stream.pipeThrough). The codecs are streaming-safe — framing handles split/merged chunks and ignores incomplete trailing splits appropriately.

src/stream-ndjson.ts
import { Effect, Schema, Stream } from "effect"
import { Ndjson } from "effect/unstable/encoding"
const Person = Schema.Struct({ name: Schema.String, age: Schema.Number })
type Person = typeof Person.Type
// simulated source: UTF-8 lines over a Readable or an async iterable
declare const byteSource: Stream.Stream<Uint8Array, never>
declare const stringSource: Stream.Stream<string, never>
// decode → validate → work → re-encode pipeline
const decoded: Stream.Stream<Person, import("effect/Schema").SchemaError | Ndjson.NdjsonError> =
byteSource.pipe(
// byte → string per TextDecoder
Stream.decodeText(),
// string chunk handling splits newlines (splitLines internal)
// ndjson decode schemaString: string → unknown JSON per line → Schema typed Person
Stream.pipeThroughChannel(Ndjson.decodeSchemaString(Person)({ ignoreEmptyLines: true }))
)
// filter + re-encode (string variant)
const adultsEncoded: Stream.Stream<string, typeof Person.Encoded> =
decoded.pipe(
Stream.filter((p) => p.age >= 18),
Stream.pipeThroughChannel(Ndjson.encodeSchemaString(Person)())
)
// bytes variant — encode/decode operate on Uint8Array ndjson
const adultBytes: Stream.Stream<Uint8Array, import("effect/Schema").SchemaError | Ndjson.NdjsonError> =
decoded.pipe(
Stream.filter((p) => p.age >= 18),
Stream.pipeThroughChannel(Ndjson.encodeSchema(Person)())
)
// string-only ndjson (no byte transcoding)
const stringOnly: Stream.Stream<Person, import("effect/Schema").SchemaError | Ndjson.NdjsonError> =
stringSource.pipe(
Stream.pipeThroughChannel(Ndjson.decodeSchemaString(Person)())
)
// error handling — NdjsonError carries kind: "Pack" | "Unpack" and cause
const resilient: Stream.Stream<Person, never> =
byteSource.pipe(
Stream.decodeText(),
Stream.pipeThroughChannel(Ndjson.decodeSchemaString(Person)()),
Stream.catchTag("NdjsonError", (e) =>
Stream.fromEffect(Effect.logWarning(`ndjson malformed: ${e.kind}`)).pipe(
Stream.flatMap(() => Stream.empty)
)
),
Stream.catchTag("SchemaError", () => Stream.empty) // drop malformed rows
)
// splitLines / decodeText: Stream-stage helpers for text framing
const lines: Stream.Stream<string> = byteSource.pipe(
Stream.decodeText(), // Uint8Array → string (streaming TextDecoder)
Stream.splitLines // split on \n, handles chunk boundaries
)
Ndjson channel constructor Direction Input Stream element Output element Errors
Ndjson.decodeString(opts?) Channel<string → unknown> string (line chunks) unknown per line NdjsonError
Ndjson.decode(opts?) Channel<Uint8Array → unknown> Uint8Array (UTF-8) unknown NdjsonError
Ndjson.decodeSchemaString(S) Channel<string → S.Type> string S.Type NdjsonError | SchemaError
Ndjson.decodeSchema(S) Channel<Uint8Array → S.Type> Uint8Array S.Type NdjsonError | SchemaError
Ndjson.decodeSchemaString(S)({ ignoreEmptyLines: true }) same same same skips empty lines before JSON parse
Ndjson.encodeString() Channel<unknown → string> unknown trailing-newline string NdjsonError
Ndjson.encode() Channel<unknown → Uint8Array> unknown UTF-8 bytes NdjsonError
Ndjson.encodeSchemaString(S)() Channel<S.Type → string> S.Type trailing-newline string SchemaError | NdjsonError
Ndjson.encodeSchema(S)() Channel<S.Type → Uint8Array> S.Type bytes SchemaError | NdjsonError

SchemaBinary — compact row-oriented binary frames

Section titled “SchemaBinary — compact row-oriented binary frames”

For density and schema-evolvable wire format, use SchemaBinary channels. The binary protocol is array-of-structs aware: a batch of records declares its shape once; repeated strings are dictionary-compressed across the batch.

src/stream-schemabinary.ts
import { Effect, Schema, Stream } from "effect"
import { SchemaBinary } from "effect/unstable/encoding"
const Event = Schema.Struct({
id: Schema.String.pipe(SchemaBinary.fieldId(1)),
score: Schema.Number.pipe(SchemaBinary.fieldId(2))
})
type Event = typeof Event.Type
// Stream<Event> → Stream<Uint8Array> (one frame per upstream chunk, many rows per frame)
const asBytes: Stream.Stream<Uint8Array, import("effect/Schema").SchemaError> =
Stream.make({ id: "a", score: 10 }, { id: "b", score: 20 }).pipe(
Stream.pipeThroughChannel(SchemaBinary.encode(Event)())
)
// Uint8Array → Stream<Event> (framings may split/merge across chunks transparently)
declare const byteFrames: Stream.Stream<Uint8Array, never>
const asEvents: Stream.Stream<Event, import("effect/Schema").SchemaError> =
byteFrames.pipe(
Stream.pipeThroughChannel(SchemaBinary.decode(Event)())
)
// decode → transform → re-encode — the realistic hop
const repack: Stream.Stream<Uint8Array, import("effect/Schema").SchemaError> =
asEvents.pipe(
Stream.filter((e) => e.score > 5),
Stream.pipeThroughChannel(SchemaBinary.encode(Event)())
)
// combining codecs on a Duplex channel (send typed values both ways)
// See `NodeStream.pipeThroughDuplex` for Node; `HttpClient` wraps similar for HTTP
Channel Direction Type header
SchemaBinary.encode(S)(opts?) S.Type → Uint8Array encoded side wire layout
SchemaBinary.decode(S)(opts?) Uint8Array → S.Type decoded side

SchemaBinary also exposes toCodec(S) for one-frame sync encodeUnknownSync/decodeUnknownSync outside streaming contexts. Wire options: fingerprint: true (compact positional layout with hash; smaller but peers must match schema exactly) vs default (self-describing row shape for evolution).

Resource safety — scoped callbacks and channels

Section titled “Resource safety — scoped callbacks and channels”

Stream’s callback source is resource-safe by construction. The constructor takes a function whose body runs in Scope:

src/stream-resource.ts
import { Effect, Queue, Stream } from "effect"
const riskySource: Stream.Stream<number> =
Stream.callback<number>((queue) =>
Effect.gen(function* () {
const handle = scheduleEvery(100, (n: number) => Queue.offerUnsafe(queue, n))
// associate handle teardown with this stream execution's lifetime
yield* Effect.addFinalizer(() => Effect.sync(() => clearHandle(handle)))
yield* Effect.log("source started")
})
)
// handle cleared when: stream ends (natural completion), stream fails, or downstream
// is interrupted (timeout, race loser, fiber interrupt). In all cases, LIFO finalizers unwind.
function scheduleEvery(ms: number, f: (n: number) => void) {
let n = 0
const id = setInterval(() => f(n++), ms)
return id as unknown as object
}
function clearHandle(handle: object) { clearInterval(handle as unknown as number) }
// for ad-hoc lifecycle outside `callback`, use acquireRelease scoped streams
const fromFileLike = Stream.scoped(
Effect.acquireRelease(
Effect.sync(() => ({ read: () => "line", close: () => {} })),
(res) => Effect.sync(() => res.close())
).pipe(Effect.map((res) => Stream.fromIterable([res.read()])))
)

Any Channel pulled into a stream — Channel.fromTransform, Channel.decodeText, Channel.splitLines, Ndjson.*, SchemaBinary.*, NodeStream.fromReadableChannel — inherits the same scope discipline: their finalizers close, flush, and detach on termination.

Where queues distribute work (one item to one worker) and pubsubs broadcast information (one item to all subscribers), streams express processing over either:

src/stream-pubsub-queue.ts
import { Effect, PubSub, Queue, Stream } from "effect"
declare const queue: Queue.Queue<Task, TaskError>
interface Task { id: string }
class TaskError extends Error {}
// one consumer at a time → many consumers each with their own Stream.fromQueue
const fromQForEachWorker = Stream.fromQueue(queue).pipe(
Stream.runForEach((t) => Effect.log(`work ${t.id}`))
)
// fan-out with the same bus as in ch.15's OrderEvents service
declare const pubsub: PubSub.PubSub<OrderEvent>
type OrderEvent = { orderId: string }
const fromPS: Stream.Stream<OrderEvent> = Stream.fromPubSub(pubsub)
// the stream's subscription is scoped: subscribe on pull start, unsubscribe on end/failure/interrupt
// `replay: N` in PubSub construction still applies — late subscribers see buffered history (see ch.15)
// streams also feed back into pubsubs/queues as sinks
const toQueue: Effect.Effect<void, never> =
Stream.fromIterable([1, 2, 3]).pipe(Stream.runIntoQueue(queue, { shutdown: true }))
// runIntoQueue / runIntoPubSub / toPubSub / toQueue are first-class
const completion: Effect.Effect<void> =
Stream.fromIterable([1, 2, 3]).pipe(
Stream.tap((n) => Effect.log(`enqueued ${n}`)),
Stream.runIntoQueue(queue, { shutdown: true }) // end() after drain
)

Stream.fromQueue(queue).pipe(Stream.make(1,2,3), Stream.runIntoQueue(...)) round-trips are zero-convention: producing and consuming under stream scopes gives you backpressure-composable pipelines over durable shared state when you need it, and purely stream-internal pipelines when you don’t.

Error handling — catch at the stream level

Section titled “Error handling — catch at the stream level”

Stream exposes catch* combinators mirroring Effect — without needing to interrupt iteration:

src/stream-errors.ts
import { Effect, Schema, Stream } from "effect"
class TransientError extends Schema.TaggedError<TransientError>()(
"TransientError", { status: Schema.Number }
) {}
class FatalError extends Schema.TaggedError<FatalError>()("FatalError", { reason: Schema.String }) {}
declare const source: Stream.Stream<number, TransientError | FatalError>
// per-tag recovery — drop or replace failed branches
const recovered: Stream.Stream<number, FatalError> =
source.pipe(
Stream.catchTag("TransientError", (e) =>
Effect.logWarning(`transient ${e.status}, retrying`).pipe(
Stream.fromEffect, Stream.flatMap(() => Stream.empty)
)
)
)
// multi-tag and predicate forms
const multi: Stream.Stream<number, FatalError> =
source.pipe(Stream.catchTags({
TransientError: (e) => Stream.empty // absorb
}))
// when diagnoses share a TaggedError union, use reason-typed catching
const reason: Stream.Stream<number, never> =
source.pipe(Stream.catchReason("TransientError", (r) => Stream.succeed(-1)))
// generic catch-all / cause-level
const causeAware = source.pipe(
Stream.catchCause((cause) => Stream.succeed(-1))
)
// combine with queue error-stashes or dead-letter pattern from ch.15
Combinator Catches
Stream.catch(f) any E
Stream.catchFilter(pred, f) filtered
Stream.catchTag("Tag", f) Data.TaggedError _tag
Stream.catchTags({ Tag: f, ... }) many tags → never (or residual E)
Stream.catchReason("Union", reason, f) Cause reason-typed
Stream.catchReasons("Union", { Reason: f }) many reasons
Stream.catchCause(f) / catchCauseIf full Cause<E>

In v4 these derive from Filter + scheduling — the same machinery behind Effect.catch*.

Choosing: Stream vs Queue vs Array vs async iterable

Section titled “Choosing: Stream vs Queue vs Array vs async iterable”

A real program uses all four. Use the smallest abstraction that carries the needed guarantees:

Decision: which collection-like thing?
Rendering diagram…
Criterion Prefer
Whole collection in memory, no time, no concurrency Array / Chunk (plain)
Map A → Effect<B> with concurrency, then collect Effect.forEach(items, f, { concurrency })
Producer outpaces consumer (or vice versa) with bounded memory Queue (bounded suspend / dropping / sliding) + consumers; pull as Stream.fromQueue
Broadcast PubSub; consume as Stream.fromPubSub when pipeline combinators needed
Time-shaped source (tick, schedule, timeout), composable windowing, encoding channels Stream
Need filtering/mapEffect-then-consume-as-pipeline Stream (operators compose regardless of source)
Interop with outside world (Readable, EventTarget, async SDK) Stream constructors (NodeStream.fromReadable, fromEventListener, fromAsyncIterable)
Need leftover-aware folding Sink via Stream.run(sink) (Sink.head, take, sum, …)

Realistic pipeline — paginated API → stream → enrichment → filter → collect

Section titled “Realistic pipeline — paginated API → stream → enrichment → filter → collect”

Put the whole surface together. Spec: the platform exposes a cursor-paginated listing API (GET /items?cursor=…) returning Array<Listed> plus nextCursor. Each listed item is then enriched through a rate-limited lookup effect; only qualifying items are kept. The pipeline must paginate until exhaustion, enrich with bounded concurrency, filter, and materialize.

src/pipeline.ts
import { Effect, Option, Schedule, Schema, Stream } from "effect"
// ── domain ──
const Listed = Schema.Struct({
id: Schema.String,
kind: Schema.String
})
type Listed = typeof Listed.Type
const Enriched = Schema.Struct({
id: Schema.String,
kind: Schema.String,
score: Schema.Number
})
type Enriched = typeof Enriched.Type
class ApiError extends Schema.TaggedError<ApiError>()(
"ApiError", { status: Schema.Number, message: Schema.String, retryable: Schema.Boolean }
) {}
class EnrichError extends Schema.TaggedError<EnrichError>()(
"EnrichError", { id: Schema.String, message: Schema.String }
) {}
// ── external effects — would be HttpClient or platform fetch in real code ──
declare const listPage: (cursor: string | null) => Effect.Effect<
{ items: ReadonlyArray<Listed>; nextCursor: string | null },
ApiError
>
declare const enrich: (item: Listed) => Effect.Effect<Enriched, EnrichError>
// ── schedule reused for both pagination retries and enrich retries ──
const apiRetry = Schedule.min([
Schedule.exponential("200 millis"),
Schedule.spaced("4 seconds")
]).pipe(Schedule.jittered, Schedule.setInputType<ApiError>(), Schedule.while(({ input }) => input.retryable))
const retryListPage = (cursor: string | null): Effect.Effect<ReturnType<typeof listPage> extends Effect.Effect<infer A, any, any> ? A : never, ApiError> =>
listPage(cursor).pipe(Effect.retry(apiRetry))
// ── streams ──
// page cursor is `string | null`; Option models exhaustion
const listedStream: Stream.Stream<Listed, ApiError> =
Stream.paginate(null as string | null, Effect.fn(function* (cursor) {
const page = yield* retryListPage(cursor)
const next = page.nextCursor === null ? Option.none<string | null>() : Option.some(page.nextCursor)
return [page.items, next] as const
})).pipe(
// Each paginate step emitted a chunk `Array<Listed>` — flatten to per-element stream
Stream.flatMap((items) => Stream.fromIterable(items))
)
// Enrichment that retries transient fetch errors but fails through non-retryable ones
const enrichedStream: Stream.Stream<Enriched, ApiError | EnrichError> =
listedStream.pipe(
Stream.mapEffect(
(item) =>
enrich(item).pipe(
// Enrichment has its own local retry for transient lookup failures
Effect.retry(
Schedule.exponential("100 millis").pipe(
Schedule.upTo({ times: 3 }),
Schedule.setInputType<EnrichError>()
// enrichment here: assume all EnrichError are retryable — tighten with while() if not
)
)
),
{ concurrency: 4 } // ← bound parallelism to the lookup service's rate limit
)
)
// Filter to qualified records, then materialize
const qualifiedPipeline: Effect.Effect<Array<Enriched>, ApiError | EnrichError> =
enrichedStream.pipe(
Stream.filter((e) => e.score > 42),
Stream.take(100), // guardrail: never materialize unboundedly in this report context
Stream.runCollect
)
// ── execution entrypoint ──
export const runPipeline = Effect.fn("runPipeline")(function* () {
const qualified = yield* qualifiedPipeline
yield* Effect.log(`pipeline complete: ${qualified.length} qualified of at most 100`)
return qualified
})
// ── virtual-time test — no sleeps hit the wall clock ──
import { TestClock } from "effect/testing"
import { it, expect } from "vitest"
it("paginate→enrich→filter collects under virtual time", async () => {
// stub implementations via Effect.provideService would go here in real test
// listPage/enrich supplied by a Test layer returning deterministic pages
const program = Effect.gen(function* () {
const fiber = yield* Effect.fork(runPipeline)
// exponential backoffs for apiRetry (200ms, 400ms, ...) are virtual —
// advance past enough ticks that pagination + enrichment retries resolve
yield* TestClock.adjust("10 seconds")
return yield* fiber.await
})
const result = await Effect.runPromise(Effect.provide(program, TestClock.layer()))
expect(Array.isArray(result)).toBe(true)
})
// bounded and alternative sinks
const drained: Effect.Effect<void, ApiError | EnrichError> =
enrichedStream.pipe(Stream.runDrain)
  • paginate over fromEffectSchedule. Pagination cursors are data-driven — the next cursor is inside the response, not a clock. paginate threads Option<S> cursor state; fromEffectSchedule drives on ticks alone.
  • flatMap(...fromIterable) after paginate. paginate steps in pages (Array<Listed> is one element per step). Consumers usually want elements; flattening at the boundary makes later mapEffect per-item reasoning natural. If page-level batching matters downstream, defer flattening.
  • mapEffect with concurrency: 4. Lookup is I/O-bound; 4 is a plausible per-route limit without collapsing to serial. Adjust with the same Semaphore discipline from ch.13 if the enrichment shares a pool with other code.
  • retry inside mapEffect, not around the stream. The lookup wrapped individually retries per-item transient failures; a retry around the whole stream would re-paginate and duplicate successfully enriched records. Scope failure handling to the element.
  • take(100) before runCollect. Materializing an unbounded remote source is how you OOM. Bounding at consumption site is the discipline — treat unconstrained runCollect like unbounded Promise.all and justify it.
  • TestClock.adjust drives the whole sleep fabric. Both listPage backoffs and the enrichment exponential("100 millis") retries are Clock.sleep underneath; moving virtual time once flushes both without real delays. Where paginate is cursor-driven, TestClock matters only for the retry delays, not for cursor advancement.