Documentation
¶
Overview ¶
Codecs encode a batch entry's rows into its data file: ndjson and gzipped ndjson here, and others registered by the packages that implement them.
Package spool hands batches of stream writes from any process to the one that owns a store, as directories of data files committed by rename.
A spool directory holds tmp/ (batches being written), incoming/ (published batches, oldest first by name), failed/ (batches that could not be ingested, with an error.json) and trash/ (ingested batches awaiting collection). A batch is published by writing it under tmp/ and renaming it into incoming/: the rename is its only commit point, so a reader never sees half a batch.
Loading a published batch back: the manifest is checked, every data file is verified against its digest and row count, and the rows are decoded.
The manifest: a published batch's description, with explicit JSON names and string enums so a build reading another build's batch reads it the same.
Row values as a sqlite kind table stores them, so a spooled row and the same row appended directly store identical values.
Publishing a batch: its data files and manifest are written and made durable under tmp/, then the batch directory is renamed into incoming/.
The bulk Writer: a producer's writes gathered into batches cut at a row and byte cap, each delivered whole and in order.
Index ¶
- Constants
- Variables
- func Digest(data []byte) string
- func Formats() []string
- func Normalize(schema recordstore.KindSchema, row recordstore.Row) (recordstore.Row, error)
- func Register(format string, codec Codec)
- type Codec
- type CollectOptions
- type Dir
- func (d Dir) Collect(now time.Time, options CollectOptions) error
- func (d Dir) Fail(name string, reason error) error
- func (d Dir) Incoming() ([]string, error)
- func (d Dir) Load(name string) (recordstore.Batch, error)
- func (d Dir) Path() string
- func (d Dir) Publish(batch recordstore.Batch, format string) (string, error)
- func (d Dir) Publisher(format string) func(context.Context, recordstore.Batch) error
- func (d Dir) Trash(name string) error
- type Producer
- type Writer
- type WriterOptions
Constants ¶
const ( FormatNDJSON = "ndjson" FormatNDJSONGzip = "ndjson.gz" )
const ManifestFormat = 1
ManifestFormat is the manifest format this build writes and reads.
Variables ¶
var ErrManifestFormat = errors.New("spool manifest format unsupported")
ErrManifestFormat reports a batch published in a manifest format this build does not read.
Functions ¶
func Normalize ¶
func Normalize(schema recordstore.KindSchema, row recordstore.Row) (recordstore.Row, error)
Normalize converts row into the values sqlitetable.Value stores for schema's columns: a structured column's value as its JSON text, a time as sqlitetable.TimeLayout, and any other value the driver cannot store as its JSON text. A key schema does not declare is an error, but for a dynamic kind, whose key gets the value its inferred column stores — a null none, an object or an array itself, for the owner to infer its column from.
Types ¶
type Codec ¶
type Codec interface {
Encode(w io.Writer, rows []recordstore.Row) error
Decode(r io.Reader, fn func(recordstore.Row) error) error
}
Codec writes and reads the rows of one data file. Decode hands numbers over as json.Number, so an int64 survives the trip whole.
type CollectOptions ¶
type CollectOptions struct {
// TmpAge is how long a batch may stay unpublished in tmp, by its last
// change, before it is taken as abandoned. Zero is 24 hours.
TmpAge time.Duration
// FailedAge is how long a failed batch is kept for inspection. Zero is 30
// days.
FailedAge time.Duration
}
CollectOptions say how long abandoned and failed batches are kept.
type Dir ¶
type Dir struct {
// contains filtered or unexported fields
}
Dir is a spool directory.
func OpenDir ¶
OpenDir opens the spool directory at path, creating it and its subdirectories private to the user.
func (Dir) Collect ¶
func (d Dir) Collect(now time.Time, options CollectOptions) error
Collect removes every trashed batch, the failed batches older than FailedAge and the tmp entries older than TmpAge, as of now.
func (Dir) Fail ¶
Fail moves a batch that cannot be ingested to failed, with reason recorded beside it in error.json.
func (Dir) Load ¶
func (d Dir) Load(name string) (recordstore.Batch, error)
Load reads the published batch name. An error means the batch as a whole cannot be ingested — ErrManifestFormat for a manifest format this build does not read — and belongs in failed/.
func (Dir) Publish ¶
Publish writes batch into the spool, its rows in format, and returns the name it was published under in incoming/. Every kind the batch appends to must have its schema in batch.Schemas: the rows are normalized against it, and the owner ingests them under it. Nothing is published unless all of it is.
type Producer ¶
type Producer struct {
// contains filtered or unexported fields
}
Producer numbers one process's batches: every batch it identifies gets the next seq, so an owner can ingest a producer's batches in order.
func NewProducer ¶
NewProducer identifies the batches of this process as instance, made by build.
func (*Producer) Next ¶
func (p *Producer) Next() recordstore.Producer
Next is the identity of the producer's next batch.
type Writer ¶
type Writer struct {
// contains filtered or unexported fields
}
Writer gathers writes into batches. A write it has taken is delivered by a later Append that fills a batch or by Flush; a bulk write spanning batches is not atomic. It is not safe for concurrent use.
func NewWriter ¶
func NewWriter(options WriterOptions) (*Writer, error)
NewWriter returns a Writer delivering through options.Deliver.
func (*Writer) Append ¶
Append gathers rows for stream, cutting and delivering a batch whenever one fills.
type WriterOptions ¶
type WriterOptions struct {
// Producer identifies every batch the writer delivers. Required.
Producer *Producer
// Schema resolves the kinds appended to; a batch carries the schema of
// every kind it appends to. Required.
Schema recordstore.SchemaResolver
// Deliver hands a batch on: Dir.Publisher, or a BatchAppender's
// AppendBatch. A batch it fails is delivered again, with the same id and
// seq, before any later one. Required.
Deliver func(context.Context, recordstore.Batch) error
// MaxRows and MaxBytes cap a batch's rows and their JSON size; zero is
// 50,000 rows and 64MiB. A batch always takes at least one row.
MaxRows int
MaxBytes int
}
WriterOptions configure a Writer.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package parquet is the spool's parquet codec: importing it registers the "parquet" format, for producers exporting large batches in columnar form.
|
Package parquet is the spool's parquet codec: importing it registers the "parquet" format, for producers exporting large batches in columnar form. |