variable / @typed/fx/Fx

flatMapConcurrently

Maps each element of an Fx to a new Fx, running them concurrently with a limit.

Package version
2.0.0-beta.7
Category
Concurrent work
Since
1.0.0

Follow one execution

Maps each element of an Fx to a new Fx, running them concurrently with a limit. This scenario uses concurrent source deliveries: a pending source callback can wait for admission while another callback is already active.

Fx timelineflatMapConcurrently waits when every permit is occupiedflatMapConcurrently
OperatorflatMapConcurrently(load, 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.

A non-positive, fractional, infinite, or unsafe-integer limit fails through the Sink with Cause.IllegalArgumentError. Source and inner failures and services remain typed. The required Scope owns waiting and active fibers; source completion waits for all of them, and interruption cancels both groups and runs active inner finalizers.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

Import

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

Signatures

export declare const flatMapConcurrently: {
    <A, B, E2, R2>(f: (a: A) => Fx<B, E2, R2>, concurrency: number): <E, R>(self: Fx<A, E, R>) => Fx<B, E | E2 | Cause.IllegalArgumentError, R | R2 | Scope.Scope>;
    <A, E, R, B, E2, R2>(self: Fx<A, E, R>, f: (a: A) => Fx<B, E2, R2>, concurrency: number): Fx<B, E | E2 | Cause.IllegalArgumentError, R | R2 | Scope.Scope>;
};

Why

This is the admission-controlled form of {@link flatMap }. It bounds active work while retaining every source value, which is appropriate when a remote service or local resource has a known concurrency budget.

Concurrency, ordering, and buffering

Each source value creates one child fiber. A semaphore admits at most concurrency inners into execution; excess fibers wait for a permit rather than being dropped. Values from each inner retain local order, but admitted inners may interleave and output is not reordered into source order.

Ownership and lifetime

A non-positive, fractional, infinite, or unsafe-integer limit fails through the Sink with Cause.IllegalArgumentError. Source and inner failures and services remain typed. The required Scope owns waiting and active fibers; source completion waits for all of them, and interruption cancels both groups and runs active inner finalizers.

Examples

import { Fx } from "@typed/fx"
import { Effect } from "effect"

const jobs = Fx.fromIterable([
  { id: "one", wait: "20 millis" as const },
  { id: "two", wait: "20 millis" as const },
  { id: "three", wait: "1 millis" as const }
])
const bounded = Fx.flatMapConcurrently(
  jobs,
  ({ id, wait }) => Fx.at(id, wait),
  2
)

Effect.runPromise(Effect.scoped(Fx.collectAll(bounded))).then(console.log)
// "three" waits for one of the first two permits despite its shorter delay

Other public imports

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

Source