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
- Variables
- type CSVConfig
- type Config
- type DocumentHandler
- type Logger
- type NDJSONConfig
- type Option
- func WithBatchSize(n int) Option
- func WithCSVDelimiter(r rune) Option
- func WithCSVHasHeader(b bool) Option
- func WithChunkSize(n int) Option
- func WithInputFormat(f OutputFormat) Option
- func WithLogger(l Logger) Option
- func WithMaxOpenReaders(n int) Option
- func WithMaxRowErrors(n int, onSkip func(err error)) Option
- func WithOutputFormat(f OutputFormat) Option
- func WithParallelWorkers(n int) Option
- func WithReadBufferSize(n int) Option
- func WithSingleFileOutput(b bool) Option
- func WithTransforms(rules ...PathTransformRule) Option
- type OutputFormat
- type ParquetConfig
- type PathTransformRule
- type PathTransformer
- type Result
- type RunConfig
- type Source
- type Stats
- type TooManyRowErrorsError
- type Transformer
Constants ¶
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 ¶
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.
type DocumentHandler ¶
DocumentHandler receives one decoded document at a time. It may be called concurrently and must be safe for that.
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 ¶
WithBatchSize sets how many decoded documents a decode worker batches per send. n<=0 falls back to the default.
func WithCSVDelimiter ¶
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 ¶
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 ¶
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 ¶
WithLogger sets the logger that receives ProcessFile diagnostics.
func WithMaxOpenReaders ¶
WithMaxOpenReaders caps simultaneously-open Parquet row-group readers. n<=0 falls back to the configured worker count, not a fixed default.
func WithMaxRowErrors ¶
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 ¶
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 ¶
WithReadBufferSize sets the chunked readers' (NDJSON, CSV, TSV) read buffer size, in bytes. n<=0 falls back to the default.
func WithSingleFileOutput ¶
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.
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 ¶
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 ¶
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.