options

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: 7 Imported by: 0

Documentation

Overview

Package options defines streamio's Config and the functional Option funcs that build it, along with the transform rules and format constants ProcessFile's callers configure a run with.

Index

Constants

View Source
const (

	// DispatchQueueDepthPerWorker sizes the channel connecting decode workers to dispatch
	// workers, per decode worker, for both Parquet and NDJSON. Not user-configurable — there's
	// no corresponding WithXxx option — so it's a plain exported constant, not a Config field.
	DispatchQueueDepthPerWorker = 4
)

Variables

View Source
var ErrTooManyRowErrors = errors.New("streamio: too many row errors")

ErrTooManyRowErrors is TooManyRowErrorsError's sentinel: wrapped so a caller can check for this failure with errors.Is without depending on the concrete type.

Functions

This section is empty.

Types

type CSVConfig

type CSVConfig struct {
	// Delimiter separates fields on a row. It defaults to ',' for CSV and '\t' for TSV.
	Delimiter rune
	// HasHeader says whether the first row names the columns rather than holding data.
	HasHeader bool
}

CSVConfig holds settings specific to reading or writing CSV or TSV.

type Config

type Config struct {
	// Logger receives ProcessFile diagnostics, when set.
	Logger Logger
	// Transformer changes each decoded record before it is encoded, when set.
	Transformer Transformer
	// OnRowErrorSkip, if set, is called once per row MaxRowErrors allows to be skipped, with the
	// decode error that row raised. It must be safe for concurrent use: one decode worker per
	// goroutine may call it.
	OnRowErrorSkip func(err error)
	// OptionErr is a deferred validation-error slot: an option func that fails validation stores
	// its error here instead of panicking, and New's caller checks it once after every option has
	// run.
	OptionErr error
	// InputFormat is the source file's format, or empty to detect it from the file extension.
	InputFormat OutputFormat
	// OutputFormat is the format documents are delivered in.
	OutputFormat OutputFormat
	// Run holds settings for the decode/dispatch pool, shared by every format.
	Run RunConfig
	// NDJSON holds settings only the NDJSON format reads.
	NDJSON NDJSONConfig
	// Parquet holds settings only the Parquet format reads.
	Parquet ParquetConfig
	// CSV holds settings the CSV and TSV formats read; the two differ only in CSV.Delimiter.
	CSV CSVConfig
	// MaxRowErrors caps how many rows may fail to decode and be skipped before the run fails,
	// mirroring BigQuery's load-job max_bad_records: its zero value means fail immediately on the
	// first row error, matching every version of streamio before this option existed. Only a
	// decoder error that identifies itself as safely skippable (see formatio.RowError) counts
	// against this cap — currently NDJSON and CSV/TSV field errors, not Parquet or a CSV/TSV syntax
	// error, which always fail the run regardless of MaxRowErrors.
	MaxRowErrors int
}

Config holds the options passed to Process.

func New

func New(opts ...Option) Config

New returns a Config with the given options applied over its defaults.

type DocumentHandler

type DocumentHandler func(ctx context.Context, doc []byte) error

DocumentHandler receives one decoded document at a time. It may be called concurrently and must be safe for that.

type Logger

type Logger interface {
	Printf(format string, args ...any)
}

Logger provides the logging capabilities for the module.

type NDJSONConfig

type NDJSONConfig struct {
	// ChunkSize is the reader's fixed work-queue chunk size, in bytes.
	ChunkSize int
}

NDJSONConfig holds settings specific to reading or writing NDJSON.

type Option

type Option func(*Config)

Option configures an optional aspect of a Config.

func WithBatchSize

func WithBatchSize(n int) Option

WithBatchSize sets how many decoded documents a decode worker batches per send. n<=0 falls back to the default.

func WithCSVDelimiter

func WithCSVDelimiter(r rune) Option

