Hyperlinkv0.8.0-beta.28

Stream

Stream.pipeThroughChannelconsteffect/Stream.ts:9049
<R2, E, E2, A, A2>(
  channel: Channel.Channel<
    Arr.NonEmptyReadonlyArray<A2>,
    E2,
    unknown,
    Arr.NonEmptyReadonlyArray<A>,
    E,
    unknown,
    R2
  >
): <R>(self: Stream<A, E, R>) => Stream<A2, E2, R2 | R>
<R, R2, E, E2, A, A2>(
  self: Stream<A, E, R>,
  channel: Channel.Channel<
    Arr.NonEmptyReadonlyArray<A2>,
    E2,
    unknown,
    Arr.NonEmptyReadonlyArray<A>,
    E,
    unknown,
    R2
  >
): Stream<A2, E2, R | R2>

Pipes this stream through a channel that consumes and emits chunked elements.

Details

The channel receives NonEmptyReadonlyArray chunks and can transform both the output elements and error type.

Example (Piping through a channel)

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

type NumberChunk = readonly [number, ...Array<number>]

const doubleChunks = Channel.identity<NumberChunk, never, unknown>().pipe(
  Channel.map((chunk) => Array.map(chunk, (n) => n * 2))
)

const program = Effect.gen(function*() {
  const result = yield* Stream.fromArray([1, 2, 3]).pipe(
    Stream.rechunk(2),
    Stream.pipeThroughChannel(doubleChunks),
    Stream.runCollect
  )
  yield* Console.log(result)
})

Effect.runPromise(program)
// => [2, 4, 6]
Pipe
Source effect/Stream.ts:904912 lines
export const pipeThroughChannel: {
  <R2, E, E2, A, A2>(
    channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A2>, E2, unknown, Arr.NonEmptyReadonlyArray<A>, E, unknown, R2>
  ): <R>(self: Stream<A, E, R>) => Stream<A2, E2, R2 | R>
  <R, R2, E, E2, A, A2>(
    self: Stream<A, E, R>,
    channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A2>, E2, unknown, Arr.NonEmptyReadonlyArray<A>, E, unknown, R2>
  ): Stream<A2, E2, R | R2>
} = dual(2, <R, R2, E, E2, A, A2>(
  self: Stream<A, E, R>,
  channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A2>, E2, unknown, Arr.NonEmptyReadonlyArray<A>, E, unknown, R2>
): Stream<A2, E2, R | R2> => fromChannel(Channel.pipeTo(self.channel, channel)))