sql

package
v1.25.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 8, 2026 License: Apache-2.0 Imports: 16 Imported by: 0

Documentation

Index

Constants

View Source
const AllTables = "all"

AllTables is what a per-table constraint switch takes to mean every table at once.

View Source
const DialectFieldBlockNumber = "_block_number_"
View Source
const DialectFieldBlockTimestamp = "_block_timestamp_"
View Source
const DialectFieldDeleted = "_deleted_"
View Source
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.

View Source
const DialectFieldVersion = "_version_"
View Source
const DialectTableBlock = "_blocks_"
View Source
const DialectTableCursor = "_cursors_"

Variables

This section is empty.

Functions

func Quoted

func Quoted(value string) string

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 NewBaseDialect(registry map[string]*schema.Table, logger *zap.Logger) *BaseDialect

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

func (d *BaseDialect) GetTable(table string) *schema.Table

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

func (c *Context) BlockNumber() int

func (*Context) SetNumber

func (c *Context) SetNumber(id int)

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

type EnumValue struct {
	Number int32
	Name   string
}

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.

func (EnumValue) String added in v1.22.0

func (v EnumValue) String() string

String is the textual rendering: the enum's name, or its number when the value is not one the descriptor knows about, which is what a proto3 reader must tolerate.

type ForeignKey

type ForeignKey struct {
	Name         string
	Table        string
	Field        string
	ForeignTable string
	ForeignField string
}

func (*ForeignKey) String

func (f *ForeignKey) String() string

type Inserter

type Inserter interface {
	Insert(table string, values []any) error
}

type Parent

type Parent struct {
	// contains filtered or unexported fields
}

type SinkInfo

type SinkInfo struct {
	SchemaHash string `json:"schema_hash"`
}

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

func ParseWriteMode(in string) (WriteMode, error)

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.

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.

Jump to

Keyboard shortcuts

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