Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion proxy/search/stream_search.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ func (si *Ingestor) StreamSearch(
mergedStream = exec.NewDistributedAggregator(producers, sr.Agg.Func, sr.Agg.Quantiles)
} else {
const seqIdColIdx = 0
mergedDocsStream := exec.NewNMergedProducers(producers, seqIdColIdx, "", query.DataTypeSeqID, sr.Order)
mergedDocsStream := exec.NewNMergedProducers(producers, query.SeqIDColumn(seqIdColIdx), sr.Order)
mergedStream = exec.NewLimiter(mergedDocsStream, uint32(sr.Size), uint32(offset))
}

Expand Down
119 changes: 119 additions & 0 deletions query/column.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
package query

import (
insaneJSON "github.com/ozontech/insane-json"

"github.com/ozontech/seq-db/seq"
)

type Column[T any] struct {
idx int
dataType DataType
get func(*RecordVals) T
}

func (c Column[T]) Idx() int {
return c.idx
}

func (c Column[T]) DataType() DataType {
return c.dataType
}

// Val returns value of r's column at c.idx index
func (c Column[T]) Val(r *Record) T {
return c.get(r.Vals[c.idx])
}

// RawData returns raw data of r's column at c.idx index
func (c Column[T]) RawData(r *Record) []byte {
return r.Vals[c.idx].RawData()
}

func SeqIDColumn(idx int) Column[seq.ID] {
return Column[seq.ID]{
idx: idx,
dataType: DataTypeSeqID,
get: (*RecordVals).AsSeqID,
}
}

func BytesColumn(idx int) Column[[]byte] {
return Column[[]byte]{
idx: idx,
dataType: DataTypeBytes,
get: (*RecordVals).AsBytes,
}
}

func StringColumn(idx int) Column[string] {
return Column[string]{
idx: idx,
dataType: DataTypeString,
get: (*RecordVals).AsString,
}
}

func DocColumn(idx int) Column[*insaneJSON.Root] {
return Column[*insaneJSON.Root]{
idx: idx,
dataType: DataTypeDocument,
get: (*RecordVals).AsDoc,
}
}

func Float64Column(idx int) Column[float64] {
return Column[float64]{
idx: idx,
dataType: DataTypeFloat64,
get: (*RecordVals).AsFloat64,
}
}

func Uint64Column(idx int) Column[uint64] {
return Column[uint64]{
idx: idx,
dataType: DataTypeUint64,
get: (*RecordVals).AsUint64,
}
}

func Int64Column(idx int) Column[int64] {
return Column[int64]{
idx: idx,
dataType: DataTypeInt64,
get: (*RecordVals).AsInt64,
}
}

func Uint32Column(idx int) Column[uint32] {
return Column[uint32]{
idx: idx,
dataType: DataTypeUint32,
get: (*RecordVals).AsUint32,
}
}

func Int32Column(idx int) Column[int32] {
return Column[int32]{
idx: idx,
dataType: DataTypeInt32,
get: (*RecordVals).AsInt32,
}
}

func Float64ArrayColumn(idx int) Column[[]float64] {
return Column[[]float64]{
idx: idx,
dataType: DataTypeFloat64Array,
get: (*RecordVals).AsFloat64Array,
}
}

func StringArrayColumn(idx int) Column[[]string] {
return Column[[]string]{
idx: idx,
dataType: DataTypeStringArray,
get: (*RecordVals).AsStringArray,
}
}
51 changes: 34 additions & 17 deletions query/exec/aggregator.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,23 @@ const (
ExecutorStateDone
)

var (
aggColTokenIn = query.StringColumn(0)
aggColMinIn = query.Float64Column(1)
aggColMaxIn = query.Float64Column(2)
aggColSumIn = query.Float64Column(3)
aggColTotalIn = query.Uint64Column(4)
aggColTsIn = query.Uint64Column(6)
aggColSamplesIn = query.Float64ArrayColumn(7)
aggColValuesIn = query.StringArrayColumn(8)
)

var (
aggColTokenOut = query.StringColumn(0)
aggColValueOut = query.Float64Column(1)
aggColTsOut = query.Uint64Column(2)
)

