Hyperlinkv0.8.0-beta.28

Stream

Stream.flattenEffectconsteffect/Stream.ts:2078
<
  Arg extends
    | Stream<Effect.Effect<any, any, any>, any, any>
    | {
        readonly concurrency?: number | "unbounded" | undefined
        readonly unordered?: boolean | undefined
      }
    | undefined = {
    readonly concurrency?: number | "unbounded" | undefined
    readonly unordered?: boolean | undefined
  }
>(
  selfOrOptions?: Arg,
  options?:
    | {
        readonly concurrency?: number | "unbounded" | undefined
        readonly unordered?: boolean | undefined
      }
    | undefined
): [Arg] extends [
  Stream<
    Effect.Effect<infer _A, infer _EX, infer _RX>,
    infer _E,
    infer _R
  >
]
  ? Stream<_A, _EX | _E, _RX | _R>
  : <A, EX, RX, E, R>(
      self: Stream<Effect.Effect<A, EX, RX>, E, R>
    ) => Stream<A, EX | E, RX | R>

Flattens a stream of Effect values into a stream of their results.

When to use

Use when stream elements already are effects and their successes should become stream elements while their failures enter the stream error channel.

Example (Flattening a stream of Effect values into a stream of their results)

import { Console, Effect, Stream } from "effect"

const stream = Stream.make(Effect.succeed(1), Effect.succeed(2), Effect.succeed(3))

const program = Effect.gen(function*() {
  const result = yield* Stream.runCollect(stream.pipe(Stream.flattenEffect()))
  yield* Console.log(result)
})

Effect.runPromise(program)
// Output: [1, 2, 3]
mapping
Source effect/Stream.ts:207826 lines
export const flattenEffect: <
  Arg extends Stream<Effect.Effect<any, any, any>, any, any> | {
    readonly concurrency?: number | "unbounded" | undefined
    readonly unordered?: boolean | undefined
  } | undefined = {
    readonly concurrency?: number | "unbounded" | undefined
    readonly unordered?: boolean | undefined
  }
>(
  selfOrOptions?: Arg,
  options?: {
    readonly concurrency?: number | "unbounded" | undefined
    readonly unordered?: boolean | undefined
  } | undefined
) => [Arg] extends [Stream<Effect.Effect<infer _A, infer _EX, infer _RX>, infer _E, infer _R>] ?
  Stream<_A, _EX | _E, _RX | _R>
  : <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>) => Stream<A, EX | E, RX | R> = dual(
    (args) => isStream(args[0]),
    <A, E, R, EX, RX>(
      self: Stream<Effect.Effect<A, EX, RX>, E, R>,
      options?: {
        readonly concurrency?: number | "unbounded" | undefined
        readonly unordered?: boolean | undefined
      } | undefined
    ): Stream<A, EX | E, RX | R> => mapEffect(self, identity, options)
  )