Follow one execution
Partitions the stream into arrays, emitting when n is reached or duration elapses after the first element of the current group. The size must be a positive safe integer. A group can retain up to n values, so callers own the memory policy for valid sizes. Invalid sizes fail with Cause.IllegalArgumentError. Matches Effect Stream.groupedWithin.
groupedWithinTick 0. Input: a.
Read this diagram
Follow each lane from left to right. Events stacked vertically share a tick; the green cursor marks the current time across every lane.
- A value
The text inside the pill is the emitted value. - Work starts
The raised chevron starts an inner run (^in the source). - The run returns
The vertical bar ends this lane’s run. - A cause is delivered
The exclamation mark belongs to this lane. - Work is interrupted
The cross marks cancellation of this run. - Current time
The line and diamond move together across all lanes. - Happening now
A highlighted event is at the current tick. - Still ahead
Muted, dashed values have not happened yet. - Time continues
The lane’s arrow is not a return marker. An empty stretch can be quiet work that is still running.
Illustrated ticks start at 0. At 1×, one illustrated tick takes one second; captions specify real durations when timing matters. A cause or interruption belongs to its lane, and other work may continue. Scroll horizontally to inspect the rest of a long timeline.
Each run retains at most n values and at most one timer fiber in the required Scope. A source failure is Sink delivery, so flushNow still runs afterward unless the consumer interrupts the run. Flushing or interruption cancels the timer; invalid sizes fail before subscription.
Import
import { groupedWithin } from "@typed/fx/Fx";Signatures
export declare const groupedWithin: {
(n: number, duration: Duration.Input): <A, E, R>(self: Fx<A, E, R>) => Fx<NonEmptyReadonlyArray<A>, E | Cause.IllegalArgumentError, R | Scope.Scope>;
<A, E, R>(self: Fx<A, E, R>, n: number, duration: Duration.Input): Fx<NonEmptyReadonlyArray<A>, E | Cause.IllegalArgumentError, R | Scope.Scope>;
};Why
groupedWithin bounds a batch by both cardinality and time. The timer starts with the first value,
and whichever boundary wins emits the ordered non-empty group and resets both buffer and timer.
A partial group is also flushed when the source run returns, including after delivered failure.
Ownership and lifetime
Each run retains at most n values and at most one timer fiber in the required Scope. A source
failure is Sink delivery, so flushNow still runs afterward unless the consumer interrupts the
run. Flushing or interruption cancels the timer; invalid sizes fail before subscription.
Examples
import { Cause, Effect } from "effect"
import { Fx, Sink } from "@typed/fx"
const source = Fx.make<number, string>((sink) =>
sink.onSuccess(1).pipe(Effect.andThen(sink.onFailure(Cause.fail("boom"))))
)
const program = Fx.groupedWithin(source, 2, "1 hour").run(
Sink.make(Effect.logError, (group) => Effect.log(group))
) // logs the failure, then flushes [1]Other public imports
These import paths expose the same declaration. Each page retains its own public name and signature.