Documentation
¶
Index ¶
- func NewNMergedProducers(producers []query.RecordProducer, colIdx int, field string, ...) query.RecordProducer
- type DistributedAggregator
- type DocFilter
- type DocProjector
- type Eq
- type ExecutorState
- type FieldsFilter
- type Filter
- type FilterExpr
- type Gt
- type Limiter
- type Lt
- type Merger
- type SearcherDataSource
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func NewNMergedProducers ¶
func NewNMergedProducers( producers []query.RecordProducer, colIdx int, field string, dataType query.DataType, order seq.DocsOrder, ) query.RecordProducer
Types ¶
type DistributedAggregator ¶
type DistributedAggregator struct {
// contains filtered or unexported fields
}
func NewDistributedAggregator ¶
func NewDistributedAggregator( inputs []query.RecordProducer, aggFunc seq.AggFunc, quantiles []float64, ) *DistributedAggregator
func (*DistributedAggregator) Finalize ¶
func (a *DistributedAggregator) Finalize() *query.Summary
func (*DistributedAggregator) Next ¶
func (a *DistributedAggregator) Next() *query.Record
type DocFilter ¶
type DocFilter struct {
// contains filtered or unexported fields
}
func NewDocFilter ¶
func NewDocFilter( field string, filter FilterExpr[string], ) *DocFilter
type DocProjector ¶
type DocProjector struct {
// contains filtered or unexported fields
}
func NewDocProjector ¶
func NewDocProjector( input query.RecordProducer, colIdx int, filter *FieldsFilter, ) *DocProjector
func (*DocProjector) Finalize ¶
func (p *DocProjector) Finalize() *query.Summary
func (*DocProjector) Next ¶
func (p *DocProjector) Next() *query.Record
type Eq ¶
type Eq[T comparable] struct { // contains filtered or unexported fields }
func NewEq ¶
func NewEq[T comparable]( pred T, ) *Eq[T]
type ExecutorState ¶
type ExecutorState byte
const ( ExecutorStateReadingInput ExecutorState = iota ExecutorStateProcessingData ExecutorStateProducingOutput ExecutorStateDone )
type FieldsFilter ¶
type Filter ¶
type Filter[T any] struct { // contains filtered or unexported fields }
func NewFilter ¶
func NewFilter[T any]( input query.RecordProducer, colIdx int, expr FilterExpr[T], withTotal bool, ) *Filter[T]
type FilterExpr ¶
type Limiter ¶
type Limiter struct {
// contains filtered or unexported fields
}
func NewLimiter ¶
func NewLimiter( input query.RecordProducer, limit uint32, offset uint32, ) *Limiter
type SearcherDataSource ¶
type SearcherDataSource struct {
// contains filtered or unexported fields
}
SearcherDataSource is limitless: for the documents path it walks the matched set in fixed-size batches via cursor pagination:
- repeatedly call Searcher.SearchDocs with a constant per-batch limit and a cursor (OffsetId) advanced from the previous batch's last ID;
- the Fetcher.FetchDocs step is performed on demand — only the current document is fetched right before it is returned from Next().
The aggregation path is a scan-all request (one SearchDocs, no cursor): aggregations require a full scan and the proxy-side DistributedAggregator merges complete results.
func NewSearcherDataSource ¶
func NewSearcherDataSource( ctx context.Context, tr *querytracer.Tracer, searchParams processor.SearchParams, fracManager *fracmanager.FracManager, searcher *fracmanager.Searcher, fetcher *fracmanager.Fetcher, ) *SearcherDataSource
func (*SearcherDataSource) Ctx ¶
func (s *SearcherDataSource) Ctx() context.Context
func (*SearcherDataSource) Finalize ¶
func (s *SearcherDataSource) Finalize() *query.Summary
func (*SearcherDataSource) Next ¶
func (s *SearcherDataSource) Next() *query.Record
Click to show internal directories.
Click to hide internal directories.