Versions in this module Expand all Collapse all v0 v0.3.0 Aug 5, 2020 Changes in this version + func NewDataSourceBuilderFactory(partitions int) physical.DataSourceBuilderFactory + func NewDataSourceBuilderFactoryFromConfig(dbConfig map[string]interface{}) (physical.DataSourceBuilderFactory, error) + type DataSource struct + func (ds *DataSource) Get(ctx context.Context, variables octosql.Variables, streamID *execution.StreamID) (execution.RecordStream, *execution.ExecutionOutput, error) + type QueueElement struct + Type isQueueElement_Type + XXX_NoUnkeyedLiteral struct{} + XXX_sizecache int32 + XXX_unrecognized []byte + func (*QueueElement) Descriptor() ([]byte, []int) + func (*QueueElement) ProtoMessage() + func (*QueueElement) XXX_OneofWrappers() []interface{} + func (m *QueueElement) GetError() string + func (m *QueueElement) GetRecord() *execution.Record + func (m *QueueElement) GetType() isQueueElement_Type + func (m *QueueElement) Reset() + func (m *QueueElement) String() string + func (m *QueueElement) XXX_DiscardUnknown() + func (m *QueueElement) XXX_Marshal(b []byte, deterministic bool) ([]byte, error) + func (m *QueueElement) XXX_Merge(src proto.Message) + func (m *QueueElement) XXX_Size() int + func (m *QueueElement) XXX_Unmarshal(b []byte) error + type QueueElement_Error struct + Error string + type QueueElement_Record struct + Record *execution.Record + type RecordStream struct + func (rs *RecordStream) Close(ctx context.Context, storage storage.Storage) error + func (rs *RecordStream) Next(ctx context.Context) (*execution.Record, error) + func (rs *RecordStream) RunWorker(ctx context.Context) error + func (rs *RecordStream) RunWorkerInternal(ctx context.Context, tx storage.StateTransaction) error