Documentation
¶
Overview ¶
Package streamio converts between document stream formats (NDJSON, Parquet), decoding an input file into canonical records and re-encoding it in the requested output format.
Index ¶
- Constants
- Variables
- type DocumentHandler
- type Format
- type Logger
- 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 Format) Option
- func WithLogger(l Logger) Option
- func WithMaxOpenReaders(n int) Option
- func WithMaxRowErrors(n int, onSkip func(err error)) Option
- func WithOutputFormat(f Format) Option
- func WithParallelWorkers(n int) Option
- func WithReadBufferSize(n int) Option
- func WithSingleFileOutput(b bool) Option
- func WithTransforms(rules ...TransformRule) Option
- type Result
- type Source
- type Stats
- type TooManyRowErrorsError
- type TransformRule
Constants ¶
const ( // FormatJSON selects newline-delimited JSON. FormatJSON = options.FormatJSON // FormatParquet selects Apache Parquet. FormatParquet = options.FormatParquet // FormatCSV selects comma-separated values. FormatCSV = options.FormatCSV // FormatTSV selects tab-separated values. FormatTSV = options.FormatTSV // FormatArrow selects Apache Arrow IPC (the file/random-access variant). FormatArrow = options.FormatArrow )
Variables ¶
var ErrNoConversionPath = errors.New("streamio: no conversion path")
ErrNoConversionPath is returned when the requested output format can't be produced from the input file's native format: the pair needs a decoder for the input and an encoder for the output, and at least one of the two isn't registered. Wrapped with both formats by ProcessFile.
var ErrTooManyRowErrors = options.ErrTooManyRowErrors
ErrTooManyRowErrors is TooManyRowErrorsError's sentinel: check for this failure with errors.Is without depending on the concrete type.
Functions ¶
This section is empty.
Types ¶
type DocumentHandler ¶
DocumentHandler is invoked once per output document. It may be called concurrently by multiple workers, so implementations must be safe for that.
type Format ¶
type Format = options.OutputFormat
Format identifies an input or output document encoding.
func DetectInputFormat ¶
DetectInputFormat returns the native format of the file at path, derived from its extension (case-insensitive): any extension not named below means NDJSON. ProcessFile calls this itself when WithInputFormat isn't given; it's exported so a caller that needs the format ahead of time (for example, to build format-specific options before calling ProcessFile) can run the same detection instead of duplicating or guessing at it.
type Option ¶
Option configures an optional aspect of a ProcessFile run.
func WithBatchSize ¶
WithBatchSize sets how many documents are batched together per dispatch. 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.
func WithChunkSize ¶
WithChunkSize sets the fixed work-queue chunk size, in bytes, used to split NDJSON input across workers. n<=0 falls back to the default.
func WithInputFormat ¶
WithInputFormat sets the source file's format explicitly, instead of letting ProcessFile infer it from the file extension.
func WithLogger ¶
WithLogger sets the logger that receives ProcessFile diagnostics.
func WithMaxOpenReaders ¶
WithMaxOpenReaders caps the number of simultaneously open Parquet row-group readers. n<=0 falls back to the configured number of parallel workers, 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), matching streamio's original behavior. onSkip, if non-nil, is called once per row skipped under the cap. Only NDJSON and CSV/TSV field errors count against the cap — a CSV/TSV syntax error and every Parquet decode error always fail the run regardless of n.
Exceeding the cap fails the run with a *TooManyRowErrorsError wrapping every row error collected up to and including the one that exceeded it; Result.Stats.RowsSkipped is still populated on that failure path.
func WithOutputFormat ¶
WithOutputFormat sets the format documents are delivered in.
func WithParallelWorkers ¶
WithParallelWorkers sets the number of concurrent decode workers. n<=0 falls back to the default.
func WithReadBufferSize ¶
WithReadBufferSize sets the read buffer size, in bytes, used while scanning the input file. 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. Forces decoding to a single worker.
func WithTransforms ¶
func WithTransforms(rules ...TransformRule) Option
WithTransforms applies a set of rename/drop rules to every decoded record before it is encoded into the requested output format.
type Result ¶
Result reports the outcome of a ProcessFile run.
func ProcessFile ¶
func ProcessFile( ctx context.Context, path string, handler DocumentHandler, opts ...Option, ) (Result, error)
ProcessFile reads the file at path and calls handler once per document, delivering each one in cfg.OutputFormat (options.FormatJSON unless WithOutputFormat says otherwise); see Option for tuning.
func ProcessReaderAt ¶
func ProcessReaderAt( ctx context.Context, src Source, handler DocumentHandler, opts ...Option, ) (Result, error)
ProcessReaderAt reads src and calls handler once per document, the same way ProcessFile does for an on-disk file — src need not be backed by a real file, only support ReadAt over its declared Size, which is what every format's chunked-parallel decoder actually requires.
type Source ¶
Source is a sized, randomly-addressable byte source ProcessReaderAt reads from: an *os.File, an in-memory buffer, or anything else providing ReadAt over a known span. Name labels it in diagnostics and error text and need not be a real file path.
type TooManyRowErrorsError ¶
type TooManyRowErrorsError = options.TooManyRowErrorsError
TooManyRowErrorsError is returned when WithMaxRowErrors' limit is exceeded: it wraps every row error collected up to and including the one that exceeded it.
type TransformRule ¶
type TransformRule = options.PathTransformRule
TransformRule is a single rename or drop rule built by RenamePath or DropPath.
func DropPath ¶
func DropPath(path string) TransformRule
DropPath builds a transform rule that removes one field path.
func RenamePath ¶
func RenamePath(from, to string) TransformRule
RenamePath builds a transform rule that renames one field path.
Directories
¶
| Path | Synopsis |
|---|---|
|
cmd
|
|
|
streamio
command
Command streamio-cli converts a document file between the formats streamio supports (NDJSON, Parquet, CSV, TSV, and Arrow IPC).
|
Command streamio-cli converts a document file between the formats streamio supports (NDJSON, Parquet, CSV, TSV, and Arrow IPC). |
|
internal
|
|
|
arrowio
Package arrowio supplies streamio's Apache Arrow IPC (file format) capabilities: decoding into canonical records and encoding records back into Arrow IPC documents.
|
Package arrowio supplies streamio's Apache Arrow IPC (file format) capabilities: decoding into canonical records and encoding records back into Arrow IPC documents. |
|
csvio
Package csvio owns streamio's CSV and TSV support: one package for both, since they differ only in field delimiter.
|
Package csvio owns streamio's CSV and TSV support: one package for both, since they differ only in field delimiter. |
|
formatio
Package formatio defines the per-format capabilities streamio composes: streaming a file's own bytes (RawSource), decoding into canonical records (RecordDecoder), and encoding records as another format's documents (RecordEncoder).
|
Package formatio defines the per-format capabilities streamio composes: streaming a file's own bytes (RawSource), decoding into canonical records (RecordDecoder), and encoding records as another format's documents (RecordEncoder). |
|
jsonio
Package jsonio owns streamio's JSON output path: encoding canonical record.Records as JSON documents, and decoding NDJSON back into them.
|
Package jsonio owns streamio's JSON output path: encoding canonical record.Records as JSON documents, and decoding NDJSON back into them. |
|
options
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.
|
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. |
|
parquetio
Package parquetio supplies streamio's Parquet capabilities: raw row-group passthrough, decoding into canonical records, and encoding records back into Parquet documents.
|
Package parquetio supplies streamio's Parquet capabilities: raw row-group passthrough, decoding into canonical records, and encoding records back into Parquet documents. |
|
pool
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).
|
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). |
|
record
Package record defines streamio's canonical, format-neutral row: the intermediate every cross-format conversion goes through between a decoder and an encoder.
|
Package record defines streamio's canonical, format-neutral row: the intermediate every cross-format conversion goes through between a decoder and an encoder. |