streamio

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

README

streamio

A Go library and CLI for converting between NDJSON, Parquet, CSV, TSV, and Arrow IPC at scale, built for pipelines that move large datasets between storage and processing systems rather than one-off scripting. Any format can be read and written, decoding and dispatch run concurrently on a shared worker pool, and when the input and output formats match, bytes are streamed straight through with no decode step at all — reliably processing files far larger than available memory, fast.

Use cases

  • Bulk-loading columnar files into a datastore.
  • Converting exported NDJSON logs into columnar formats for analytics.
  • Reshaping CSV/TSV exports into a format a downstream system expects.
  • Any pipeline where the bottleneck is decode/encode throughput on multi-gigabyte files, not one-time convenience.

Install

Prebuilt CLI binaries for Linux and macOS (amd64/arm64) and Windows (amd64) are on the releases page. Download the archive for your platform, extract it, and put the streamio binary on your PATH.

To build from source instead:

go install github.com/AlexAslan/streamio/cmd/streamio@latest

Getting started

Convert a file (output format defaults to json; input format is detected from the file extension):

./streamio convert --in data.ndjson --out data.parquet --out-format parquet

--out-format accepts json, parquet, csv, tsv, or arrow. See ./streamio convert --help for every flag, including CSV/TSV delimiter and header options, batch/worker sizing, and --max-row-errors.

As a library:

go get github.com/AlexAslan/streamio
import "github.com/AlexAslan/streamio"

result, err := streamio.ProcessFile(ctx, "data.parquet", handler,
    streamio.WithOutputFormat(streamio.FormatJSON),
)

See Configuration for the full streamio.Option reference.

Docs

  • Architecture — entry point, routing, and the decode/dispatch pool.
  • Formats — NDJSON, Parquet, and Arrow IPC read/write routes.
  • record/formatio — the generic cross-format seam every format package implements.
  • Configuration — streamio.Option reference and non-file-path input.
  • Benchmarks — measured throughput and memory across every format pair.
  • Package layout — where everything lives.

Contributing

See CONTRIBUTING.md. Security issues: SECURITY.md.

License

MIT. Release archives also include THIRD_PARTY_NOTICES.md with the licenses of bundled dependencies.

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

View Source
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

View Source
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.

View Source
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

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

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

func DetectInputFormat(path string) Format

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 Logger

type Logger = options.Logger

Logger receives ProcessFile diagnostics.

type Option

type Option = options.Option

Option configures an optional aspect of a ProcessFile run.

func WithBatchSize

func WithBatchSize(n int) Option

WithBatchSize sets how many documents are batched together per dispatch. 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.

func WithChunkSize

func WithChunkSize(n int) Option

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

func WithInputFormat(f Format) Option

WithInputFormat sets the source file's format explicitly, instead of letting ProcessFile infer it 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 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

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), 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

func WithOutputFormat(f Format) Option

WithOutputFormat sets the format documents are delivered in.

func WithParallelWorkers

func WithParallelWorkers(n int) Option

WithParallelWorkers sets the number of concurrent decode workers. n<=0 falls back to the default.

func WithReadBufferSize

func WithReadBufferSize(n int) Option

WithReadBufferSize sets the read buffer size, in bytes, used while scanning the input file. 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. 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

type Result = options.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

type Source = options.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 Stats

type Stats = options.Stats

Stats holds counters describing a completed ProcessFile run.

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.

Jump to

Keyboard shortcuts

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