Documentation
¶
Index ¶
- func ResolveDecodeWorkers(workers int) int
- func SetupDatabaseSchema(ctx context.Context, dsnString string, schemaName string, ...) (protosql.Database, error)
- type Holder
- type Sinker
- func (s *Sinker) HandleBlockRangeCompletion(ctx context.Context, cursor *sink.Cursor) error
- func (s *Sinker) HandleBlockScopedData(ctx context.Context, data *pbsubstreamsrpc.BlockScopedData, isLive *bool, ...) (err error)
- func (s *Sinker) HandleBlockUndoSignal(ctx context.Context, undoSignal *pbsubstreamsrpc.BlockUndoSignal, ...) (err error)
- func (s *Sinker) LogStats()
- func (s *Sinker) Run(ctx context.Context) error
- type SinkerFactoryClickhouse
- type SinkerFactoryFunc
- type SinkerFactoryOptions
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func ResolveDecodeWorkers ¶ added in v1.22.0
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 Sinker ¶
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
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)
type SinkerFactoryClickhouse ¶
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 ¶
func (o SinkerFactoryOptions) Defaults() SinkerFactoryOptions
Directories
¶
| Path | Synopsis |
|---|---|
|
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. |