flx

package module
v0.1.7 Latest Latest
Warning

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

Go to latest
Published: Apr 8, 2026 License: MIT Imports: 14 Imported by: 0

README

flx

flx ��һ���Է��� Stream[T] Ϊ���ĵ���ʽ�����붯̬�������ƿ⡣

������ go-zero �� fx ���������ؽ����������� fx API ������Ŀ����ֱ�ӣ�

  • �� Stream[T] �ṩ���Ͱ�ȫ����ʽ����
  • �����̶����������޲�������̬������ǿ����������
  • ����ʽ context.Context ������ʽ context option
  • �� API �������ִ� Go

����

  • ��һ�����������ˡ�ӳ�䡢���顢�ۺ�
  • ��Ҫ�� pipeline �п��Ʋ�����
  • ��Ҫ�������ж�̬���� worker ����
  • ��Ҫ������ʱȡ������ worker
  • ��Ҫ����ʽ����֮�⸴�� retry / timeout / parallel ����

�汾Ҫ��

  • Go 1.26.1+

��ǰģ��·���� github.com/ezra-sullivan/flx��ʾ������ֱ��ʹ�ã�

import "github.com/ezra-sullivan/flx"

worker ������������ API �����Ƽ��������Ĺ����Ӱ����룺

import "github.com/ezra-sullivan/flx/pipeline/control"

���������������µ�ģ��·����ʾ���еĵ���·����Ҫͬ���滻��

����ʾ��

package main

import (
	"fmt"

	"github.com/ezra-sullivan/flx"
)

func main() {
	out := flx.Map(
		flx.Values(1, 2, 3, 4, 5).
			Filter(func(v int) bool { return v%2 == 0 }),
		func(v int) int { return v * 10 },
	)

	out.ForEach(func(v int) {
		fmt.Println(v)
	})
}

������

20
40

ʵսʾ��

�ֿ��ﻹ����һ����ֱ�����е� HTTP ͼƬ���� example��

go run ./examples/http_image_pipeline

����ʾ���᣺

  • �� https://picsum.photos/v2/list ��ҳ��ȡͼƬ
  • ���ش�ͼ�����浽 examples/http_image_pipeline/.output/original/��Ĭ�� 5 �ţ�
  • �Ѵ�ͼ����Сͼ�����浽 examples/http_image_pipeline/.output/processed/
  • ���ؽ׶λ��� flx.DoWithRetryCtx(...) �Զ����ԣ�Ĭ�� 3 ��
  • Ĭ��ֻ�����غʹ����������ϴ��������ϴ����� examples/http_image_pipeline/main.go ���� UploadEndpoint

����ʵ�ֶ��ڲֿ��

  • Stage ���ϰ棺examples/http_image_pipeline/stage_pipeline.go
  • ԭ�� flx API �棺examples/http_image_pipeline/native_pipeline.go

����˵���� doc/examples/http-image-pipeline.md��

API ��״

flx ���á�ͬ���Ͳ������������������Ͳ���ʹ�ð������ͺ������Ļ������ơ�

���磺

  • ͬ���Ͳ�����Filter��Sort��Head��Skip
  • �����Ͳ�����flx.Map��flx.FlatMap��flx.MapContext
  • stage ������װ��flx.Stage��flx.StageErr��flx.FlatStage��flx.FlatStageErr��flx.Tap

����һ�� pipeline ������ stage ֮�䱣��ͬһ�� Stream[T] ���ͣ����Լ����÷�����������д�ø������������� stage�������� flx.Stage(...).Through(...).Through(...)��

s := flx.Values(1, 2, 3).Filter(func(v int) bool { return v > 1 })
out := flx.Map(s, func(v int) string { return fmt.Sprintf("v=%d", v) })

�������Ʋ���Ϊ��׷�������ر𣬶�����Ϊ Go 1.26 ��Ȼ��֧�֡��������Ͳ����ķ�������

��������

