variable / @typed/fx/Fx

grouped

Partitions the stream into non-empty arrays of size n. The final array may be smaller if there are leftover elements. 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.grouped.

Package version
2.0.0-beta.13
Category
Stateful transforms
Since
1.0.0

Follow one execution

Partitions the stream into non-empty arrays of size n. The final array may be smaller if there are leftover elements. 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.grouped.

Fx timelinegrouped emits full batches and flushes the final partial batchgrouped
Operatorgrouped(2)
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 releases its buffer after the final flush or interruption. Invalid sizes deliver failure before source acquisition. A terminal observer may interrupt before the post-failure flush, but a Sink that handles the Cause can receive the partial group afterward.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

Import

import { grouped } from "@typed/fx/Fx";

Signatures

export declare const grouped: {
    (n: number): <A, E, R>(self: Fx<A, E, R>) => Fx<NonEmptyReadonlyArray<A>, E | Cause.IllegalArgumentError, R>;
    <A, E, R>(self: Fx<A, E, R>, n: number): Fx<NonEmptyReadonlyArray<A>, E | Cause.IllegalArgumentError, R>;
};

Why

grouped exposes fixed-size batching without changing source order. Full groups contain exactly n values. Because failure is delivered to the Sink while Fx.run remains infallible, any partial group is flushed after the source run returns even if a failure was delivered first.

Ownership and lifetime

Each run retains at most n values and releases its buffer after the final flush or interruption. Invalid sizes deliver failure before source acquisition. A terminal observer may interrupt before the post-failure flush, but a Sink that handles the Cause can receive the partial group afterward.

Examples

import { Cause, Effect, Ref } from "effect"
import { Fx, Sink } from "@typed/fx"
const program = Effect.gen(function* () {
  const deliveries = yield* Ref.make<Array<string>>([])
  const source = Fx.make<number, string>((sink) =>
    sink.onSuccess(1).pipe(Effect.andThen(sink.onFailure(Cause.fail("boom"))))
  )
  yield* Fx.grouped(source, 2).run(Sink.make(
    () => Ref.update(deliveries, (xs) => [...xs, "failure"]),
    (group) => Ref.update(deliveries, (xs) => [...xs, `group:${group.join(",")}`])
  ))
  return yield* Ref.get(deliveries) // ["failure", "group:1"]
})

Other public imports

These import paths expose the same declaration. Each page retains its own public name and signature.

Source