Hyperlinkv0.8.0-beta.28

Stream

Stream.partitionEffectconsteffect/Stream.ts:4578
<A, Pass, Fail, EX, RX>(
  filter: Filter.FilterEffect<NoInfer<A>, Pass, Fail, EX, RX>,
  options?: {
    readonly capacity?: number | "unbounded" | undefined
    readonly concurrency?: number | "unbounded" | undefined
  }
): <E, R>(
  self: Stream<A, E, R>
) => Effect.Effect<
  [passes: Stream<Pass, E | EX>, fails: Stream<Fail, E | EX>],
  never,
  R | RX | Scope.Scope
>
<A, E, R, Pass, Fail, EX, RX>(
  self: Stream<A, E, R>,
  filter: Filter.FilterEffect<NoInfer<A>, Pass, Fail, EX, RX>,
  options?: {
    readonly capacity?: number | "unbounded" | undefined
    readonly concurrency?: number | "unbounded" | undefined
  }
): Effect.Effect<
  [passes: Stream<Pass, E | EX>, fails: Stream<Fail, E | EX>],
  never,
  R | RX | Scope.Scope
>

Splits a stream with an effectful Filter, returning scoped streams for filter successes and failures.

When to use

Use when you need to classify each stream element with an effectful Filter and consume both passing and failing mapped values as streams.

Details

The returned streams are backed by queues in the current scope and should be consumed while that scope remains open. The first stream emits success values from the filter, and the second emits failure values.

Source effect/Stream.ts:457853 lines
export const partitionEffect: {
  <A, Pass, Fail, EX, RX>(filter: Filter.FilterEffect<NoInfer<A>, Pass, Fail, EX, RX>, options?: {
    readonly capacity?: number | "unbounded" | undefined
    readonly concurrency?: number | "unbounded" | undefined
  }): <E, R>(self: Stream<A, E, R>) => Effect.Effect<
    [
      passes: Stream<Pass, E | EX>,
      fails: Stream<Fail, E | EX>
    ],
    never,
    R | RX | Scope.Scope
  >
  <A, E, R, Pass, Fail, EX, RX>(
    self: Stream<A, E, R>,
    filter: Filter.FilterEffect<NoInfer<A>, Pass, Fail, EX, RX>,
    options?: {
      readonly capacity?: number | "unbounded" | undefined
      readonly concurrency?: number | "unbounded" | undefined
    }
  ): Effect.Effect<
    [
      passes: Stream<Pass, E | EX>,
      fails: Stream<Fail, E | EX>
    ],
    never,
    R | RX | Scope.Scope
  >
} = dual(
  (args) => isStream(args[0]),
  <A, E, R, Pass, Fail, EX, RX>(
    self: Stream<A, E, R>,
    filter: Filter.FilterEffect<NoInfer<A>, Pass, Fail, EX, RX>,
    options?: {
      readonly capacity?: number | "unbounded" | undefined
      readonly concurrency?: number | "unbounded" | undefined
    }
  ): Effect.Effect<
    [
      passes: Stream<Pass, E | EX>,
      fails: Stream<Fail, E | EX>
    ],
    never,
    R | RX | Scope.Scope
  > =>
    Effect.map(
      partitionQueue<Result.Result<Pass, Fail>, E | EX, R | RX, Pass, Fail>(
        mapEffect(self, (a) => filter(a as NoInfer<A>), options),
        (result) => result,
        options
      ),
      ([passes, fails]) => [fromQueue(passes), fromQueue(fails)] as const
    )
)