Stream �����������
  • Values / From / FromChan / Concat
  • Filter / Buffer / Sort / Reverse
  • Head / Tail / Skip
  • Map / MapErr
  • FlatMap / FlatMapErr
  • MapContext / FlatMapContext
  • Stage / StageErr / FlatStage / FlatStageErr / Tap
  • DistinctBy / GroupBy / Chunk
  • DistinctByCount / GroupByCount
  • DistinctByWindow / GroupByWindow
  • Reduce
�ս�����ѯ
  • Done / DoneErr
  • ForEach / ForEachErr
  • Parallel / ParallelErr
  • Count / Collect
  • First / Last
  • AllMatch / AnyMatch / NoneMatch
  • Max / Min
  • Err
��������������
  • control.WithWorkers
  • control.WithUnlimitedWorkers
  • control.WithDynamicWorkers
  • control.WithForcedDynamicWorkers
  • control.WithInterruptibleWorkers�����ݱ�����
  • control.WithErrorStrategy
  • control.ErrorStrategyFailFast
  • control.ErrorStrategyCollect
  • control.ErrorStrategyLogAndContinue
�������ߺ���
  • Parallel / ParallelErr / ParallelWithErrorStrategy
  • DoWithRetry / DoWithRetryCtx
  • DoWithTimeout / DoWithTimeoutCtx

�ڴ��߽�

  • DistinctBy ��ά����ǰ����ȫ���Ѽ� key ���ϣ�����Ψһ key �����������ڴ�Ҳ������������
  • GroupBy ���Ȼ����������������ٰ� key �������飻�����ʺϳ������ݼ����޽�����
  • ������ API ���ʺ�����ȷ�߽������������ݡ�
  • DistinctByCount / GroupByCount ���ð����������зֵ� tumbling window���ʺ���Ҫ��Ӳ�ڴ��߽���ȥ��/���鳡����
  • DistinctByWindow / GroupByWindow ���� processing-time tumbling window������ʽ���� context.Context�����Ǹ��ʺ�����ʱ�鵵��������������

��̬����

��������

����ʱ������������ worker��ֻӰ���������� slot �� worker��

��Ҫ���⵼�룺github.com/ezra-sullivan/flx/pipeline/control

ctrl := control.NewConcurrencyController(4)

out := flx.Map(
	flx.Values(1, 2, 3, 4, 5),
	func(v int) int { return v * 10 },
	control.WithDynamicWorkers(ctrl),
)

ctrl.SetWorkers(8)
ctrl.SetWorkers(2)

out.Done()
ǿ������

����ʱȡ������ worker��Ҫ��ʹ�� MapContext* �� FlatMapContext*��

��Ҫ���⵼�룺github.com/ezra-sullivan/flx/pipeline/control

ctx := context.Background()
ctrl := control.NewConcurrencyController(4)

out := flx.FlatMapContext(
	ctx,
	flx.Values("a", "b", "c"),
	func(ctx context.Context, v string, pipe chan<- string) {
		select {
		case <-ctx.Done():
			return
		default:
		}

		flx.SendContext(ctx, pipe, v+"!")
	},
	control.WithForcedDynamicWorkers(ctrl),
)

ctrl.SetWorkers(1)
out.Done()

control.WithInterruptibleWorkers ��Ȼ���ã����´��뽨��ͳһд�� control.WithForcedDynamicWorkers��

MapContext / MapContextErr �ĵ�������������Ҳ����Ӧ ctx.Done()���������� FlatMapContext* ���Լ������η���ֵ����Ȼ������ʽʹ�� SendContext��

������������

Ĭ�ϴ��������� control.ErrorStrategyFailFast������ worker ���ش����� panic��

  • fail-fast������ȡ����ǰ�����������ս��׶α�¶����
  • collect������ִ�У������ϲ�����
  • log-and-continue����¼��־������

ҵ���������������Ҫ�ȶ������մ����߽磬����ʹ�û��������� source �� *Err �ս����������� DoneErr / CollectErr��

out := flx.MapErr(flx.Values("1", "x", "3"), strconv.Atoi)
items, err := out.CollectErr()

����˵����

  • From �������������� panic�������� stream ����״̬
  • DoneErr / CollectErr ������������ *Err �ս�����������ʽ�õ���������
  • FirstErr / AllMatchErr / AnyMatchErr / NoneMatchErr ������������·�����н������������أ����ں�̨ drain ����
  • ���ĸ���· *Err API ���ص��ǵ�ǰ�������գ����غ��ŷ����� fail-fast error ����֤�����ڷ���ֵ��
  • First / AllMatch / AnyMatch / NoneMatch ������·�ս�����Ҳ��ѭ fail-fast ���壻���������Ѿ���¼ fail-fast ���������ǻ��������� *Err �ս�����һ�� panic

flx Ĭ�ϰѴ�����ģΪ stream ״̬�������ǹٷ��ṩһ�� value + error �� item ������ �������� fx Ǩ�ƹ�����֮ǰ�������ﴫ struct{ Value T; Err error } ������������������д���� flx ����Ȼ���Ա���������Ӧ�ñ���Ϊҵ�����ݽ�ģ���ʺϡ����ֳɹ�������ͳһ�ռ�ʧ����ij��������������� MapErr / CollectErr / DoneErr ����������ͨ����

�� go-zero fx ����Ҫ����

  • fx.Just -> flx.Values
  • fx.Range -> flx.FromChan
  • stream.Map(...) -> flx.Map(stream, ...)
  • stream.Walk(...) -> flx.FlatMap(stream, ...)
  • stream.WalkCtx... -> flx.MapContext... / flx.FlatMapContext...
  • stream.Merge() -> stream.Collect() / stream.CollectErr()
  • fx.WithDynamicWorkersCtx -> control.WithForcedDynamicWorkers
  • fx.SendCtx -> flx.SendContext
  • flx ���ṩ�ٷ� Result[T] / ItemError[T]����Ҫ��������ʱ���Զ���ҵ���ṹ��

�������� fx Ǩ�ƣ��������ϱ� README��doc/quickstart.md �� doc/guide.md �е� API ����������˵���𲽵�����

�ĵ�����

Documentation

Overview

Package flx provides generic stream processing with dynamic concurrency control, explicit context-aware transforms, and reusable retry, timeout, and parallel execution helpers.

Module path: github.com/ezra-sullivan/flx

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrCanceled is an alias for context.Canceled.
	ErrCanceled = context.Canceled
	// ErrTimeout is an alias for context.DeadlineExceeded.
	ErrTimeout = context.DeadlineExceeded
	// ErrNilContext reports that a required parent context was nil.
	ErrNilContext = errors.New("flx: nil context")
	// ErrNegativeTimeout reports that a timeout duration was negative.
	ErrNegativeTimeout = errors.New("flx: timeout must not be negative")
	// ErrInvalidRetryTimes reports that a retry count was zero or negative.
	ErrInvalidRetryTimes = errors.New("flx: retry times must be greater than 0")
	// ErrNegativeRetryInterval reports that a retry interval was negative.
	ErrNegativeRetryInterval = errors.New("flx: retry interval must not be negative")
	// ErrNegativeRetryTimeout reports that a total retry timeout was negative.
	ErrNegativeRetryTimeout = errors.New("flx: retry timeout must not be negative")
	// ErrNegativeAttemptTimeout reports that a per-attempt timeout was negative.
	ErrNegativeAttemptTimeout = errors.New("flx: attempt timeout must not be negative")
	// ErrAttemptTimeoutRequiresRetryCtx reports that attempt timeouts only work
	// with the context-aware retry API.
	ErrAttemptTimeoutRequiresRetryCtx = errors.New("flx: WithAttemptTimeout requires DoWithRetryCtx")
	// ErrRetryAttemptTimeout reports that one retry attempt exceeded its own
	// attempt timeout.
	ErrRetryAttemptTimeout = errors.New("flx: retry attempt timeout")
)
View Source
var (
	// ErrInvalidWindowCount reports that a count-based window size was less than
	// one.
	ErrInvalidWindowCount = errors.New("flx: window count must be greater than 0")
	// ErrInvalidWindowDuration reports that a time-based window duration was not
	// positive.
	ErrInvalidWindowDuration = errors.New("flx: window duration must be positive")
)

