Documentation
¶
Overview ¶
Package pool runs the decode/dispatch worker pool shared by every streamio route: each format only supplies a formatio.RawSource (RunRaw) or a decoder/encoder pair (RunRecords).
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func RunRaw ¶
func RunRaw( ctx context.Context, cfg options.Config, src formatio.RawSource, sink options.DocumentHandler, ) (options.Result, error)
RunRaw runs the pool over the raw-passthrough path: cfg.Run.Workers goroutines call src.DecodeRaw against as many dispatch goroutines calling sink, until done or an error stops every worker.
src is not closed here — the caller that built it owns it.
func RunRecords ¶
func RunRecords( ctx context.Context, cfg options.Config, dec formatio.RecordDecoder, newEncoder func() (formatio.RecordEncoder, error), transformer options.Transformer, sink options.DocumentHandler, ) (options.Result, error)
RunRecords runs the pool over the generic cross-format path: decode the input into canonical records and encode each batch in the requested output format, using the same fan-out/fan-in skeleton RunRaw uses.
Decode workers come from dec (split across up to cfg.Run.Workers when it supports that, one otherwise), each paired with its own encoder from newEncoder since an encoder holds scratch and isn't safe for concurrent use. dec is not closed here — the caller that built it owns it.
Types ¶
type DecodeStats ¶
type DecodeStats = formatio.DecodeStats
DecodeStats is formatio.DecodeStats under the name the pool's own callers use.