2020-03-10 05:00:52 +00:00
|
|
|
package storageflux
|
2018-10-05 22:02:31 +00:00
|
|
|
|
2020-03-10 05:00:52 +00:00
|
|
|
//go:generate env GO111MODULE=on go run github.com/benbjohnson/tmpl -data=@types.tmpldata table.gen.go.tmpl
|
2018-10-05 22:02:31 +00:00
|
|
|
|
|
|
|
import (
|
2019-07-03 14:26:08 +00:00
|
|
|
"errors"
|
2018-10-31 23:21:54 +00:00
|
|
|
"sync/atomic"
|
2018-10-05 22:02:31 +00:00
|
|
|
|
2018-12-17 23:25:49 +00:00
|
|
|
"github.com/apache/arrow/go/arrow/array"
|
2018-10-05 22:02:31 +00:00
|
|
|
"github.com/influxdata/flux"
|
2018-12-17 23:25:49 +00:00
|
|
|
"github.com/influxdata/flux/arrow"
|
2018-10-05 22:02:31 +00:00
|
|
|
"github.com/influxdata/flux/execute"
|
2018-12-17 23:25:49 +00:00
|
|
|
"github.com/influxdata/flux/memory"
|
2020-07-28 22:59:11 +00:00
|
|
|
"github.com/influxdata/influxdb/v2/models"
|
2018-10-05 22:02:31 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
type table struct {
|
|
|
|
bounds execute.Bounds
|
|
|
|
key flux.GroupKey
|
|
|
|
cols []flux.ColMeta
|
|
|
|
|
|
|
|
// cache of the tags on the current series.
|
|
|
|
// len(tags) == len(colMeta)
|
|
|
|
tags [][]byte
|
|
|
|
defs [][]byte
|
|
|
|
|
2020-10-22 21:34:22 +00:00
|
|
|
done chan struct{}
|
|
|
|
empty bool
|
2018-10-05 22:02:31 +00:00
|
|
|
|
2019-05-30 17:31:54 +00:00
|
|
|
colBufs *colReader
|
2018-10-05 22:02:31 +00:00
|
|
|
|
|
|
|
err error
|
|
|
|
|
2019-07-03 14:26:08 +00:00
|
|
|
cancelled, used int32
|
2019-11-27 14:31:53 +00:00
|
|
|
cache *tagsCache
|
2019-07-03 14:26:08 +00:00
|
|
|
alloc *memory.Allocator
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func newTable(
|
2018-10-31 23:21:54 +00:00
|
|
|
done chan struct{},
|
2018-10-05 22:02:31 +00:00
|
|
|
bounds execute.Bounds,
|
|
|
|
key flux.GroupKey,
|
|
|
|
cols []flux.ColMeta,
|
|
|
|
defs [][]byte,
|
2019-11-27 14:31:53 +00:00
|
|
|
cache *tagsCache,
|
2019-03-06 21:29:45 +00:00
|
|
|
alloc *memory.Allocator,
|
2018-10-05 22:02:31 +00:00
|
|
|
) table {
|
|
|
|
return table{
|
2019-05-30 17:31:54 +00:00
|
|
|
done: done,
|
|
|
|
bounds: bounds,
|
|
|
|
key: key,
|
|
|
|
tags: make([][]byte, len(cols)),
|
|
|
|
defs: defs,
|
|
|
|
cols: cols,
|
2019-11-27 14:31:53 +00:00
|
|
|
cache: cache,
|
2019-05-30 17:31:54 +00:00
|
|
|
alloc: alloc,
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (t *table) Key() flux.GroupKey { return t.key }
|
|
|
|
func (t *table) Cols() []flux.ColMeta { return t.cols }
|
|
|
|
func (t *table) Err() error { return t.err }
|
2020-04-24 18:31:46 +00:00
|
|
|
func (t *table) Empty() bool { return t.colBufs == nil || t.colBufs.l == 0 }
|
2018-10-05 22:02:31 +00:00
|
|
|
|
2018-10-31 23:21:54 +00:00
|
|
|
func (t *table) Cancel() {
|
|
|
|
atomic.StoreInt32(&t.cancelled, 1)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (t *table) isCancelled() bool {
|
|
|
|
return atomic.LoadInt32(&t.cancelled) != 0
|
|
|
|
}
|
|
|
|
|
2020-10-22 21:34:22 +00:00
|
|
|
func (t *table) init(advance func() bool) {
|
|
|
|
t.empty = !advance() && t.err == nil
|
|
|
|
}
|
|
|
|
|
2019-07-03 14:26:08 +00:00
|
|
|
func (t *table) do(f func(flux.ColReader) error, advance func() bool) error {
|
|
|
|
// Mark this table as having been used. If this doesn't
|
|
|
|
// succeed, then this has already been invoked somewhere else.
|
|
|
|
if !atomic.CompareAndSwapInt32(&t.used, 0, 1) {
|
|
|
|
return errors.New("table already used")
|
|
|
|
}
|
|
|
|
defer t.closeDone()
|
|
|
|
|
|
|
|
if !t.Empty() {
|
|
|
|
t.err = f(t.colBufs)
|
|
|
|
t.colBufs.Release()
|
|
|
|
|
|
|
|
for !t.isCancelled() && t.err == nil && advance() {
|
|
|
|
t.err = f(t.colBufs)
|
|
|
|
t.colBufs.Release()
|
|
|
|
}
|
|
|
|
t.colBufs = nil
|
|
|
|
}
|
|
|
|
|
|
|
|
return t.err
|
|
|
|
}
|
|
|
|
|
|
|
|
func (t *table) Done() {
|
|
|
|
// Mark the table as having been used. If this has already
|
|
|
|
// been done, then nothing needs to be done.
|
|
|
|
if atomic.CompareAndSwapInt32(&t.used, 0, 1) {
|
|
|
|
defer t.closeDone()
|
|
|
|
}
|
|
|
|
|
|
|
|
if t.colBufs != nil {
|
|
|
|
t.colBufs.Release()
|
|
|
|
t.colBufs = nil
|
2019-05-30 17:31:54 +00:00
|
|
|
}
|
2019-07-03 14:26:08 +00:00
|
|
|
}
|
2019-05-30 17:31:54 +00:00
|
|
|
|
2019-07-03 14:26:08 +00:00
|
|
|
// allocateBuffer will allocate a suitable buffer for the
|
|
|
|
// table implementations to use. If the existing buffer
|
|
|
|
// is not used anymore, then it may be reused.
|
|
|
|
//
|
|
|
|
// The allocated buffer can be accessed at colBufs or
|
|
|
|
// through the returned colReader.
|
|
|
|
func (t *table) allocateBuffer(l int) *colReader {
|
|
|
|
if t.colBufs == nil || atomic.LoadInt64(&t.colBufs.refCount) > 0 {
|
|
|
|
// The current buffer is still being used so we should
|
|
|
|
// generate a new one.
|
|
|
|
t.colBufs = &colReader{
|
|
|
|
key: t.key,
|
|
|
|
colMeta: t.cols,
|
|
|
|
cols: make([]array.Interface, len(t.cols)),
|
|
|
|
}
|
2019-05-30 17:31:54 +00:00
|
|
|
}
|
2019-07-03 14:26:08 +00:00
|
|
|
t.colBufs.refCount = 1
|
|
|
|
t.colBufs.l = l
|
2019-05-30 17:31:54 +00:00
|
|
|
return t.colBufs
|
|
|
|
}
|
|
|
|
|
|
|
|
type colReader struct {
|
|
|
|
refCount int64
|
|
|
|
|
|
|
|
key flux.GroupKey
|
|
|
|
colMeta []flux.ColMeta
|
|
|
|
cols []array.Interface
|
|
|
|
l int
|
|
|
|
}
|
|
|
|
|
|
|
|
func (cr *colReader) Retain() {
|
|
|
|
atomic.AddInt64(&cr.refCount, 1)
|
|
|
|
}
|
|
|
|
func (cr *colReader) Release() {
|
|
|
|
if atomic.AddInt64(&cr.refCount, -1) == 0 {
|
|
|
|
for _, col := range cr.cols {
|
|
|
|
col.Release()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (cr *colReader) Key() flux.GroupKey { return cr.key }
|
|
|
|
func (cr *colReader) Cols() []flux.ColMeta { return cr.colMeta }
|
|
|
|
func (cr *colReader) Len() int { return cr.l }
|
2019-05-28 22:24:26 +00:00
|
|
|
|
2019-05-30 17:31:54 +00:00
|
|
|
func (cr *colReader) Bools(j int) *array.Boolean {
|
|
|
|
execute.CheckColType(cr.colMeta[j], flux.TBool)
|
|
|
|
return cr.cols[j].(*array.Boolean)
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
2019-05-30 17:31:54 +00:00
|
|
|
func (cr *colReader) Ints(j int) *array.Int64 {
|
|
|
|
execute.CheckColType(cr.colMeta[j], flux.TInt)
|
|
|
|
return cr.cols[j].(*array.Int64)
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
2019-05-30 17:31:54 +00:00
|
|
|
func (cr *colReader) UInts(j int) *array.Uint64 {
|
|
|
|
execute.CheckColType(cr.colMeta[j], flux.TUInt)
|
|
|
|
return cr.cols[j].(*array.Uint64)
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
2019-05-30 17:31:54 +00:00
|
|
|
func (cr *colReader) Floats(j int) *array.Float64 {
|
|
|
|
execute.CheckColType(cr.colMeta[j], flux.TFloat)
|
|
|
|
return cr.cols[j].(*array.Float64)
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
2019-05-30 17:31:54 +00:00
|
|
|
func (cr *colReader) Strings(j int) *array.Binary {
|
|
|
|
execute.CheckColType(cr.colMeta[j], flux.TString)
|
|
|
|
return cr.cols[j].(*array.Binary)
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
2019-05-30 17:31:54 +00:00
|
|
|
func (cr *colReader) Times(j int) *array.Int64 {
|
|
|
|
execute.CheckColType(cr.colMeta[j], flux.TTime)
|
|
|
|
return cr.cols[j].(*array.Int64)
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// readTags populates b.tags with the provided tags
|
|
|
|
func (t *table) readTags(tags models.Tags) {
|
|
|
|
for j := range t.tags {
|
|
|
|
t.tags[j] = t.defs[j]
|
|
|
|
}
|
|
|
|
|
|
|
|
if len(tags) == 0 {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
for _, tag := range tags {
|
|
|
|
j := execute.ColIdx(string(tag.Key), t.cols)
|
2020-07-29 20:31:29 +00:00
|
|
|
// In the case of group aggregate, tags that are not referenced in group() are not included in the result, but
|
|
|
|
// readTags () still get a complete tag list. Here is just to skip the tags that should not present in the result.
|
|
|
|
if j < 0 {
|
|
|
|
continue
|
|
|
|
}
|
2018-10-05 22:02:31 +00:00
|
|
|
t.tags[j] = tag.Value
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// appendTags fills the colBufs for the tag columns with the tag value.
|
2019-05-30 17:31:54 +00:00
|
|
|
func (t *table) appendTags(cr *colReader) {
|
2018-10-05 22:02:31 +00:00
|
|
|
for j := range t.cols {
|
|
|
|
v := t.tags[j]
|
|
|
|
if v != nil {
|
2019-11-27 14:31:53 +00:00
|
|
|
cr.cols[j] = t.cache.GetTag(string(v), cr.l, t.alloc)
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// appendBounds fills the colBufs for the time bounds
|
2019-05-30 17:31:54 +00:00
|
|
|
func (t *table) appendBounds(cr *colReader) {
|
2019-11-27 14:31:53 +00:00
|
|
|
start, stop := t.cache.GetBounds(t.bounds, cr.l, t.alloc)
|
|
|
|
cr.cols[startColIdx], cr.cols[stopColIdx] = start, stop
|
2018-10-05 22:02:31 +00:00
|
|
|
}
|
|
|
|
|
2018-10-31 23:21:54 +00:00
|
|
|
func (t *table) closeDone() {
|
|
|
|
if t.done != nil {
|
|
|
|
close(t.done)
|
|
|
|
t.done = nil
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-12-17 23:25:49 +00:00
|
|
|
func (t *floatTable) toArrowBuffer(vs []float64) *array.Float64 {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewFloat(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
|
|
|
func (t *floatGroupTable) toArrowBuffer(vs []float64) *array.Float64 {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewFloat(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
2020-10-22 21:34:22 +00:00
|
|
|
func (t *floatWindowTable) mergeValues(intervals []int64) *array.Float64 {
|
|
|
|
b := arrow.NewFloatBuilder(t.alloc)
|
|
|
|
b.Resize(len(intervals))
|
|
|
|
t.appendValues(intervals, b.Append, b.AppendNull)
|
|
|
|
return b.NewFloat64Array()
|
|
|
|
}
|
2018-12-17 23:25:49 +00:00
|
|
|
func (t *integerTable) toArrowBuffer(vs []int64) *array.Int64 {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewInt(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
|
|
|
func (t *integerGroupTable) toArrowBuffer(vs []int64) *array.Int64 {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewInt(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
2020-10-22 21:34:22 +00:00
|
|
|
func (t *integerWindowTable) mergeValues(intervals []int64) *array.Int64 {
|
|
|
|
b := arrow.NewIntBuilder(t.alloc)
|
|
|
|
b.Resize(len(intervals))
|
|
|
|
appendNull := b.AppendNull
|
|
|
|
if t.fillValue != nil {
|
|
|
|
appendNull = func() { b.Append(*t.fillValue) }
|
|
|
|
}
|
|
|
|
t.appendValues(intervals, b.Append, appendNull)
|
|
|
|
return b.NewInt64Array()
|
|
|
|
}
|
2018-12-17 23:25:49 +00:00
|
|
|
func (t *unsignedTable) toArrowBuffer(vs []uint64) *array.Uint64 {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewUint(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
|
|
|
func (t *unsignedGroupTable) toArrowBuffer(vs []uint64) *array.Uint64 {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewUint(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
2020-10-22 21:34:22 +00:00
|
|
|
func (t *unsignedWindowTable) mergeValues(intervals []int64) *array.Uint64 {
|
|
|
|
b := arrow.NewUintBuilder(t.alloc)
|
|
|
|
b.Resize(len(intervals))
|
|
|
|
t.appendValues(intervals, b.Append, b.AppendNull)
|
|
|
|
return b.NewUint64Array()
|
|
|
|
}
|
2018-12-17 23:25:49 +00:00
|
|
|
func (t *stringTable) toArrowBuffer(vs []string) *array.Binary {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewString(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
|
|
|
func (t *stringGroupTable) toArrowBuffer(vs []string) *array.Binary {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewString(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
2020-10-22 21:34:22 +00:00
|
|
|
func (t *stringWindowTable) mergeValues(intervals []int64) *array.Binary {
|
|
|
|
b := arrow.NewStringBuilder(t.alloc)
|
|
|
|
b.Resize(len(intervals))
|
|
|
|
t.appendValues(intervals, b.AppendString, b.AppendNull)
|
|
|
|
return b.NewBinaryArray()
|
|
|
|
}
|
2018-12-17 23:25:49 +00:00
|
|
|
func (t *booleanTable) toArrowBuffer(vs []bool) *array.Boolean {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewBool(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
|
|
|
func (t *booleanGroupTable) toArrowBuffer(vs []bool) *array.Boolean {
|
2019-03-06 21:29:45 +00:00
|
|
|
return arrow.NewBool(vs, t.alloc)
|
2018-12-17 23:25:49 +00:00
|
|
|
}
|
2020-10-22 21:34:22 +00:00
|
|
|
func (t *booleanWindowTable) mergeValues(intervals []int64) *array.Boolean {
|
|
|
|
b := arrow.NewBoolBuilder(t.alloc)
|
|
|
|
b.Resize(len(intervals))
|
|
|
|
t.appendValues(intervals, b.Append, b.AppendNull)
|
|
|
|
return b.NewBooleanArray()
|
|
|
|
}
|