// aggKey identifies a single timeseries bin: the grouping token plus the
// (floored) timestamp. ts is 0 (DummyMID) for non-timeseries aggregations, so
// all samples for the same token collapse into one bucket.
Expand Down Expand Up @@ -137,27 +154,27 @@ func (a *DistributedAggregator) drainInput(input query.RecordProducer) {
}

key := aggKey{
token: r.Vals[0].AsString(),
ts: r.Vals[6].AsUint64(),
token: aggColTokenIn.Val(r),
ts: aggColTsIn.Val(r),
}

a.mu.Lock()

s, exists := a.buckets[key]
if !exists {
s = seq.NewSamplesContainers()
s.Min = r.Vals[1].AsFloat64()
s.Max = r.Vals[2].AsFloat64()
s.Min = aggColMinIn.Val(r)
s.Max = aggColMaxIn.Val(r)
} else {
s.Min = min(s.Min, r.Vals[1].AsFloat64())
s.Max = max(s.Max, r.Vals[2].AsFloat64())
s.Min = min(s.Min, aggColMinIn.Val(r))
s.Max = max(s.Max, aggColMaxIn.Val(r))
}

s.Sum += r.Vals[3].AsFloat64()
s.Total += int64(r.Vals[4].AsUint64())
s.Sum += aggColSumIn.Val(r)
s.Total += int64(aggColTotalIn.Val(r))

