stream

package
v0.27.3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 6, 2026 License: Apache-2.0 Imports: 6 Imported by: 0

Documentation

Overview

Package stream contains domain-neutral, typed lazy stream operators.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type StreamOp

type StreamOp[In, Out any] func(iter.Seq2[In, error]) iter.Seq2[Out, error]

StreamOp transforms a lazy stream without erasing its item types.

func Batch

func Batch[T any](size int) StreamOp[T, []T]

Batch groups successful items into bounded batches. A final partial batch is emitted when the upstream ends normally.

func BatchByTime added in v0.23.0

func BatchByTime[T any](ctx context.Context, size int, timeout time.Duration) StreamOp[T, []T]

BatchByTime groups items up to a maximum size or until the timeout interval elapses.

func Collect

func Collect[T any](maxItems int) StreamOp[T, []T]

Collect materializes a bounded stream into one item. A non-positive limit is rejected at iteration time so the operator remains lazy.

func Debounce added in v0.23.0

func Debounce[T any](ctx context.Context, interval time.Duration) StreamOp[T, T]

Debounce drops items followed by another item within the interval.

The operator is a feeder goroutine plus a timing loop:

  • the feeder reads the upstream iterator and pushes items onto a channel, closing the channel when the iterator is exhausted or the context is canceled;
  • the loop restarts a timer on every received item and yields the item only when the timer fires without a newer item arriving.

Splitting the two makes each side's branch count legible: the feeder has one select and one error check, the loop has one select and the pending/timer pair. Neither alone exceeds the complexity budget.

func Filter

func Filter[T any](pred func(T) bool) StreamOp[T, T]

Filter forwards only items accepted by pred.

func FlatMap

func FlatMap[In, Out any](fn func(In) iter.Seq2[Out, error]) StreamOp[In, Out]

FlatMap expands each successful item sequentially. The inner stream must honor the downstream yield result; no goroutine is created per item.

func Map

func Map[In, Out any](fn func(In) Out) StreamOp[In, Out]

Map transforms every successful item. It does not intercept errors.

func MapE

func MapE[In, Out any](fn func(In) (Out, error)) StreamOp[In, Out]

MapE transforms every successful item and can stop with an error.

func Reduce

func Reduce[T, Acc any](initial Acc, fn func(Acc, T) (Acc, error)) StreamOp[T, Acc]

Reduce consumes the complete stream and emits one accumulated value.

func Take

func Take[T any](n int) StreamOp[T, T]

Take forwards at most n successful items. n <= 0 yields an empty stream.

func Throttle added in v0.23.0

func Throttle[T any](interval time.Duration) StreamOp[T, T]

Throttle drops items that arrive faster than the specified interval. Zero allocations.

func Window added in v0.23.0

func Window[T any](size, step int) StreamOp[T, []T]

Window groups items into overlapping sliding windows.

func WithContext

func WithContext[T any](ctx context.Context) StreamOp[T, T]

WithContext stops an operation when ctx is canceled before or during iteration. It preserves upstream errors and downstream early stop.

func WorkerPool added in v0.27.0

func WorkerPool[In, Out any](
	ctx context.Context,
	concurrency int,
	bufferCapacity int,
	worker func(context.Context, In) (Out, error),
) StreamOp[In, Out]

WorkerPool processes upstream items concurrently using a pool of workers. The worker callback runs for each item, up to `concurrency` in parallel. Results are yielded as they complete — not in input order — so callers that need ordering must sort downstream.

Termination:

  • Upstream ends → feeder closes workChan → workers drain it and exit → resultChan closes after the WaitGroup finishes.
  • Downstream stops consuming (yield returns false) → cancel pool → workers and feeder observe cancellation and return.
  • ctx is canceled → same as above.

Panic isolation: a panic inside worker or inside the upstream iterator is recovered and delivered as an error on the result stream. The pool keeps running other workers; a panicking item produces one error and the pool continues with the next item.

The function is safe to call with nil worker (yields a single error) and with concurrency < 1 (clamped to 1). bufferCapacity < 1 defaults to 2×concurrency.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL