Documentation
¶
Overview ¶
Package pipeline provides composable stages for building concurrent data processing pipelines using channels. It handles context cancellation, error propagation, and goroutine lifecycle management automatically.
Index ¶
- type Composer
- type Config
- type Context
- type InputChannel
- type MultiChannelReceiver
- func (m MultiChannelReceiver[T]) At(index int) <-chan T
- func (m MultiChannelReceiver[T]) Iter() iter.Seq[<-chan T]
- func (m MultiChannelReceiver[T]) Len() int
- func (m MultiChannelReceiver[T]) SinkAt(ctx context.Context, index int) []T
- func (m MultiChannelReceiver[T]) SinkAtIter(ctx context.Context, index int) iter.Seq[T]
- type MultiChannelSender
- func (m MultiChannelSender[T]) At(index int) chan<- T
- func (m MultiChannelSender[T]) Iter() iter.Seq[chan<- T]
- func (m MultiChannelSender[T]) Len() int
- func (m MultiChannelSender[T]) Link(ctx Context, index int, in <-chan T) error
- func (m MultiChannelSender[T]) LinkAll(ctx Context, in MultiChannelReceiver[T]) error
- func (m MultiChannelSender[T]) Send(ctx context.Context, index int, values ...T) error
- func (m MultiChannelSender[T]) SendRoundRobin(ctx context.Context, values ...T) error
- func (m MultiChannelSender[T]) SendToAll(ctx context.Context, values ...T) error
- type OutputChannel
- type Pipeline
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Composer ¶ added in v3.2.0
func (Composer[In, Out]) Inputs ¶ added in v3.2.0
func (c Composer[In, Out]) Inputs() MultiChannelReceiver[In]
func (Composer[In, Out]) Outputs ¶ added in v3.2.0
func (c Composer[In, Out]) Outputs() MultiChannelSender[Out]
type Config ¶ added in v3.2.0
type Config[In any, Out any] struct { Name string InputChannels int InputBufferSize int OutputChannels int OutputBufferSize int Composer func(Composer[In, Out]) }
Config defines the make up of a pipeline and is required for construction of it
type InputChannel ¶ added in v3.2.0
type InputChannel[T any] interface { ~chan T | ~<-chan T }
type MultiChannelReceiver ¶
type MultiChannelReceiver[T any] []<-chan T
func NewMultiChannelReceiver ¶ added in v3.2.0
func NewMultiChannelReceiver[T any, C InputChannel[T]](in ...C) MultiChannelReceiver[T]
func (MultiChannelReceiver[T]) At ¶
func (m MultiChannelReceiver[T]) At(index int) <-chan T
func (MultiChannelReceiver[T]) Iter ¶
func (m MultiChannelReceiver[T]) Iter() iter.Seq[<-chan T]
func (MultiChannelReceiver[T]) Len ¶
func (m MultiChannelReceiver[T]) Len() int
func (MultiChannelReceiver[T]) SinkAt ¶ added in v3.2.0
func (m MultiChannelReceiver[T]) SinkAt(ctx context.Context, index int) []T
func (MultiChannelReceiver[T]) SinkAtIter ¶ added in v3.2.0
type MultiChannelSender ¶
type MultiChannelSender[T any] []chan<- T
func NewMultiChannelSender ¶ added in v3.2.0
func NewMultiChannelSender[T any, C OutputChannel[T]](out ...C) MultiChannelSender[T]
func (MultiChannelSender[T]) At ¶
func (m MultiChannelSender[T]) At(index int) chan<- T
func (MultiChannelSender[T]) Iter ¶
func (m MultiChannelSender[T]) Iter() iter.Seq[chan<- T]
func (MultiChannelSender[T]) Len ¶
func (m MultiChannelSender[T]) Len() int
func (MultiChannelSender[T]) Link ¶ added in v3.2.0
func (m MultiChannelSender[T]) Link(ctx Context, index int, in <-chan T) error
func (MultiChannelSender[T]) LinkAll ¶ added in v3.2.0
func (m MultiChannelSender[T]) LinkAll(ctx Context, in MultiChannelReceiver[T]) error
func (MultiChannelSender[T]) Send ¶
func (m MultiChannelSender[T]) Send(ctx context.Context, index int, values ...T) error
func (MultiChannelSender[T]) SendRoundRobin ¶ added in v3.1.0
func (m MultiChannelSender[T]) SendRoundRobin(ctx context.Context, values ...T) error
type OutputChannel ¶ added in v3.2.0
type OutputChannel[T any] interface { ~chan T | ~chan<- T }
type Pipeline ¶
Pipeline coordinates concurrent processing stages, managing their lifecycle and propagating errors and cancellation signals across all stages.
func NewPipeline ¶
func NewPipeline[In any, Out any](ctx context.Context, cfg Config[In, Out]) (*Pipeline[In, Out], context.Context)
NewPipeline creates a new Pipeline and a derived context for coordinating pipeline stages. The returned context is cancelled when any stage encounters an error. Use the returned Pipeline to register stages and wait for completion.
func (*Pipeline[In, Out]) CloseAllInputs ¶
func (p *Pipeline[In, Out]) CloseAllInputs()
CloseAllInputs will close all of the input channels
func (*Pipeline[In, Out]) Inputs ¶
func (p *Pipeline[In, Out]) Inputs() MultiChannelSender[In]
func (*Pipeline[In, Out]) Outputs ¶
func (p *Pipeline[In, Out]) Outputs() MultiChannelReceiver[Out]