Functions

func DoWithRetry

func DoWithRetry(fn func() error, opts ...RetryOption) error

DoWithRetry executes fn until it succeeds or the retry budget is exhausted.

func DoWithRetryCtx

func DoWithRetryCtx(ctx context.Context, fn func(context.Context, int) error, opts ...RetryOption) error

DoWithRetryCtx executes fn until it succeeds or the retry budget is exhausted, passing the current attempt context and zero-based attempt index.

func DoWithTimeout

func DoWithTimeout(fn func() error, timeout time.Duration, opts ...TimeoutOption) error

DoWithTimeout runs fn with a derived timeout context and returns its result.

func DoWithTimeoutCtx

func DoWithTimeoutCtx(fn func(context.Context) error, timeout time.Duration, opts ...TimeoutOption) error

DoWithTimeoutCtx runs fn with a derived timeout context and passes that context into the callback.

func Parallel

func Parallel(fns ...func())

Parallel runs each function in its own goroutine and panics if the chosen fail-fast strategy records an error.

func ParallelErr

func ParallelErr(fns ...func() error) error

ParallelErr runs each function in its own goroutine and returns the joined worker errors without applying stream fail-fast semantics.

func ParallelWithErrorStrategy

func ParallelWithErrorStrategy(strategy control.ErrorStrategy, fns ...func()) error

ParallelWithErrorStrategy runs each function in its own goroutine and applies strategy to worker panics and returned errors.

func Reduce

func Reduce[T, R any](s Stream[T], fn func(<-chan T) (R, error)) (R, error)

Reduce hands s's source channel to fn, drains any remaining items after fn returns, and joins fn's returned error with the stream error state.

func SendContext

func SendContext[T any](ctx context.Context, pipe chan<- T, item T) bool

SendContext sends item to pipe unless ctx has already been canceled.

Types

type Group added in v0.1.2

type Group[K comparable, T any] = streaming.Group[K, T]

Group holds one grouping key plus the items assigned to that key.

type RetryOption

type RetryOption func(*retryOptions)

RetryOption mutates the behavior of one retry call.

func WithAttemptTimeout

func WithAttemptTimeout(timeout time.Duration) RetryOption

WithAttemptTimeout sets a timeout for each individual retry attempt.

func WithIgnoreErrors

func WithIgnoreErrors(ignoreErrors []error) RetryOption

WithIgnoreErrors treats matching errors as successful completion.

func WithInterval

func WithInterval(interval time.Duration) RetryOption

WithInterval sets the delay between failed attempts.

func WithRetry

func WithRetry(times int) RetryOption

WithRetry sets the maximum number of attempts, including the first one.

func WithTimeout

func WithTimeout(timeout time.Duration) RetryOption

WithTimeout sets an overall timeout for the full retry loop.

type Stream

type Stream[T any] struct {
	// contains filtered or unexported fields
}

Stream is a lazy sequence of values backed by a channel plus shared error state that records upstream worker failures.

func Chunk

func Chunk[T any](s Stream[T], n int) Stream[[]T]

Chunk groups items into slices of size n, emitting a final short chunk when the source ends.

func Concat

func Concat[T any](s Stream[T], others ...Stream[T]) Stream[T]

Concat merges s with others and returns a stream that emits items from all inputs as they arrive.

func DistinctBy

func DistinctBy[T any, K comparable](s Stream[T], fn func(T) K) Stream[T]

DistinctBy keeps the first item for each key produced by fn.

func DistinctByCount added in v0.1.2

func DistinctByCount[T any, K comparable](s Stream[T], n int, fn func(T) K) Stream[T]

DistinctByCount keeps the first item for each key within windows of n input items, resetting the seen-key set after every window.

func DistinctByWindow added in v0.1.2

func DistinctByWindow[T any, K comparable](ctx context.Context, s Stream[T], every time.Duration, fn func(T) K) Stream[T]

DistinctByWindow keeps the first item for each key within a time window that starts when the first item in that window arrives.

func FlatMap

func FlatMap[T, U any](s Stream[T], fn func(T, chan<- U), opts ...control.Option) Stream[U]

FlatMap calls fn for each item and lets fn emit zero or more output values.

func FlatMapContext

func FlatMapContext[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T, chan<- U), opts ...control.Option) Stream[U]

FlatMapContext is FlatMap with a caller-provided context passed into each worker.

func FlatMapContextErr

func FlatMapContextErr[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T, chan<- U) error, opts ...control.Option) Stream[U]

FlatMapContextErr is FlatMapErr with a caller-provided context passed into each worker.

func FlatMapErr

func FlatMapErr[T, U any](s Stream[T], fn func(T, chan<- U) error, opts ...control.Option) Stream[U]

FlatMapErr calls fn for each item and records any returned worker error in the stream state.

func FlatStage added in v0.1.5

func FlatStage[I, O any](
	ctx context.Context,
	in Stream[I],
	fn func(context.Context, I, chan<- O),
	opts ...control.Option,
) Stream[O]

FlatStage applies fn to each item in in and lets fn emit zero or more output values. It is a thin semantic wrapper around FlatMapContext for stage-oriented pipelines.

func FlatStageErr added in v0.1.5

func FlatStageErr[I, O any](
	ctx context.Context,
	in Stream[I],
	fn func(context.Context, I, chan<- O) error,
	opts ...control.Option,
) Stream[O]

FlatStageErr applies fn to each item in in, lets fn emit zero or more output values, and records returned worker errors in the stream state. It is a thin semantic wrapper around FlatMapContextErr.

func From

func From[T any](generate func(chan<- T)) Stream[T]

From adapts a producer callback into a stream. Panics from generate are captured in the stream state and surfaced by terminal operations.

func FromChan

func FromChan[T any](source <-chan T) Stream[T]

FromChan wraps source as a Stream without changing its production semantics.

func GroupBy

func GroupBy[T any, K comparable](s Stream[T], fn func(T) K) Stream[[]T]

GroupBy drains s, groups items by fn, and emits groups in first-seen key order.

func GroupByCount added in v0.1.2

func GroupByCount[T any, K comparable](s Stream[T], n int, fn func(T) K) Stream[Group[K, T]]

GroupByCount groups items by fn within windows of n input items and emits one Group per key in first-seen order for each window.

func GroupByWindow added in v0.1.2

func GroupByWindow[T any, K comparable](ctx context.Context, s Stream[T], every time.Duration, fn func(T) K) Stream[Group[K, T]]

GroupByWindow groups items by fn within a time window that starts when the first item in that window arrives and flushes on timer tick or source close.

func Map

func Map[T, U any](s Stream[T], fn func(T) U, opts ...control.Option) Stream[U]

Map applies fn to each item in s and emits the mapped values.

func MapContext

func MapContext[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T) U, opts ...control.Option) Stream[U]

MapContext is Map with a caller-provided context passed into each worker.

func MapContextErr

func MapContextErr[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T) (U, error), opts ...control.Option) Stream[U]

MapContextErr is MapErr with a caller-provided context passed into each worker.

func MapErr

func MapErr[T, U any](s Stream[T], fn func(T) (U, error), opts ...control.Option) Stream[U]

MapErr applies fn to each item in s and records any returned error in the stream state.

func Stage added in v0.1.5

func Stage[I, O any](
	ctx context.Context,
	in Stream[I],
	fn func(context.Context, I) O,
	opts ...control.Option,
) Stream[O]

Stage applies fn to each item in in and emits the mapped values. It is a thin semantic wrapper around MapContext for pipelines that want explicit stage-shaped call sites without introducing a second runtime model.

func StageErr added in v0.1.5

func StageErr[I, O any](
	ctx context.Context,
	in Stream[I],
	fn func(context.Context, I) (O, error),
	opts ...control.Option,
) Stream[O]

StageErr applies fn to each item in in and records returned worker errors in the stream state. It is a thin semantic wrapper around MapContextErr for stage-oriented pipelines that still want stream-level error handling.

func Tap added in v0.1.5

