pool

package
v0.0.7 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 10 Imported by: 0

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

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.

Jump to

Keyboard shortcuts

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