db_proto

package
v1.23.0 Latest Latest
Warning

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

Go to latest
Published: Sep 22, 2026 License: Apache-2.0 Imports: 20 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func ResolveDecodeWorkers added in v1.22.0

func ResolveDecodeWorkers(workers int) int

ResolveDecodeWorkers turns a requested worker count into the one that will be used.

One per core, less one for the goroutine draining the gRPC stream, capped at 8: TestClientDecodeScaling measures 4.52x at eight workers and 5.06x at fifteen, so the seven extra cores together buy 11%. Taking them would cost the rest of the machine for almost nothing.

func SetupDatabaseSchema

func SetupDatabaseSchema(
	ctx context.Context,
	dsnString string,
	schemaName string,
	outputModuleName string,
	rootMessageDescriptor protoreflect.MessageDescriptor,
	options SinkerFactoryOptions,
	logger *zap.Logger,
	tracer logging.Tracer,
) (protosql.Database, error)

SetupDatabaseSchema constructs the target Database for the given DSN and output message descriptor and ensures its schema exists: on a fresh database it creates the schema, tables and system tables and records the sink info hash; on an existing one it reconciles the sink info hash when a compatible migration is available. It returns the ready (but not opened) Database so callers can either start a sinker (run path) or stop right after schema setup (setup command). It is idempotent: on a database that is already set up with a matching hash, it performs no schema changes.

Types

type Holder

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

type Sinker

type Sinker struct {
	*sink.Sinker
	// contains filtered or unexported fields
}

func NewSinker

func NewSinker(rootMessageDescriptor protoreflect.MessageDescriptor, sink *sink.Sinker, db sql.Database, useTransaction bool, constraints sql.ConstraintPolicy, blockBatchSize int, decodeWorkers int, stats *stats.Stats, logger *zap.Logger) *Sinker

NewSinker builds the from-proto sinker. decodeWorkers bounds how many blocks are unmarshalled and walked concurrently at flush time; zero picks one per core, less one, capped at eight.

func (*Sinker) HandleBlockRangeCompletion added in v1.22.0

func (s *Sinker) HandleBlockRangeCompletion(ctx context.Context, cursor *sink.Cursor) error

HandleBlockRangeCompletion flushes whatever is still held when a bounded run reaches its stop block.

Blocks accumulate until the batch is full, so a run that ends mid-batch left its last blocks in memory and never wrote them, nor the cursor covering them. With --block-batch-size larger than the range, that meant an empty database and a run that looked successful.

func (*Sinker) HandleBlockScopedData

func (s *Sinker) HandleBlockScopedData(ctx context.Context, data *pbsubstreamsrpc.BlockScopedData, isLive *bool, cursor *sink.Cursor) (err error)

func (*Sinker) HandleBlockUndoSignal

func (s *Sinker) HandleBlockUndoSignal(ctx context.Context, undoSignal *pbsubstreamsrpc.BlockUndoSignal, cursor *sink.Cursor) (err error)

func (*Sinker) LogStats

func (s *Sinker) LogStats()

func (*Sinker) Run

func (s *Sinker) Run(ctx context.Context) error

type SinkerFactoryClickhouse

type SinkerFactoryClickhouse struct {
	SinkInfoFolder  string
	CursorFilePath  string
	QueryRetryCount int
	QueryRetrySleep time.Duration
}

type SinkerFactoryFunc

type SinkerFactoryFunc func(ctx context.Context, dsnString, schemaName string, logger *zap.Logger, tracer logging.Tracer) (*Sinker, error)

func SinkerFactory

func SinkerFactory(
	baseSink *sink.Sinker,
	outputModuleName string,
	rootMessageDescriptor protoreflect.MessageDescriptor,
	options SinkerFactoryOptions,
) SinkerFactoryFunc

type SinkerFactoryOptions

type SinkerFactoryOptions struct {
	UseProtoOption bool
	// Constraints says which constraints the schema gets and when; the zero value means
	// all of them, created once the backfill is done.
	Constraints     protosql.ConstraintPolicy
	UseTransactions bool
	// WriteMode says how a sealed spool segment reaches the database. The zero value
	// resolves per driver and schema.
	WriteMode protosql.WriteMode
	// DecodeWorkers bounds how many blocks are unmarshalled and walked concurrently at
	// flush time. Zero picks one per core, less one, capped at 8.
	DecodeWorkers int
	// DecodeBatchSize is how many blocks are held in memory and decoded together. Zero
	// picks four per decode worker. It sizes the CPU stage, not the database write.
	DecodeBatchSize int
	Encoding        bytes.Encoding
	// Spool, when set, holds rows on disk and applies them from a background goroutine.
	Spool      *spool.Options
	Clickhouse SinkerFactoryClickhouse
}

func (SinkerFactoryOptions) Defaults

Directories

Path Synopsis
sql
postgres/pgcopy
Package pgcopy writes the PostgreSQL binary COPY format ("PGCOPY").
Package pgcopy writes the PostgreSQL binary COPY format ("PGCOPY").
spool
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