func Tap[T any](
	ctx context.Context,
	in Stream[T],
	fn func(context.Context, T) error,
	opts ...control.Option,
) Stream[T]

Tap runs fn for each item in in and re-emits the original item when fn succeeds. If fn returns an error, Tap records that error in the stream state and does not forward the failed item.

func Values

func Values[T any](items ...T) Stream[T]

Values returns a stream that emits items in order and then closes.

func (Stream[T]) AllMatch

func (s Stream[T]) AllMatch(predicate func(T) bool) bool

AllMatch reports whether every item satisfies predicate. It drains the remainder of the stream after the first mismatch so delayed fail-fast errors can still surface.

func (Stream[T]) AllMatchErr

func (s Stream[T]) AllMatchErr(predicate func(T) bool) (bool, error)

AllMatchErr reports whether every item satisfies predicate and returns the current error state when it short-circuits.

func (Stream[T]) AnyMatch

func (s Stream[T]) AnyMatch(predicate func(T) bool) bool

AnyMatch reports whether any item satisfies predicate. It drains the remainder of the stream after the first match so delayed fail-fast errors can still surface.

func (Stream[T]) AnyMatchErr

func (s Stream[T]) AnyMatchErr(predicate func(T) bool) (bool, error)

AnyMatchErr reports whether any item satisfies predicate and returns the current error state when it short-circuits.

func (Stream[T]) Buffer

func (s Stream[T]) Buffer(n int) Stream[T]

Buffer inserts a channel buffer of size n between s and the returned stream.

func (Stream[T]) Collect

func (s Stream[T]) Collect() []T

Collect drains the stream into a slice and panics on fail-fast errors.

func (Stream[T]) CollectErr

func (s Stream[T]) CollectErr() ([]T, error)

CollectErr drains the stream into a slice and returns the final error state.

func (Stream[T]) Concat

func (s Stream[T]) Concat(others ...Stream[T]) Stream[T]

Concat merges s with others while preserving per-stream item order. Items from different input streams may interleave based on runtime scheduling.

func (Stream[T]) Count

func (s Stream[T]) Count() int

Count drains the stream and returns the number of items it produced.

func (Stream[T]) CountErr

func (s Stream[T]) CountErr() (int, error)

CountErr drains the stream, returns the item count, and returns the final error state.

func (Stream[T]) Done

func (s Stream[T]) Done()

Done drains the stream and panics if a fail-fast error was recorded.

func (Stream[T]) DoneErr

func (s Stream[T]) DoneErr() error

DoneErr drains the stream and returns the final error state.

func (Stream[T]) Err

func (s Stream[T]) Err() error

Err returns the currently accumulated stream error without draining the stream.

func (Stream[T]) Filter

func (s Stream[T]) Filter(fn func(T) bool, opts ...control.Option) Stream[T]

Filter keeps only the items for which fn returns true.

func (Stream[T]) First

func (s Stream[T]) First() (T, bool)

First returns the first item and then drains the rest of the stream so any delayed fail-fast error is observed before the call returns.

func (Stream[T]) FirstErr

func (s Stream[T]) FirstErr() (T, bool, error)

FirstErr returns the first item and the current error state, then drains the remaining source asynchronously.

func (Stream[T]) ForAll

func (s Stream[T]) ForAll(fn func(<-chan T))

ForAll hands the raw source channel to fn, then drains any leftovers and applies fail-fast panic behavior.

func (Stream[T]) ForAllErr

func (s Stream[T]) ForAllErr(fn func(<-chan T)) error

ForAllErr hands the raw source channel to fn, then drains any leftovers and returns the final error state.

func (Stream[T]) ForEach

func (s Stream[T]) ForEach(fn func(T))

ForEach calls fn for every item in the stream and panics if a fail-fast error was recorded.

func (Stream[T]) ForEachErr

func (s Stream[T]) ForEachErr(fn func(T)) error

ForEachErr calls fn for every item in the stream and returns the final error state.

func (Stream[T]) Head

func (s Stream[T]) Head(n int64) Stream[T]

