Documentation
¶
Overview ¶
Package spool holds a from-proto sink's rows on local disk, pre-encoded into whatever the target database can take unchanged, and applies whole segments from a separate goroutine.
The point is not raw throughput — pre-encoding is worth about 8% on its own, see sink/sql/db_proto/benchmarks. It is that the stream stops waiting for the database. Substreams throughput is paid for, so a slow or stalled database should cost disk, not download progress, and blocks already paid for should survive a restart rather than be streamed again.
What the bytes on disk look like, and how a sealed segment reaches the server, are the driver's business: this package owns the segment lifecycle, the disk budget and the sizing, and delegates the rest through Codec and Applier.
Index ¶
- func SanitizeFileName(name string) string
- func WriteManifest(dir string, manifest *Manifest) error
- type Applier
- type Codec
- type Format
- type FrameReader
- type FrameWriter
- type Manifest
- type Options
- type SealReason
- type SegmentWriter
- type Spool
- func (b *Spool) AppliedBlock() uint64
- func (b *Spool) BlocksBuffered() int64
- func (b *Spool) BytesOnDisk() int64
- func (b *Spool) Close(ctx context.Context) error
- func (b *Spool) Drain(ctx context.Context) error
- func (b *Spool) Insert(table string, values []any) error
- func (b *Spool) MaybeSeal(ctx context.Context) error
- func (b *Spool) RecordBlock(blockNum uint64)
- func (b *Spool) RecordCursor(cursor string)
- func (b *Spool) Seal(ctx context.Context, reason SealReason) error
- func (b *Spool) Stats() Stats
- type Stats
- type TableRecord
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func SanitizeFileName ¶
SanitizeFileName keeps a table name usable as a path component.
func WriteManifest ¶
WriteManifest is exported for a codec that needs to re-seal a segment.
Types ¶
type Applier ¶
type Applier interface {
// EnsureSchema creates whatever bookkeeping the applier needs, once, before any
// segment is written.
EnsureSchema(ctx context.Context) error
// AlreadyApplied reports a segment recovery found on disk that the database already
// holds, so it can be dropped rather than replayed.
//
// A driver with transactions answers this from its own record of applied segments. One
// without answers it from how far the stored cursor got: replaying a segment it had in
// fact applied duplicates exactly the rows re-streaming it would have duplicated, which
// is the guarantee such a driver already ships with.
AlreadyApplied(ctx context.Context, manifest *Manifest) (bool, error)
// Apply loads one segment and advances the cursor with it.
Apply(ctx context.Context, dir string, manifest *Manifest) error
}
Applier sends sealed segments to the database.
type Codec ¶
type Codec interface {
// Format is what gets recorded in the manifest.
Format() Format
// OpenSegment prepares the writers for a new segment in dir, which already exists.
OpenSegment(dir string) (SegmentWriter, error)
// Verify checks a sealed segment's files against its manifest. Only the manifest is
// fsynced on seal, so a machine crash can leave a data file short and this is what
// catches it before anything reaches the database.
Verify(dir string, manifest *Manifest) error
}
Codec owns what a segment looks like on disk.
type Format ¶
type Format string
Format names an on-disk layout. It is stored in the manifest so a segment written by an earlier run is read back the way it was written.
const ( // FormatPGCopy is one file per table in PostgreSQL's binary COPY wire format. It is // also what a segment written before formats existed carries, hence the zero value. FormatPGCopy Format = "" // FormatTuples is one file per table of rendered SQL value tuples, which the applier // wraps into multi-row INSERTs. Tables are applied in foreign key order. FormatTuples Format = "tuples" // FormatRowLog is a single interleaved file of (table, rendered tuple) in walk order. // // It exists because grouping rows by table is exactly what row-insert mode cannot do. // That mode is the fallback for a schema whose foreign keys form a cycle, and a cycle // has no table order that keeps a parent ahead of its children — only the walk does, // so only the walk's own order can be replayed. FormatRowLog Format = "rowlog" // FormatValues is one file per table of typed row values, appended straight back into // the driver's own column builders at apply time. ClickHouse inserts columnar and // typed rather than as SQL text, so rendering to literals would change both the insert // path and how types are handled. FormatValues Format = "values" )
type FrameReader ¶
type FrameReader struct {
// contains filtered or unexported fields
}
FrameReader reads back what FrameWriter produced.
func OpenFrameReader ¶
func OpenFrameReader(path string) (*FrameReader, error)
func (*FrameReader) Close ¶
func (r *FrameReader) Close() error
func (*FrameReader) ReadField ¶
func (r *FrameReader) ReadField() (string, error)
ReadField returns the next field, or io.EOF once the file is exhausted.
type FrameWriter ¶
type FrameWriter struct {
// contains filtered or unexported fields
}
FrameWriter writes length-prefixed records.
Framing rather than lines because a rendered tuple carries SQL literals, and a text column holding a newline would end a line in the middle of a row. The length check that recovery does against the manifest then covers a torn write, the same way the binary COPY trailer does for FormatPGCopy.
func NewFrameWriter ¶
func NewFrameWriter(file *os.File) *FrameWriter
func (*FrameWriter) Bytes ¶
func (w *FrameWriter) Bytes() int64
func (*FrameWriter) Close ¶
func (w *FrameWriter) Close() error
func (*FrameWriter) Rows ¶
func (w *FrameWriter) Rows() int64
func (*FrameWriter) WriteRecord ¶
func (w *FrameWriter) WriteRecord(fields ...string) error
WriteRecord appends one record made of the given fields, each length-prefixed.
type Manifest ¶
type Manifest struct {
FirstBlock uint64 `json:"first_block"`
LastBlock uint64 `json:"last_block"`
Cursor string `json:"cursor"`
Tables []TableRecord `json:"tables"`
Sealed bool `json:"sealed"`
// Format says how the rows are laid out. Absent means FormatPGCopy, which is what
// every segment written before formats existed carries.
Format Format `json:"format,omitempty"`
// LogFile and LogBytes describe the single interleaved file of FormatRowLog. The
// table records still carry the column layout, since replaying a row needs it.
LogFile string `json:"log_file,omitempty"`
LogBytes int64 `json:"log_bytes,omitempty"`
}
Manifest is a segment's commit record. Its presence, parseable and sealed, is what makes a segment eligible to be applied; anything else is a torn write from a crash.
func (*Manifest) BlockCount ¶
BlockCount is how many blocks the segment covers. A segment carrying nothing but a cursor — every block in the flush produced no output — covers none.
func (*Manifest) CursorOnly ¶
CursorOnly reports a segment that carries no rows. It exists so that a stretch of blocks whose module output is empty still advances the cursor, rather than being streamed, and paid for, again on restart.
type Options ¶
type Options struct {
// Dir is where segments live. Required.
Dir string
// MaxBytes is the disk quota, and the only bound on how far ahead of the database the
// stream may run. Writes are held once the spool holds this much waiting to be
// applied, which is what turns a slow database into backpressure rather than a full
// disk. Zero picks 8GiB.
MaxBytes int64
// WriteTargetDuration is how long one commit to the database should take. The sizer
// measures each commit and steers the next segment toward it. Zero picks 3s.
WriteTargetDuration time.Duration
// SegmentMaxBytes is the ceiling the sizer may choose, whatever the target duration
// would allow. Zero picks 512MiB. The floor is segmentFloorBytes and not
// configurable, see sizer.go.
SegmentMaxBytes int64
// MaxIdle commits the open segment once no new row has reached it for this long, short
// of its size target. Without it a stalled stream sits on those rows indefinitely,
// leaving the cursor where it was. Zero picks 10s; negative disables idle sealing.
MaxIdle time.Duration
}
Options configures the on-disk spool. Everything not set here is derived at runtime.
type SealReason ¶
type SealReason int
SealReason says what closed a segment. It is counted rather than logged because the mix is the diagnosis: segments sealed at their size target are the intended path, and a run dominated by any other reason is committing short of what the sizer chose.
const ( // SealBySize is the intended path: the segment reached the size the sizer chose. SealBySize SealReason = iota // SealByIdle is the idle timer, which bounds what a stalled stream leaves unsealed. SealByIdle // SealByDrain is an undo, which has to see everything queued reach the database. SealByDrain // SealByClose is shutdown. SealByClose )
func (SealReason) String ¶
func (r SealReason) String() string
String is what the stats panel labels the count with.
type SegmentWriter ¶
type SegmentWriter interface {
// WriteRow appends one row of the named table.
WriteRow(table string, values []any) error
// PendingBytes is how much has been written so far, from the writers' own counters
// rather than a stat() per call. It is what the sizer compares against.
PendingBytes() int64
// Seal closes every stream and fills in the manifest's file records.
Seal(manifest *Manifest) error
// Discard closes and abandons a segment that will never be applied.
Discard()
}
SegmentWriter accumulates one segment's rows.
type Spool ¶
type Spool struct {
// contains filtered or unexported fields
}
Spool holds rows on disk and applies whole segments in the background.
Rows go in through Insert on the sinker's goroutine; sealed segments leave through a queue that one applier goroutine drains. The disk quota is what bounds how far ahead of the database the stream may run — there is deliberately no second ceiling on the number of queued segments, which would otherwise be the limit that binds while the operator watches the one they set.
func New ¶
func New(ctx context.Context, options Options, codec Codec, applier Applier, schema string, logger *zap.Logger) (*Spool, error)
New prepares the spool directory and starts the applier.
Recovery runs first: anything already on disk is either replayed or discarded before a single new row is written, so the applied state and the resume cursor cannot disagree.
func (*Spool) AppliedBlock ¶
AppliedBlock is the last block committed to the database, zero before the first one.
func (*Spool) BlocksBuffered ¶
BlocksBuffered is how many blocks are waiting, sealed or still being written.
func (*Spool) BytesOnDisk ¶
BytesOnDisk is how much is buffered waiting for the database, including the segment still being written. Counting only sealed segments would badly under-report: a segment grows to hundreds of megabytes before it is handed over.
func (*Spool) Close ¶
Close seals whatever is left, drains the applier and reports the first failure.
The drain is bounded by ctx. Sealing is what Close exists for — a segment left unsealed is discarded on the next start and its blocks streamed again — and that happens first, under the same deadline. Waiting for the queue to reach the database is a convenience on top: every segment in it is sealed and on disk, so recovery replays whatever the deadline cut short. Blocking a shutdown indefinitely on a database that has stopped responding buys nothing that the next start does not redo.
func (*Spool) Drain ¶
Drain seals what is being written and returns once every segment queued before it has been applied, so the caller can act on a database that holds everything the spool has accepted so far. An undo needs that: rows still in flight would otherwise land after the delete that was supposed to remove them.
func (*Spool) Insert ¶
Insert buffers one row. It blocks when the spool is full, which is the backpressure that keeps a slow database from filling the disk.
func (*Spool) MaybeSeal ¶
MaybeSeal hands the current segment to the applier once it is big enough. It is called at every sink flush, so it is also where backpressure is applied.
func (*Spool) RecordBlock ¶
RecordBlock notes which block the rows now being written belong to.
func (*Spool) RecordCursor ¶
RecordCursor notes the cursor covering everything written so far. The segment carries it so that applying the segment and advancing the cursor commit together.
A flush that produced no rows still has to advance the cursor, or a long stretch of blocks whose module output is empty is streamed, and paid for, again on restart. Such a flush opens a segment carrying nothing but the cursor.
func (*Spool) Seal ¶
func (b *Spool) Seal(ctx context.Context, reason SealReason) error
Seal closes the current segment and queues it, holding the stream if the disk quota is reached. The reason is only counted, so that the mix of size, idle, drain and shutdown seals is visible without a log line per segment.
func (*Spool) Stats ¶
Stats reads the counters and the open segment together.
The open segment is included in the queued totals for the same reason BytesOnDisk includes it: a segment grows to hundreds of megabytes before it is handed over, so counting only what is sealed badly under-reports what is waiting.
type Stats ¶
type Stats struct {
// Applied, cumulative: each counts only what a segment that reached the database
// carried, never what is still queued. Bytes is that segment's size on disk, in
// whatever format the codec writes — what the sizer steers and the disk quota bounds,
// not what crossed the network.
Segments int64
Blocks int64
Rows int64
Bytes int64
ApplyDuration time.Duration
// ApplierBusy and ApplierIdle split the applier goroutine's wall clock. Their ratio
// is the answer to "is the database or the stream the limit": an applier that is busy
// essentially all the time is the ceiling, one that mostly waits is not.
ApplierBusy time.Duration
ApplierIdle time.Duration
// QuotaWait is how long the stream was held because the spool was full. Any non-zero
// value means the database is gating download.
QuotaWait time.Duration
// Sealed counts segments by what closed them, indexed by SealReason.
Sealed [sealReasonCount]int64
// In flight right now.
QueueDepth int
SegmentTarget int64
OpenBytes int64
OpenBlocks int64
OpenAge time.Duration
// Queued totals, sealed segments plus the one being written.
BlocksBuffered int64
BytesOnDisk int64
AppliedBlock uint64
Quota int64
}
Stats is one consistent read of what the spool has committed and what it is still holding.
It is a snapshot rather than a set of accessors because the numbers only mean anything together: rows without the time they took is not a rate, and a gap without the applier's occupancy does not say whose fault it is.
The cumulative fields count what this process's applier goroutine has committed, so the caller derives intervals by differencing two snapshots. Segments replayed by recovery are not among them: recovery applies them before the goroutine starts, and counting another run's work as this one's throughput would misreport the rate for as long as the averages take to age out.
func (Stats) ApplierBusyRatioForTest ¶
ApplierBusyRatioForTest is the share of the applier's wall clock spent applying. It exists for the test that guards against the in-flight segment going uncounted; the panel differences two snapshots instead, which a single ratio cannot express.
type TableRecord ¶
type TableRecord struct {
Name string `json:"name"`
Schema string `json:"schema"`
Relation string `json:"relation"`
File string `json:"file"`
Columns []string `json:"columns"`
Rows int64 `json:"rows"`
Bytes int64 `json:"bytes"`
}
TableRecord describes one table's PGCOPY file inside a segment.
Schema and Name are the identifiers the server actually stored, not the logical table name the walk uses: the dialect writes DDL unquoted, so a message called BalanceChange lives in the relation balancechange, and a COPY has to say so.