kafka

package
Version: v0.3.0 Latest Latest
Warning

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

Go to latest
Published: Aug 5, 2020 License: MIT Imports: 17 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func NewDataSourceBuilderFactory

func NewDataSourceBuilderFactory(partitions int) physical.DataSourceBuilderFactory

func NewDataSourceBuilderFactoryFromConfig

func NewDataSourceBuilderFactoryFromConfig(dbConfig map[string]interface{}) (physical.DataSourceBuilderFactory, error)

NewDataSourceBuilderFactoryFromConfig creates a data source builder factory using the configuration.

Types

type DataSource

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

func (*DataSource) Get

type QueueElement

type QueueElement struct {
	// Types that are valid to be assigned to Type:
	//	*QueueElement_Record
	//	*QueueElement_Error
	Type                 isQueueElement_Type `protobuf_oneof:"type"`
	XXX_NoUnkeyedLiteral struct{}            `json:"-"`
	XXX_unrecognized     []byte              `json:"-"`
	XXX_sizecache        int32               `json:"-"`
}

func (*QueueElement) Descriptor

func (*QueueElement) Descriptor() ([]byte, []int)

func (*QueueElement) GetError

func (m *QueueElement) GetError() string

func (*QueueElement) GetRecord

func (m *QueueElement) GetRecord() *execution.Record

func (*QueueElement) GetType

func (m *QueueElement) GetType() isQueueElement_Type

func (*QueueElement) ProtoMessage

func (*QueueElement) ProtoMessage()

func (*QueueElement) Reset

func (m *QueueElement) Reset()

func (*QueueElement) String

func (m *QueueElement) String() string

func (*QueueElement) XXX_DiscardUnknown

func (m *QueueElement) XXX_DiscardUnknown()

func (*QueueElement) XXX_Marshal

func (m *QueueElement) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*QueueElement) XXX_Merge

func (m *QueueElement) XXX_Merge(src proto.Message)

func (*QueueElement) XXX_OneofWrappers

func (*QueueElement) XXX_OneofWrappers() []interface{}

XXX_OneofWrappers is for the internal use of the proto package.

func (*QueueElement) XXX_Size

func (m *QueueElement) XXX_Size() int

func (*QueueElement) XXX_Unmarshal

func (m *QueueElement) XXX_Unmarshal(b []byte) error

type QueueElement_Error

type QueueElement_Error struct {
	Error string `protobuf:"bytes,2,opt,name=error,proto3,oneof"`
}

type QueueElement_Record

type QueueElement_Record struct {
	Record *execution.Record `protobuf:"bytes,1,opt,name=record,proto3,oneof"`
}

type RecordStream

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

func (*RecordStream) Close

func (rs *RecordStream) Close(ctx context.Context, storage storage.Storage) error

func (*RecordStream) Next

func (rs *RecordStream) Next(ctx context.Context) (*execution.Record, error)

func (*RecordStream) RunWorker

func (rs *RecordStream) RunWorker(ctx context.Context) error

func (*RecordStream) RunWorkerInternal

func (rs *RecordStream) RunWorkerInternal(ctx context.Context, tx storage.StateTransaction) error

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
t or T : Toggle theme light dark auto
y or Y : Canonical URL