WithCSVDelimiter sets the field delimiter CSV or TSV reads and writes. Pass ',' for CSV or '\t' for TSV; the zero value falls back to the format's own default.

func WithCSVHasHeader

func WithCSVHasHeader(b bool) Option

WithCSVHasHeader says whether a CSV or TSV file's first row names its columns rather than holding data. It defaults to false, so a headerless file must be explicitly opted out of.

func WithChunkSize

func WithChunkSize(n int) Option

WithChunkSize sets the NDJSON reader's fixed work-queue chunk size, in bytes. n<=0 falls back to the default.

func WithInputFormat

func WithInputFormat(f OutputFormat) Option

WithInputFormat sets the source file's format explicitly. Omitting it, or passing the zero value, lets ProcessFile infer the input format from the file extension.

func WithLogger

func WithLogger(l Logger) Option

WithLogger sets the logger that receives ProcessFile diagnostics.

func WithMaxOpenReaders

func WithMaxOpenReaders(n int) Option

WithMaxOpenReaders caps simultaneously-open Parquet row-group readers. n<=0 falls back to the configured worker count, not a fixed default.

func WithMaxRowErrors

func WithMaxRowErrors(n int, onSkip func(err error)) Option

WithMaxRowErrors caps how many rows may fail to decode and be skipped before the run fails, mirroring BigQuery's load-job max_bad_records: n<=0 means fail immediately on the first row error (the default). onSkip, if non-nil, is called once per row skipped under the cap.

func WithOutputFormat

func WithOutputFormat(f OutputFormat) Option

WithOutputFormat sets the format documents are delivered in. Omitting it means FormatJSON.

func WithParallelWorkers

func WithParallelWorkers(n int) Option

WithParallelWorkers sets the number of decode workers: row-group workers for Parquet, or byte-range chunk workers for NDJSON. n<=0 falls back to the default.

func WithReadBufferSize

func WithReadBufferSize(n int) Option

WithReadBufferSize sets the chunked readers' (NDJSON, CSV, TSV) read buffer size, in bytes. n<=0 falls back to the default.

func WithSingleFileOutput

func WithSingleFileOutput(b bool) Option

WithSingleFileOutput requests one continuous output document spanning the whole run, for formats that support it, instead of one document per batch. See RunConfig.SingleFileOutput.

func WithTransforms

func WithTransforms(rules ...PathTransformRule) Option

WithTransforms applies a set of rename/drop rules to every decoded record before encoding. An invalid rule set is stored on Config.OptionErr rather than returned directly.

type OutputFormat

type OutputFormat string

OutputFormat names the representation a caller wants documents delivered in, and also a file's own native format.

const (
	// FormatJSON delivers each document as a JSON object — the native format of an NDJSON input,
	// and the default for every other input.
	FormatJSON OutputFormat = "json"

	// FormatParquet delivers Parquet bytes. Against a Parquet input this is the raw-passthrough
	// case: each document is one standalone, self-contained single-row-group Parquet file.
	FormatParquet OutputFormat = "parquet"

	// FormatCSV delivers comma-separated rows.
	FormatCSV OutputFormat = "csv"

	// FormatTSV delivers tab-separated rows.
	FormatTSV OutputFormat = "tsv"

	// FormatArrow delivers Apache Arrow IPC bytes (the file/random-access variant, not the
	// streaming-only one).
	FormatArrow OutputFormat = "arrow"
)

func (OutputFormat) String

func (f OutputFormat) String() string

String returns the format's name, or "unspecified" for the zero value. Config.New normalises the zero value to FormatJSON, so a Config in hand never reports "unspecified".

type ParquetConfig

type ParquetConfig struct {
	// MaxOpenReaders caps simultaneously-open row-group readers.
	MaxOpenReaders int
}

ParquetConfig holds settings specific to reading or writing Parquet.

type PathTransformRule

type PathTransformRule struct {
	// contains filtered or unexported fields
}