if a.aggFunc == seq.AggFuncQuantile {
for _, v := range r.Vals[7].AsFloat64Array() {
for _, v := range aggColSamplesIn.Val(r) {
s.InsertSample(v)
}
}
Expand All @@ -171,7 +188,7 @@ func (a *DistributedAggregator) drainInput(input query.RecordProducer) {
m = make(map[string]struct{})
a.values[key] = m
}
for _, v := range r.Vals[8].AsStringArray() {
for _, v := range aggColValuesIn.Val(r) {
m[v] = struct{}{}
}
}
Expand Down Expand Up @@ -199,21 +216,21 @@ func (a *DistributedAggregator) Finalize() *query.Summary {
}

func sortBuckets(aggFunc seq.AggFunc, buckets []*query.Record) {
// ts (Vals[2]) is the primary key (ASC), matching seq/qpr.go sortBuckets
// ts is the primary key (ASC), matching seq/qpr.go sortBuckets
// where MID comes first. Within the same ts buckets are ordered by value.
sortByTsValueDescNameAsc := func(left, right *query.Record) int {
return cmp.Or(
cmp.Compare(left.Vals[2].AsUint64(), right.Vals[2].AsUint64()),
cmp.Compare(right.Vals[1].AsFloat64(), left.Vals[1].AsFloat64()),
cmp.Compare(left.Vals[0].AsString(), right.Vals[0].AsString()),
cmp.Compare(aggColTsOut.Val(left), aggColTsOut.Val(right)),
cmp.Compare(aggColValueOut.Val(right), aggColValueOut.Val(left)),
cmp.Compare(aggColTokenOut.Val(left), aggColTokenOut.Val(right)),
)
}

sortByTsValueNameAsc := func(left, right *query.Record) int {
return cmp.Or(
cmp.Compare(left.Vals[2].AsUint64(), right.Vals[2].AsUint64()),
cmp.Compare(left.Vals[1].AsFloat64(), right.Vals[1].AsFloat64()),
cmp.Compare(left.Vals[0].AsString(), right.Vals[0].AsString()),
cmp.Compare(aggColTsOut.Val(left), aggColTsOut.Val(right)),
cmp.Compare(aggColValueOut.Val(left), aggColValueOut.Val(right)),
cmp.Compare(aggColTokenOut.Val(left), aggColTokenOut.Val(right)),
)
}

Expand Down
18 changes: 6 additions & 12 deletions query/exec/filter.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,14 +12,10 @@ type FilterExpr[T any] interface {
Eval(T) bool
}

type ValGetter[T any] func(*query.RecordVals) T

type Filter[T any] struct {
input query.RecordProducer

colIdx int
getVal ValGetter[T]
expr FilterExpr[T]
col query.Column[T]
expr FilterExpr[T]

// withTotal requests the accurate total of records that pass the filter.
// When true, Finalize drains the (possibly partially consumed) input to the
Expand All @@ -30,7 +26,7 @@ type Filter[T any] struct {
// draining the remaining input.
passed uint64

// roots holds every record whose colIdx val has been decoded (and thus
// roots holds every record whose val has been decoded (and thus
// Spawn'd an insaneJSON root). They are released back to the library pool in
// Finalize. Record.Release is idempotent, so records that were forwarded
// downstream (and released there too) are safe to release here as well.
Expand All @@ -39,15 +35,13 @@ type Filter[T any] struct {

func NewFilter[T any](
input query.RecordProducer,
colIdx int,
get ValGetter[T],
col query.Column[T],
expr FilterExpr[T],
withTotal bool,
) *Filter[T] {
return &Filter[T]{
input: input,
colIdx: colIdx,
getVal: get,
col: col,
expr: expr,
withTotal: withTotal,
}
Expand All @@ -60,7 +54,7 @@ func (f *Filter[T]) Next() *query.Record {
return nil
}

passes := f.expr.Eval(f.getVal(r.Vals[f.colIdx]))
passes := f.expr.Eval(f.col.Val(r))
// The decoded root is now cached; keep a reference so Finalize can release it.
f.roots = append(f.roots, r)
if passes {
Expand Down
17 changes: 8 additions & 9 deletions query/exec/filter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ func TestFilterEq(t *testing.T) {

filterExpr := NewEq[uint32](cond)

testFilter(t, 0, (*query.RecordVals).AsUint32, filterExpr, func(r *query.Record) bool {
testFilter(t, query.Uint32Column(0), filterExpr, func(r *query.Record) bool {
return r.Vals[0].AsUint32() == uint32(cond)
})
}
Expand All @@ -26,7 +26,7 @@ func TestFilterGt(t *testing.T) {

filterExpr := NewGt[uint32](cond)

testFilter(t, 0, (*query.RecordVals).AsUint32, filterExpr, func(r *query.Record) bool {
testFilter(t, query.Uint32Column(0), filterExpr, func(r *query.Record) bool {
return r.Vals[0].AsUint32() > uint32(cond)
})
}
Expand All @@ -36,7 +36,7 @@ func TestFilterLt(t *testing.T) {

filterExpr := NewLt[uint32](cond)

testFilter(t, 0, (*query.RecordVals).AsUint32, filterExpr, func(r *query.Record) bool {
testFilter(t, query.Uint32Column(0), filterExpr, func(r *query.Record) bool {
return r.Vals[0].AsUint32() < uint32(cond)
})
}
Expand All @@ -49,16 +49,15 @@ func TestDocumentFilter(t *testing.T) {

filterExpr := NewDocFilter(field, NewEq(cond))

testFilter(t, 1, (*query.RecordVals).AsDoc, filterExpr, func(r *query.Record) bool {
testFilter(t, query.DocColumn(1), filterExpr, func(r *query.Record) bool {
field := r.Vals[1].AsDoc().Dig(field)
return field.AsString() == cond
})
}

func testFilter[T any](
t *testing.T,
colIdx int,
get ValGetter[T],
col query.Column[T],
filterExpr FilterExpr[T],
wantFilterFunc func(*query.Record) bool,
) {
Expand All @@ -74,7 +73,7 @@ func testFilter[T any](
}
}

filter := NewFilter(&input, colIdx, get, filterExpr, false)
filter := NewFilter(&input, col, filterExpr, false)

outputData := make([]*query.Record, 0)
for r := filter.Next(); r != nil; r = filter.Next() {
Expand All @@ -101,7 +100,7 @@ func TestFilterTotalDrainsInput(t *testing.T) {
}
}

filter := NewFilter(&input, 0, (*query.RecordVals).AsUint32, filterExpr, true)
filter := NewFilter(&input, query.Uint32Column(0), filterExpr, true)
outputData := make([]*query.Record, 0)
for i := 0; i < len(wantData); i++ {
r := filter.Next()
Expand All @@ -126,7 +125,7 @@ func TestFilterTotalErrorPropagated(t *testing.T) {
err: assertErr,
}

filter := NewFilter(&input, 0, (*query.RecordVals).AsUint32, filterExpr, true)
filter := NewFilter(&input, query.Uint32Column(0), filterExpr, true)
for r := filter.Next(); r != nil; r = filter.Next() {
}

Expand Down
Loading
Loading