Hyperlinkv0.8.0-beta.28

Stream

Stream.groupedWithinconsteffect/Stream.ts:8214
(chunkSize: number, duration: Duration.Input): <A, E, R>(
  self: Stream<A, E, R>
) => Stream<Array<A>, E, R>
<A, E, R>(
  self: Stream<A, E, R>,
  chunkSize: number,
  duration: Duration.Input
): Stream<Array<A>, E, R>

Partitions the stream into arrays, emitting when the chunk size is reached or the duration passes.

Example (Grouping elements by size or time)

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

const program = Effect.gen(function*() {
  const values = yield* Stream.make(1, 2, 3).pipe(
    Stream.groupedWithin(2, "5 seconds"),
    Stream.runCollect
  )
  yield* Console.log(values)
})

Effect.runPromise(program)
// Output: [ [ 1, 2 ], [ 3 ] ]
grouping
Source effect/Stream.ts:821416 lines
export const groupedWithin: {
  (
    chunkSize: number,
    duration: Duration.Input
  ): <A, E, R>(self: Stream<A, E, R>) => Stream<Array<A>, E, R>
  <A, E, R>(self: Stream<A, E, R>, chunkSize: number, duration: Duration.Input): Stream<Array<A>, E, R>
} = dual(3, <A, E, R>(
  self: Stream<A, E, R>,
  chunkSize: number,
  duration: Duration.Input
): Stream<Array<A>, E, R> =>
  aggregateWithin(
    self,
    Sink.take(chunkSize),
    Schedule.spaced(duration)
  ))