PathTransformRule is one declarative path-level transform. Use RenamePath and DropPath to build rules so invalid partial rules cannot be constructed accidentally.

func DropPath

func DropPath(name string) PathTransformRule

DropPath returns a transform rule that removes one field path. Dotted paths address nested map entries, for example "attributes.debug".

func RenamePath

func RenamePath(from, to string) PathTransformRule

RenamePath returns a transform rule that changes one field path. Dotted paths address nested map entries, for example "attributes.service".

type PathTransformer

type PathTransformer struct {
	// contains filtered or unexported fields
}

PathTransformer applies compiled rename/drop rules in one pass over a record. Rules are stored as a path trie so per-record work checks only the fields that can match at each level.

func NewPathTransformer

func NewPathTransformer(rules ...PathTransformRule) (*PathTransformer, error)

NewPathTransformer compiles field rules once, before any records are processed.

func (*PathTransformer) Transform

func (t *PathTransformer) Transform(rec record.Record) (record.Record, error)

Transform applies configured rename/drop actions to rec.

type Result

type Result struct {
	Stats Stats
}

Result reports what one ProcessFile call did.

type RunConfig

type RunConfig struct {
	// ReadBufferSize is the read buffer size, in bytes, used while scanning the input file.
	ReadBufferSize int
	// BatchSize is how many decoded documents a decode worker batches into one send to a dispatch
	// worker.
	BatchSize int
	// Workers is the number of concurrent decode workers.
	Workers int
	// SingleFileOutput requests one continuous output document spanning the whole run, for formats
	// that support it, instead of one document per batch.
	//
	// It forces Workers down to a single decode worker: see formatio.FinalizableRecordEncoder.
	SingleFileOutput bool
}

RunConfig holds settings for the decode/dispatch pool that apply no matter which format is being read or written.

type Source

type Source struct {
	Reader io.ReaderAt
	Name   string
	Size   int64
}

Source is a sized, randomly-addressable byte source: what every format's chunked-parallel decoder actually needs, satisfied by an *os.File or anything else providing ReadAt over a known span. Name labels it in diagnostics and error text; it need not be a real file path.

type Stats

type Stats struct {
	// RowsRead is how many input rows were decoded, regardless of how many documents they became.
	RowsRead int64
	// RowsSkipped is how many rows WithMaxRowErrors allowed to be skipped. Always zero under the
	// default MaxRowErrors of 0.
	RowsSkipped int64
	// DocumentsDispatched is how many documents were handed to the caller's handler. It can differ
	// from RowsRead when an encoder renders a whole batch as one document (Parquet) rather than one
	// document per row (NDJSON).
	DocumentsDispatched int64
	// ReadDuration sums every worker's decode+encode time.
	ReadDuration time.Duration
	// DispatchDuration sums every worker's time spent in the caller's handler.
	DispatchDuration time.Duration
}

Stats reports what one ProcessFile call did.

type TooManyRowErrorsError

type TooManyRowErrorsError struct {
	// Errors holds every row error collected before the run stopped, in the order encountered
	// across every decode worker sharing the run's limit.
	Errors []error
	// Limit is the MaxRowErrors value the run was configured with.
	Limit int
}

TooManyRowErrorsError is returned when the number of skipped row errors exceeds MaxRowErrors: the (MaxRowErrors+1)th row error stops the run rather than being skipped, and this wraps every row error collected up to and including that one — not just a count, since Stats.RowsSkipped (still populated on this failure path) already gives the count on its own.

func (*TooManyRowErrorsError) Error

func (e *TooManyRowErrorsError) Error() string

func (*TooManyRowErrorsError) Unwrap

func (e *TooManyRowErrorsError) Unwrap() error

type Transformer

type Transformer interface {
	Transform(record.Record) (record.Record, error)
}

Transformer changes a decoded canonical record before it is encoded into the requested output format. It is only used on the record path; raw passthrough never decodes records and therefore cannot transform them.

Jump to

Keyboard shortcuts

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