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:
- Pull. Each step requests the next
Chunk<A>from upstream. Effect work between chunks is a suspension point — other fibers run, interruptions are observed,TestClockcan move time. - Chunked. Upstream emits
NonEmptyArray<A>chunks; downstream maps run per element but the runtime batches transport. You choose logical batching (Chunksize), internals amortize effect overhead. - Backpressured. A slow consumer simply does not pull — upstream work is not scheduled. No
buffer(∞)hiding an OOM. - 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).
sequenceDiagram participant C as Consumer (runCollect) participant S as Stream / Channel participant P as Producer (paginate / callback) C->>S: pull Chunk S->>P: request next P-->>S: Chunk[A] (1..n items) S-->>C: emit Note over P,S: P suspends until next pull — no buffered backlog C->>S: pull Chunk S->>P: request next P-->>S: Done S-->>C: end
Creating streams — the full catalog
Section titled “Creating streams — the full catalog”From pure values
Section titled “From pure values”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 hintconst ranged: Stream.Stream<number> = Stream.range(1, 100) // inclusiveconst 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)))From effects
Section titled “From effects”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 effectconst oneUser: Stream.Stream<User, FetchError> = Stream.fromEffect(fetchUser)
// one fallback value into an effectconst fromOption = Effect.succeed(42).pipe(Stream.fromEffect)
// repeating an effectconst 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 effectconst 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 controlconst fromPull: Stream.Stream<number> = Stream.fromPull( Effect.succeed(() => Effect.succeed([[1, 2, 3]] as const)) // pull fn)
// unwrap a lazy stream effectconst 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") ))From paginated APIs
Section titled “From paginated APIs”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 itconst 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”import { Effect, Queue, Stream } from "effect"
declare const element: HTMLElementclass LetterError extends Error { constructor(readonly cause: unknown) { super("letter") } }
// DOM / EventTargetconst 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 insideconst 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”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 sourcesconst ticks: Stream.Stream<void> = Stream.tick("1 seconds") // Stream<void> at intervalconst 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 |
Transforming — the operator catalog
Section titled “Transforming — the operator catalog”Elementwise maps
Section titled “Elementwise maps”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 appropriateconst 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 | — |
Taking, dropping, windowing
Section titled “Taking, dropping, windowing”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 effectconst 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.
Concurrency within the stream
Section titled “Concurrency within the stream”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 streamsconst 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.Consuming — the run* family and Sinks
Section titled “Consuming — the run* family and Sinks”A Stream does nothing until run. Different collectors turn its emission into an Effect.
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 foldsconst 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)
// conversionsconst 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.
Ndjson — newline-delimited JSON streams
Section titled “Ndjson — newline-delimited JSON streams”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 iterabledeclare 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 ndjsonconst 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 causeconst 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 framingconst 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.
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 hopconst 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:
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 streamsconst 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.
PubSub / Queue → Stream interop
Section titled “PubSub / Queue → Stream interop”Where queues distribute work (one item to one worker) and pubsubs broadcast information (one item to all subscribers), streams express processing over either:
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.fromQueueconst 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 servicedeclare 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 sinksconst 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:
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 branchesconst 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 formsconst multi: Stream.Stream<number, FatalError> = source.pipe(Stream.catchTags({ TransientError: (e) => Stream.empty // absorb }))
// when diagnoses share a TaggedError union, use reason-typed catchingconst reason: Stream.Stream<number, never> = source.pipe(Stream.catchReason("TransientError", (r) => Stream.succeed(-1)))
// generic catch-all / cause-levelconst 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:
flowchart TD
Q1{"Items already in memory<br/>and small?"}
Q1 -- yes --> ARR["Array / Chunk<br/>map / filter directly"]
Q1 -- no --> Q2{"Producer/consumer<br/>rate mismatch?"}
Q2 -- no --> Q3{"Time involved?<br/>polling / ticks / delays / windowing"}
Q3 -- yes --> STR1["Stream<br/>pull + schedule"]
Q3 -- no --> Q4{"Effectful per element?"}
Q4 -- yes --> EFF["Effect.forEach(items, f)<br/>with concurrency"]
Q4 -- no --> ARR2["Array map"]
Q2 -- yes --> Q5{"Each item processed<br/>once, by one consumer?"}
Q5 -- yes --> QU["Queue + workers<br/>or Stream.fromQueue"]
Q5 -- no --> Q6{"All consumers see every item?"}
Q6 -- yes --> PS["PubSub + Stream.fromPubSub"]
Q6 -- no --> MIX["Queue for work<br/>PubSub for events<br/>(both services)"]
EFF -. streamable .-> STR2["Wrap as Stream<br/>when pipeline grows"]
| 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.
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 exhaustionconst 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 onesconst 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 materializeconst 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 sinksconst drained: Effect.Effect<void, ApiError | EnrichError> = enrichedStream.pipe(Stream.runDrain)Why each choice
Section titled “Why each choice”paginateoverfromEffectSchedule. Pagination cursors are data-driven — the next cursor is inside the response, not a clock.paginatethreadsOption<S>cursor state;fromEffectScheduledrives on ticks alone.flatMap(...fromIterable)afterpaginate.paginatesteps in pages (Array<Listed>is one element per step). Consumers usually want elements; flattening at the boundary makes latermapEffectper-item reasoning natural. If page-level batching matters downstream, defer flattening.mapEffectwithconcurrency: 4. Lookup is I/O-bound; 4 is a plausible per-route limit without collapsing to serial. Adjust with the sameSemaphorediscipline from ch.13 if the enrichment shares a pool with other code.retryinsidemapEffect, 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)beforerunCollect. Materializing an unbounded remote source is how you OOM. Bounding at consumption site is the discipline — treat unconstrainedrunCollectlike unboundedPromise.alland justify it.TestClock.adjustdrives the whole sleep fabric. BothlistPagebackoffs and the enrichmentexponential("100 millis")retries areClock.sleepunderneath; moving virtual time once flushes both without real delays. Wherepaginateis cursor-driven,TestClockmatters only for the retry delays, not for cursor advancement.