connpubsub

package
v0.0.0-...-205417f Latest Latest
Warning

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

Go to latest
Published: Feb 28, 2026 License: AGPL-3.0 Imports: 19 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type PubSubConnector

type PubSubConnector struct {
	*metadataStore.PostgresMetadata
	// contains filtered or unexported fields
}

func NewPubSubConnector

func NewPubSubConnector(
	ctx context.Context,
	env map[string]string,
	config *protos.PubSubConfig,
) (*PubSubConnector, error)

func (*PubSubConnector) Close

func (c *PubSubConnector) Close() error

func (*PubSubConnector) ConnectionActive

func (c *PubSubConnector) ConnectionActive(ctx context.Context) error

func (*PubSubConnector) CreateRawTable

func (*PubSubConnector) NewTopicCache

func (c *PubSubConnector) NewTopicCache() topicCache

func (*PubSubConnector) ReplayTableSchemaDeltas

func (c *PubSubConnector) ReplayTableSchemaDeltas(_ context.Context, _ map[string]string,
	flowJobName string, _ []*protos.TableMapping, schemaDeltas []*protos.TableSchemaDelta, _ []string,
) error

func (*PubSubConnector) SetupQRepMetadataTables

func (*PubSubConnector) SetupQRepMetadataTables(_ context.Context, _ *protos.QRepConfig) error

func (*PubSubConnector) SyncQRepRecords

func (c *PubSubConnector) SyncQRepRecords(
	ctx context.Context,
	config *protos.QRepConfig,
	partition *protos.QRepPartition,
	stream *model.QRecordStream,
) (int64, shared.QRepWarnings, error)

func (*PubSubConnector) SyncRecords

type PubSubMessage

type PubSubMessage struct {
	*pubsub.Message
	Topic string
}

Jump to

Keyboard shortcuts

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