Hyperlinkv0.8.0-beta.28

Stream

Stream.scheduleconsteffect/Stream.ts:2736
<X, E2, R2, A>(schedule: Schedule.Schedule<X, NoInfer<A>, E2, R2>): <
  E,
  R
>(
  self: Stream<A, E, R>
) => Stream<A, E | E2, R2 | R>
<A, E, R, X, E2, R2>(
  self: Stream<A, E, R>,
  schedule: Schedule.Schedule<X, NoInfer<A>, E2, R2>
): Stream<A, E | E2, R | R2>

Schedules the stream's elements according to the provided schedule.

Example (Scheduling stream elements)

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

const program = Effect.gen(function*() {
  const result = yield* Stream.make(1, 2, 3).pipe(
    Stream.schedule(Schedule.spaced("10 millis")),
    Stream.runCollect
  )

  yield* Console.log(result)
})

Effect.runPromise(program)
// Output: [ 1, 2, 3 ]
rate limiting
Source effect/Stream.ts:273618 lines
export const schedule: {
  <X, E2, R2, A>(
    schedule: Schedule.Schedule<X, NoInfer<A>, E2, R2>
  ): <E, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>
  <A, E, R, X, E2, R2>(
    self: Stream<A, E, R>,
    schedule: Schedule.Schedule<X, NoInfer<A>, E2, R2>
  ): Stream<A, E | E2, R | R2>
} = dual(2, <A, E, R, X, E2, R2>(
  self: Stream<A, E, R>,
  schedule: Schedule.Schedule<X, NoInfer<A>, E2, R2>
): Stream<A, E | E2, R | R2> =>
  self.channel.pipe(
    Channel.flattenArray,
    Channel.schedule(schedule),
    Channel.map(Arr.of),
    fromChannel
  ))