Documentation
¶
Index ¶
- Constants
- func Quoted(value string) string
- func ScalarFieldValue(fd protoreflect.FieldDescriptor, value protoreflect.Value) any
- type BaseDatabase
- type BaseDialect
- func (d *BaseDialect) AddCreateTableSql(table string, sql string)
- func (d *BaseDialect) AddForeignKeyReferencing(table string, referencedTable string, sql string)
- func (d *BaseDialect) AddForeignKeySql(table string, sql string)
- func (d *BaseDialect) AddIndexSql(table string, sql string)
- func (d *BaseDialect) AddPrimaryKeySql(table string, sql string)
- func (d *BaseDialect) AddUniqueConstraintSql(table string, sql string)
- func (d *BaseDialect) GetCreateTableSql(table string) string
- func (d *BaseDialect) GetTable(table string) *schema.Table
- func (d *BaseDialect) GetTables() []*schema.Table
- func (d *BaseDialect) TableApplyOrder() ([]string, error)
- func (d *BaseDialect) TableApplyRanks() (map[string]int, error)
- func (d *BaseDialect) TableNames() []string
- func (d *BaseDialect) UseRowIDField(table string) bool
- type BufferedInserter
- type Constraint
- type ConstraintPolicy
- func (p ConstraintPolicy) ApplyAtHead() bool
- func (p ConstraintPolicy) ApplyUpfront() bool
- func (p ConstraintPolicy) ConstraintsParallelism() int
- func (p ConstraintPolicy) Describe() string
- func (p ConstraintPolicy) SkipForeignKey(string) bool
- func (p ConstraintPolicy) SkipPrimaryKey(table string) bool
- func (p ConstraintPolicy) SkipUnique(table string) bool
- func (p ConstraintPolicy) SkipsEverything() bool
- func (p ConstraintPolicy) WithBlockNumberIndex(from ConstraintPolicy) ConstraintPolicy
- type ConstraintTiming
- type Context
- type Database
- type Dialect
- type EnumValue
- type ForeignKey
- type Inserter
- type Parent
- type SinkInfo
- type WriteMode
- type WriteStats
Constants ¶
const AllTables = "all"
AllTables is what a per-table constraint switch takes to mean every table at once.
const DialectFieldBlockNumber = "_block_number_"
const DialectFieldBlockTimestamp = "_block_timestamp_"
const DialectFieldDeleted = "_deleted_"
const DialectFieldRowID = "_row_id_"
DialectFieldRowID numbers the rows a single block writes to a single table, starting at zero. It only exists where the sorting key would otherwise not be unique, which today means a ClickHouse table whose message carries no 'order_by_fields' annotation: the sink then sorts on (_block_number_, _row_id_), and without the second column the ReplacingMergeTree would collapse every row of a block into one.
const DialectFieldVersion = "_version_"
const DialectTableBlock = "_blocks_"
const DialectTableCursor = "_cursors_"
Variables ¶
This section is empty.
Functions ¶
func ScalarFieldValue ¶ added in v1.22.0
func ScalarFieldValue(fd protoreflect.FieldDescriptor, value protoreflect.Value) any
ScalarFieldValue unwraps a non-message, non-list value for the inserters.
Enums are the only kind that cannot be handed over as-is: protoreflect yields a protoreflect.EnumNumber, a named int32 type that every dialect's type switch misses, so both of them panicked on any message carrying an enum field.
Every path that feeds a value to an inserter has to go through here, including the list elements and the dialects' inline/nested extraction, or the same panic comes back for a repeated or nested enum.
Types ¶
type BaseDatabase ¶
type BaseDatabase struct {
RootMessageDescriptor protoreflect.MessageDescriptor
// contains filtered or unexported fields
}
func NewBaseDatabase ¶
func NewBaseDatabase(moduleOutputType string, rootMessageDescriptor protoreflect.MessageDescriptor, useProtoOptions bool, logger *zap.Logger) (database *BaseDatabase, err error)
func (*BaseDatabase) WalkMessageDescriptorAndInsertWithDialect ¶
func (d *BaseDatabase) WalkMessageDescriptorAndInsertWithDialect(dm protoreflect.Message, blockNum uint64, blockTimestamp time.Time, parent *Parent, dialect Dialect, inserter Inserter) (time.Duration, error)
WalkMessageDescriptorAndInsertWithDialect turns one message into rows.
It takes a protoreflect.Message rather than a concrete *dynamicpb.Message because the walk is read-only: it calls Descriptor, Get and IsValid and nothing else. That is exactly the subset hyperpb implements, which is what lets the decoder swap in the faster parser without this file changing behaviour.
Values reached through Get may alias the parser's own memory — hyperpb strings and bytes point into its arena — so an inserter that keeps a []any past the walk keeps the message alive too. See decoder.arenas for where that lifetime is managed.
type BaseDialect ¶
type BaseDialect struct {
CreateTableSql map[string]string
PrimaryKeySql []*Constraint
ForeignKeySql []*Constraint
UniqueConstraintSql []*Constraint
// IndexSql holds the indexes the sink creates for itself rather than because the
// schema asked for one. They are built in the same pass as the constraints, being the
// same kind of expensive.
IndexSql []*Constraint
TableRegistry map[string]*schema.Table
Logger *zap.Logger
}
func NewBaseDialect ¶
func (*BaseDialect) AddCreateTableSql ¶
func (d *BaseDialect) AddCreateTableSql(table string, sql string)
func (*BaseDialect) AddForeignKeyReferencing ¶ added in v1.22.0
func (d *BaseDialect) AddForeignKeyReferencing(table string, referencedTable string, sql string)
AddForeignKeyReferencing records the same statement along with the logical table it points at, which is what TableApplyOrder needs.
func (*BaseDialect) AddForeignKeySql ¶
func (d *BaseDialect) AddForeignKeySql(table string, sql string)
func (*BaseDialect) AddIndexSql ¶ added in v1.22.0
func (d *BaseDialect) AddIndexSql(table string, sql string)
func (*BaseDialect) AddPrimaryKeySql ¶
func (d *BaseDialect) AddPrimaryKeySql(table string, sql string)
func (*BaseDialect) AddUniqueConstraintSql ¶
func (d *BaseDialect) AddUniqueConstraintSql(table string, sql string)
func (*BaseDialect) GetCreateTableSql ¶
func (d *BaseDialect) GetCreateTableSql(table string) string
func (*BaseDialect) GetTables ¶
func (d *BaseDialect) GetTables() []*schema.Table
func (*BaseDialect) TableApplyOrder ¶ added in v1.22.0
func (d *BaseDialect) TableApplyOrder() ([]string, error)
TableApplyOrder returns every table ordered so that a referenced table always comes before the tables referencing it.
Rows have to reach the server in that order whenever foreign keys are enforced, and the write paths group rows by table: a multi-row INSERT per table at flush, and one binary COPY per table in the buffer. Insertion order within a block is not enough, because the grouping loses it.
Nesting depth is not the answer even though it looks like it: a table can point at a sibling it has no ancestry with, so this is a topological sort over the foreign keys themselves.
A cycle has no valid order and is reported as an error rather than silently ordered wrong. A table referencing itself is not a cycle for these purposes — it constrains the order of rows within one table, which no table-level ordering can address — so those edges are skipped.
func (*BaseDialect) TableApplyRanks ¶ added in v1.22.0
func (d *BaseDialect) TableApplyRanks() (map[string]int, error)
TableApplyRanks is TableApplyOrder as a lookup, for sorting a set of tables that is not the whole schema. A table the dialect does not know about ranks last.
func (*BaseDialect) TableNames ¶ added in v1.22.0
func (d *BaseDialect) TableNames() []string
TableNames lists every table of the schema in a stable order: the block table first, being referenced by every other one and referencing none, then the rest sorted.
It is what TableApplyOrder starts from, and what the reorg path falls back to when no foreign key order exists.
func (*BaseDialect) UseRowIDField ¶ added in v1.22.0
func (d *BaseDialect) UseRowIDField(table string) bool
UseRowIDField defaults to false: only the ClickHouse dialect needs a tie-breaker in its sorting key.
type BufferedInserter ¶ added in v1.22.0
type BufferedInserter struct {
// contains filtered or unexported fields
}
BufferedInserter records the inserts a message walk produces instead of performing them, so the walk can run off the goroutine that owns the database.
It is not safe for concurrent use: each worker owns one. Replay puts the recorded inserts back in the exact order the walk produced them, which is what keeps a parent row ahead of its children.
func NewBufferedInserter ¶ added in v1.22.0
func NewBufferedInserter(expectedRows int) *BufferedInserter
NewBufferedInserter returns a buffer sized for roughly the given number of rows.
func (*BufferedInserter) Insert ¶ added in v1.22.0
func (b *BufferedInserter) Insert(table string, values []any) error
func (*BufferedInserter) Replay ¶ added in v1.22.0
func (b *BufferedInserter) Replay(target Inserter) error
Replay performs the buffered inserts against target, in order.
func (*BufferedInserter) Reset ¶ added in v1.22.0
func (b *BufferedInserter) Reset()
Reset drops the buffered inserts, keeping the backing array for reuse.
func (*BufferedInserter) Rows ¶ added in v1.22.0
func (b *BufferedInserter) Rows() int
Rows returns how many inserts are buffered.
type Constraint ¶
type Constraint struct {
Table string
Sql string
// ReferencedTable is the logical name of the table a foreign key points at, empty for
// any other kind of constraint. It is what the apply order is computed from: the
// SQL carries schema-qualified names, which is not what the registry is keyed by.
ReferencedTable string
}
type ConstraintPolicy ¶ added in v1.22.0
type ConstraintPolicy struct {
// Timing decides when they are created; the zero value is ConstraintsAuto.
Timing ConstraintTiming
// DisableForeignKeys leaves out every foreign key, including the one to the block
// table. Those are what a load pays most for.
DisableForeignKeys bool
// DisablePrimaryKeys and DisableUniques name the tables that go without, or AllTables.
DisablePrimaryKeys []string
DisableUniques []string
// DisableBlockNumberIndex leaves out the index on _block_number_. Every table carries
// that column and the reorg path deletes by it on every table, so without the index
// each undo is a sequential scan per table. It is only dead weight on a run that can
// never reorg.
//
// Unlike the rest of this policy it is not about what the schema declares, and not
// governed by Timing: the index is created when the sink starts.
DisableBlockNumberIndex bool
// Parallelism is how many constraints are created or dropped at once. Zero means one.
//
// They go on independent relations, so the only ordering that matters is between the
// waves — a foreign key needs the key it references to exist — and within a wave the
// server is free to build them side by side. Each statement still commits on its own,
// which is what keeps the pass restartable: a run that is killed keeps what it
// finished. It is an execution knob rather than a property of the schema.
Parallelism int
// WorkMem is what maintenance_work_mem is set to for the duration of each statement,
// empty leaving the server's own setting alone. The default is 64MB on most servers,
// at which an index build over a large table spills to an external merge sort; raising
// it for the pass alone is the cheapest thing that makes it faster.
//
// It multiplies with Parallelism, each concurrent build taking its own.
WorkMem string
}
func DisableAllConstraints ¶ added in v1.22.0
func DisableAllConstraints() ConstraintPolicy
DisableAllConstraints returns the policy that declares none of them, which is what an output with no schema annotations leaves the sink with. The block number index survives it: nothing in the annotations asks for that one, the reorg path does.
func (ConstraintPolicy) ApplyAtHead ¶ added in v1.22.0
func (p ConstraintPolicy) ApplyAtHead() bool
ApplyAtHead reports whether the sink creates them itself once the backfill is over.
func (ConstraintPolicy) ApplyUpfront ¶ added in v1.22.0
func (p ConstraintPolicy) ApplyUpfront() bool
ApplyUpfront reports whether the constraints go in before the first row.
func (ConstraintPolicy) ConstraintsParallelism ¶ added in v1.22.0
func (p ConstraintPolicy) ConstraintsParallelism() int
ConstraintsParallelism is how many statements the pass runs at once.
func (ConstraintPolicy) Describe ¶ added in v1.22.0
func (p ConstraintPolicy) Describe() string
Describe renders the policy for a log line.
func (ConstraintPolicy) SkipForeignKey ¶ added in v1.22.0
func (p ConstraintPolicy) SkipForeignKey(string) bool
SkipForeignKey reports whether the given table's foreign keys are meant to be left out.
func (ConstraintPolicy) SkipPrimaryKey ¶ added in v1.22.0
func (p ConstraintPolicy) SkipPrimaryKey(table string) bool
SkipPrimaryKey reports whether the given table is meant to go without its primary key.
func (ConstraintPolicy) SkipUnique ¶ added in v1.22.0
func (p ConstraintPolicy) SkipUnique(table string) bool
SkipUnique reports whether the given table's unique constraints are meant to be left out.
func (ConstraintPolicy) SkipsEverything ¶ added in v1.22.0
func (p ConstraintPolicy) SkipsEverything() bool
SkipsEverything reports a policy that declares no constraints, in which case there is nothing to apply and nothing to wait for. The block number index is not one of them: it is created when the sink starts, whatever the constraints say.
func (ConstraintPolicy) WithBlockNumberIndex ¶ added in v1.22.0
func (p ConstraintPolicy) WithBlockNumberIndex(from ConstraintPolicy) ConstraintPolicy
WithBlockNumberIndex carries the index switch over from another policy, the index being the one thing that survives an output having no annotations to declare anything.
type ConstraintTiming ¶ added in v1.22.0
type ConstraintTiming string
ConstraintTiming says when the constraints are created.
const ( // ConstraintsAuto has the sink create them itself at the first live block, and only // there: a stop block ends a run without saying the backfill is over, and a range is // routinely one chunk of several. It is the default, stop-the-world pass and all. A backfill that ends with no primary keys and no foreign keys has // produced a database nobody should query, and leaving it that way until the operator // remembers a second command is the worse failure — it is silent, and it looks like // success. ConstraintsAuto ConstraintTiming = "auto" // ConstraintsManual leaves it to the operator, through `sink postgres constraints // apply`. Building them locks every table while indexes are built and every foreign // key is validated, so on a large database this is how the pass goes into a // maintenance window instead. ConstraintsManual ConstraintTiming = "manual" // ConstraintsAlways creates them before the first row is written, so the database // rejects bad data from the start and the load pays for it throughout. ConstraintsAlways ConstraintTiming = "always" )
func ParseConstraintTiming ¶ added in v1.22.0
func ParseConstraintTiming(in string) (ConstraintTiming, error)
ParseConstraintTiming validates the flag value.
type Context ¶
type Context struct {
// contains filtered or unexported fields
}
func NewContext ¶
func NewContext() *Context
func (*Context) BlockNumber ¶
type Database ¶
type Database interface {
// Inserter: both implementations already pass themselves as the inserter for a
// walk, this states it.
Inserter
FetchSinkInfo(schemaName string) (*SinkInfo, error)
UpdateSinkInfoHash(schemaName string, newHash string) error
StoreSinkInfo(schemaName string, schemaHash string) error
CreateDatabase(useConstraints bool) error
// ApplyConstraints adds the schema's constraints to a database that already exists,
// skipping the ones already in place. A schema first synced without constraints has
// none of them, and only this puts them there.
ApplyConstraints() error
// EnsureBlockNumberIndexes creates the index the sink needs for its own reorg path,
// when it starts. It is not part of the constraint pass: --apply-constraints describes
// the schema and is the operator's to schedule, where this one the sink depends on to
// undo a reorg without sequentially scanning every table.
EnsureBlockNumberIndexes(ctx context.Context) error
// MissingConstraints names the constraints the policy says the schema should carry and
// the database does not have. It is what turns "this schema has no indexes" from
// something you find out by querying it into something the sink says on every start.
MissingConstraints() ([]string, error)
// DropConstraints removes the constraints this schema's DDL would create, leaving
// anything the sink did not put there alone. It is the escape hatch after
// --apply-constraints=always, and what makes a stalled backfill fast again without a
// second setup.
DropConstraints() error
// SwitchToDirectInserts leaves any buffered write path behind and inserts straight
// into the database from now on. It is a one-way switch, and a no-op for a backend
// that never buffered. The reason is what the caller knows and the database does not
// — reaching the chain head, or reaching the end of a bounded range — and it is what
// the switch reports, so the log does not claim a head the stream never saw.
// atChainHead says the sink carries on from here rather than exiting, which is what
// makes the spool's own bookkeeping worth clearing: a run that stays live seals a
// segment or two on every restart and would otherwise accumulate their records for as
// long as it lives.
SwitchToDirectInserts(ctx context.Context, reason string, atChainHead bool) error
// WalkMessageDescriptorAndInsert reads the message through protoreflect only, so it
// does not care which dynamic implementation produced it — dynamicpb and hyperpb are
// both accepted, and the decoder picks.
WalkMessageDescriptorAndInsert(dm protoreflect.Message, blockNum uint64, blockTimestamp time.Time, parent *Parent) (time.Duration, error)
// WalkMessageDescriptorAndInsertInto is the same walk, against a caller-supplied
// inserter. It touches no database state, so it is safe to call concurrently with
// a BufferedInserter per goroutine.
WalkMessageDescriptorAndInsertInto(dm protoreflect.Message, blockNum uint64, blockTimestamp time.Time, parent *Parent, inserter Inserter) (time.Duration, error)
InsertBlock(blockNum uint64, hash string, timestamp time.Time) error
HandleBlocksUndo(lastValidBlockNumber uint64) error
FetchCursor() (*sink.Cursor, error)
StoreCursor(cursor *sink.Cursor) error
BeginTransaction() error
CommitTransaction() error
RollbackTransaction()
Flush() (time.Duration, error)
DatabaseHash(schemaName string) (uint64, error)
GetDialect() Dialect
Open() error
// Close releases anything the database buffered locally, so blocks held at shutdown
// reach the server rather than being streamed again.
Close(ctx context.Context) error
// BufferStats reports what sits between the stream and the server, and what the
// applier has done to drain it. enabled is false when nothing buffers locally, in
// which case the caller knows a flush means the rows are stored and the snapshot
// says nothing.
//
// The whole snapshot rather than the few numbers the gap needs, because they are only
// readable together: a gap says the database is behind, and only the applier's
// occupancy next to it says whether that is the database being slow or the stream
// having burst.
BufferStats() (stats WriteStats, enabled bool)
}
type Dialect ¶
type Dialect interface {
SchemaHash() string
FullTableName(table *schema.Table) string
GetTable(table string) *schema.Table
GetTables() []*schema.Table
UseVersionField() bool
UseDeletedField() bool
// UseRowIDField reports whether the given table carries the DialectFieldRowID column.
// It is per-table because a schema can annotate some of its messages and not others.
UseRowIDField(table string) bool
AppendInlineFieldValues(fieldValues []any, fd protoreflect.FieldDescriptor, fv protoreflect.Value, dm protoreflect.Message) ([]any, error)
}
type EnumValue ¶ added in v1.22.0
EnumValue carries a protobuf enum field through the walk, keeping both renderings.
The dialects disagree on how an enum should be stored — PostgreSQL declares the column TEXT, ClickHouse declares it Int32 — and the walk is shared, so it cannot pick one. Passing both lets each dialect take what matches the column it created, and neither has to reach back to the field descriptor to resolve a name.
type ForeignKey ¶
type ForeignKey struct {
Name string
Table string
Field string
ForeignTable string
ForeignField string
}
func (*ForeignKey) String ¶
func (f *ForeignKey) String() string
type WriteMode ¶ added in v1.22.0
type WriteMode string
WriteMode says how a sealed spool segment reaches the database.
It is deliberately the operator's choice rather than something derived from a directory flag and the shape of the schema's foreign keys, which is how the sink used to decide. A mode that the driver or the schema cannot support is an error, not a silent downgrade to something an order of magnitude slower.
const ( // WriteModeAuto picks copy on PostgreSQL, batch-insert on ClickHouse, and row-insert // for a schema whose foreign keys cannot be ordered. WriteModeAuto WriteMode = "auto" // WriteModeCopy loads each table file with `COPY ... FROM STDIN (FORMAT BINARY)`. // Measured at ~7x the multi-row INSERT path, see sink/sql/db_proto/benchmarks. WriteModeCopy WriteMode = "copy" // WriteModeBatchInsert builds one multi-row INSERT per table, split to stay under the // driver's bind-parameter and statement-size limits. WriteModeBatchInsert WriteMode = "batch-insert" // WriteModeRowInsert issues one prepared INSERT per row, in walk order. It is the // fallback for a schema whose foreign keys form a cycle: a cycle has no table order, // so grouping rows by table cannot keep a parent ahead of its children, while the walk // itself always does. WriteModeRowInsert WriteMode = "row-insert" )
func ParseWriteMode ¶ added in v1.22.0
type WriteStats ¶ added in v1.22.0
type WriteStats struct {
spool.Stats
// DatabaseBytesWritten is what actually went down the socket to the server.
DatabaseBytesWritten int64
}
WriteStats is the spool's own account of what it committed, plus what only the driver can answer.
The split is the point: the spool measures segments on disk, in the format its codec writes, which is what the sizer steers and the quota bounds. How many bytes that turned into on the wire depends on the driver's encoding and on whether the connection compresses, so it cannot come from here — and a driver that does not measure it leaves it zero rather than reporting a wrong number.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
pgcopy
Package pgcopy writes the PostgreSQL binary COPY format ("PGCOPY").
|
Package pgcopy writes the PostgreSQL binary COPY format ("PGCOPY"). |
|
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.
|
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. |