Head returns a stream containing at most the first n items from s. The upstream source is drained after the head is satisfied so producers can exit.

func (Stream[T]) Last

func (s Stream[T]) Last() (T, bool)

Last drains the stream and returns the last item it observed.

func (Stream[T]) LastErr

func (s Stream[T]) LastErr() (T, bool, error)

LastErr drains the stream, returns the last item it observed, and returns the final error state.

func (Stream[T]) Max

func (s Stream[T]) Max(less func(T, T) bool) (T, bool)

Max drains the stream and returns the greatest item according to less.

func (Stream[T]) MaxErr

func (s Stream[T]) MaxErr(less func(T, T) bool) (T, bool, error)

MaxErr drains the stream, returns the greatest item according to less, and returns the final error state.

func (Stream[T]) Min

func (s Stream[T]) Min(less func(T, T) bool) (T, bool)

Min drains the stream and returns the least item according to less.

func (Stream[T]) MinErr

func (s Stream[T]) MinErr(less func(T, T) bool) (T, bool, error)

MinErr drains the stream, returns the least item according to less, and returns the final error state.

func (Stream[T]) NoneMatch

func (s Stream[T]) NoneMatch(predicate func(T) bool) bool

NoneMatch reports whether no item satisfies predicate. It drains the remainder of the stream after the first match so delayed fail-fast errors can still surface.

func (Stream[T]) NoneMatchErr

func (s Stream[T]) NoneMatchErr(predicate func(T) bool) (bool, error)

NoneMatchErr reports whether no item satisfies predicate and returns the current error state when it short-circuits.

func (Stream[T]) Parallel

func (s Stream[T]) Parallel(fn func(T), opts ...control.Option)

Parallel applies fn to each item using the same worker machinery as the transform operators and panics on fail-fast errors.

func (Stream[T]) ParallelErr

func (s Stream[T]) ParallelErr(fn func(T) error, opts ...control.Option) error

ParallelErr applies fn to each item using worker options and returns the final error state.

func (Stream[T]) Reverse

func (s Stream[T]) Reverse() Stream[T]

Reverse drains s, reverses the collected items, and replays them as a new stream.

func (Stream[T]) Skip

func (s Stream[T]) Skip(n int64) Stream[T]

Skip discards the first n items from s and emits the remainder.

func (Stream[T]) Sort

func (s Stream[T]) Sort(less func(T, T) bool) Stream[T]

Sort drains s, sorts all items with less, and then replays them as a new stream.

func (Stream[T]) Tail

func (s Stream[T]) Tail(n int64) Stream[T]

Tail returns a stream containing the last n items from s in original order.

func (Stream[T]) Tap added in v0.1.5

func (s Stream[T]) Tap(
	ctx context.Context,
	fn func(context.Context, T) error,
	opts ...control.Option,
) Stream[T]

Tap runs fn for each item in s and re-emits the original item when fn succeeds.

func (Stream[T]) Through added in v0.1.5

func (s Stream[T]) Through(
	ctx context.Context,
	fn func(context.Context, T) T,
	opts ...control.Option,
) Stream[T]

Through applies fn to each item in s and returns another Stream[T]. It exists to make same-type stage segments read fluently in a chain.

func (Stream[T]) ThroughErr added in v0.1.5

func (s Stream[T]) ThroughErr(
	ctx context.Context,
	fn func(context.Context, T) (T, error),
	opts ...control.Option,
) Stream[T]

ThroughErr applies fn to each item in s, records returned worker errors in the stream state, and returns another Stream[T]. It exists to make same-type stage segments read fluently in a chain.

type TimeoutOption

type TimeoutOption func() context.Context

TimeoutOption supplies the parent context for a timeout call.

func WithContext

func WithContext(ctx context.Context) TimeoutOption

WithContext sets the parent context for a timeout call.

Directories

Path Synopsis
examples
internal
pipeline
control
Package control exposes worker and concurrency control primitives used by flx stream, stage, and parallel execution APIs.
Package control exposes worker and concurrency control primitives used by flx stream, stage, and parallel execution APIs.

Jump to

Keyboard shortcuts

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