variable / @typed/fx/Fx

groupedWithin

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.

Package version
2.0.0-beta.13
Category
Time and rate
Since
1.0.0

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.

Fx timelinegroupedWithin flushes when its timer wins and again at source completiongroupedWithin
OperatorgroupedWithin(3, 2 turns)
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.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

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.

Source