WIP: implement Min/Max BSI queries

This commit is contained in:
Travis Turner 2018-04-10 11:16:27 -05:00
parent ea2921c192
commit 16539bbab1
No known key found for this signature in database
GPG key ID: 7F08008DFD9314C9
6 changed files with 445 additions and 7 deletions

View file

@ -332,14 +332,44 @@ func TestClient_ImportValue(t *testing.T) {
t.Fatal(err)
}
// Verify Sum.
sum, cnt, err := frame.FieldSum(nil, fld.Name)
if err != nil {
t.Fatal(err)
}
// Verify data.
if sum != 50 || cnt != 3 {
t.Fatalf("unexpected values: got sum=%v, count=%v; expected sum=70, cnt=3", sum, cnt)
t.Fatalf("unexpected values: got sum=%v, count=%v; expected sum=50, cnt=3", sum, cnt)
}
// Verify Min.
min, cnt, err := frame.FieldMin(nil, fld.Name)
if err != nil {
t.Fatal(err)
}
if min != -10 || cnt != 1 {
t.Fatalf("unexpected values: got min=%v, count=%v; expected min=-10, cnt=1", min, cnt)
}
// Verify Min with Filter.
filter, err := frame.FieldRange(fld.Name, pql.GT, 40)
if err != nil {
t.Fatal(err)
}
min, cnt, err = frame.FieldMin(filter, fld.Name) // TODO: change this to use the client
if err != nil {
t.Fatal(err)
}
if min != -100 || cnt != 0 {
t.Fatalf("unexpected values: got min=%v, count=%v; expected min=-100, cnt=0", min, cnt)
}
// Verify Max.
max, cnt, err := frame.FieldMax(nil, fld.Name)
if err != nil {
t.Fatal(err)
}
if max != 40 || cnt != 1 {
t.Fatalf("unexpected values: got max=%v, count=%v; expected max=40, cnt=1", max, cnt)
}
}

View file

@ -157,6 +157,12 @@ func (e *Executor) executeCall(ctx context.Context, index string, c *pql.Call, s
case "Sum":
e.Holder.Stats.CountWithCustomTags(c.Name, 1, 1.0, []string{indexTag})
return e.executeSum(ctx, index, c, slices, opt)
case "Min":
e.Holder.Stats.CountWithCustomTags(c.Name, 1, 1.0, []string{indexTag})
return e.executeFieldMin(ctx, index, c, slices, opt)
case "Max":
e.Holder.Stats.CountWithCustomTags(c.Name, 1, 1.0, []string{indexTag})
return e.executeFieldMax(ctx, index, c, slices, opt)
case "ClearBit":
return e.executeClearBit(ctx, index, c, opt)
case "Count":
@ -233,6 +239,76 @@ func (e *Executor) executeSum(ctx context.Context, index string, c *pql.Call, sl
return other, nil
}
// executeFieldMin executes a Min() call.
func (e *Executor) executeFieldMin(ctx context.Context, index string, c *pql.Call, slices []uint64, opt *ExecOptions) (MinCount, error) {
if frame, _ := c.Args["frame"]; frame == "" {
return MinCount{}, errors.New("Min(): frame required")
} else if field, _ := c.Args["field"]; field == "" {
return MinCount{}, errors.New("Min(): field required")
}
if len(c.Children) > 1 {
return MinCount{}, errors.New("Min() only accepts a single bitmap input")
}
// Execute calls in bulk on each remote node and merge.
mapFn := func(slice uint64) (interface{}, error) {
return e.executeFieldMinSlice(ctx, index, c, slice)
}
// Merge returned results at coordinating node.
reduceFn := func(prev, v interface{}) interface{} {
other, _ := prev.(MinCount)
return other.Smaller(v.(MinCount))
}
result, err := e.mapReduce(ctx, index, slices, c, opt, mapFn, reduceFn)
if err != nil {
return MinCount{}, err
}
other, _ := result.(MinCount)
if other.Count == 0 {
return MinCount{}, nil
}
return other, nil
}
// executeFieldMax executes a Max() call.
func (e *Executor) executeFieldMax(ctx context.Context, index string, c *pql.Call, slices []uint64, opt *ExecOptions) (MaxCount, error) {
if frame, _ := c.Args["frame"]; frame == "" {
return MaxCount{}, errors.New("Max(): frame required")
} else if field, _ := c.Args["field"]; field == "" {
return MaxCount{}, errors.New("Max(): field required")
}
if len(c.Children) > 1 {
return MaxCount{}, errors.New("Max() only accepts a single bitmap input")
}
// Execute calls in bulk on each remote node and merge.
mapFn := func(slice uint64) (interface{}, error) {
return e.executeFieldMaxSlice(ctx, index, c, slice)
}
// Merge returned results at coordinating node.
reduceFn := func(prev, v interface{}) interface{} {
other, _ := prev.(MaxCount)
return other.Larger(v.(MaxCount))
}
result, err := e.mapReduce(ctx, index, slices, c, opt, mapFn, reduceFn)
if err != nil {
return MaxCount{}, err
}
other, _ := result.(MaxCount)
if other.Count == 0 {
return MaxCount{}, nil
}
return other, nil
}
// executeBitmapCall executes a call that returns a bitmap.
func (e *Executor) executeBitmapCall(ctx context.Context, index string, c *pql.Call, slices []uint64, opt *ExecOptions) (*Bitmap, error) {
// Execute calls in bulk on each remote node and merge.
@ -318,7 +394,7 @@ func (e *Executor) executeBitmapCallSlice(ctx context.Context, index string, c *
}
}
// executeSumCountSlice executes calculates the sum & count for fields on a slice.
// executeSumCountSlice calculates the sum and count for fields on a slice.
func (e *Executor) executeSumCountSlice(ctx context.Context, index string, c *pql.Call, slice uint64) (SumCount, error) {
var filter *Bitmap
if len(c.Children) == 1 {
@ -342,12 +418,12 @@ func (e *Executor) executeSumCountSlice(ctx context.Context, index string, c *pq
return SumCount{}, nil
}
view := e.Holder.Fragment(index, frameName, ViewFieldPrefix+fieldName, slice)
if view == nil {
fragment := e.Holder.Fragment(index, frameName, ViewFieldPrefix+fieldName, slice)
if fragment == nil {
return SumCount{}, nil
}
vsum, vcount, err := view.FieldSum(filter, field.BitDepth())
vsum, vcount, err := fragment.FieldSum(filter, field.BitDepth())
if err != nil {
return SumCount{}, err
}
@ -357,6 +433,84 @@ func (e *Executor) executeSumCountSlice(ctx context.Context, index string, c *pq
}, nil
}
// executeFieldMinSlice calculates the min for fields on a slice.
func (e *Executor) executeFieldMinSlice(ctx context.Context, index string, c *pql.Call, slice uint64) (MinCount, error) {
var filter *Bitmap
if len(c.Children) == 1 {
bm, err := e.executeBitmapCallSlice(ctx, index, c.Children[0], slice)
if err != nil {
return MinCount{}, err
}
filter = bm
}
frameName, _ := c.Args["frame"].(string)
fieldName, _ := c.Args["field"].(string)
frame := e.Holder.Frame(index, frameName)
if frame == nil {
return MinCount{}, nil
}
field := frame.Field(fieldName)
if field == nil {
return MinCount{}, nil
}
fragment := e.Holder.Fragment(index, frameName, ViewFieldPrefix+fieldName, slice)
if fragment == nil {
return MinCount{}, nil
}
fmin, fcount, err := fragment.FieldMin(filter, field.BitDepth())
if err != nil {
return MinCount{}, err
}
return MinCount{
Min: int64(fmin) + field.Min,
Count: int64(fcount),
}, nil
}
// executeFieldMaxSlice calculates the max for fields on a slice.
func (e *Executor) executeFieldMaxSlice(ctx context.Context, index string, c *pql.Call, slice uint64) (MaxCount, error) {
var filter *Bitmap
if len(c.Children) == 1 {
bm, err := e.executeBitmapCallSlice(ctx, index, c.Children[0], slice)
if err != nil {
return MaxCount{}, err
}
filter = bm
}
frameName, _ := c.Args["frame"].(string)
fieldName, _ := c.Args["field"].(string)
frame := e.Holder.Frame(index, frameName)
if frame == nil {
return MaxCount{}, nil
}
field := frame.Field(fieldName)
if field == nil {
return MaxCount{}, nil
}
fragment := e.Holder.Fragment(index, frameName, ViewFieldPrefix+fieldName, slice)
if fragment == nil {
return MaxCount{}, nil
}
fmax, fcount, err := fragment.FieldMax(filter, field.BitDepth())
if err != nil {
return MaxCount{}, err
}
return MaxCount{
Max: int64(fmax) + field.Min,
Count: int64(fcount),
}, nil
}
// executeTopN executes a TopN() call.
// This first performs the TopN() to determine the top results and then
// requeries to retrieve the full counts for each of the top results.
@ -1619,3 +1773,37 @@ func decodeSumCount(pb *internal.SumCount) SumCount {
Count: pb.Count,
}
}
// MinCount represents a grouping of min and count for Min() calls.
type MinCount struct {
Min int64 `json:"min"`
Count int64 `json:"count"`
}
// Smaller returns the smaller of the two MinCounts.
func (mc *MinCount) Smaller(other MinCount) MinCount {
if mc.Count == 0 || other.Count < mc.Count {
return other
}
return MinCount{
Min: mc.Min,
Count: mc.Count,
}
}
// MaxCount represents a grouping of max and count for Max() calls.
type MaxCount struct {
Max int64 `json:"max"`
Count int64 `json:"count"`
}
// Larger returns the larger of the two MaxCounts.
func (mc *MaxCount) Larger(other MaxCount) MaxCount {
if mc.Count == 0 || other.Count < mc.Count {
return other
}
return MaxCount{
Max: mc.Max,
Count: mc.Count,
}
}

View file

@ -614,6 +614,70 @@ func (f *Fragment) FieldSum(filter *Bitmap, bitDepth uint) (sum, count uint64, e
return sum, count, nil
}
// FieldMin returns the min of a given field as well as the number of columns involved.
// A bitmap can be passed in to optionally filter the computed columns.
func (f *Fragment) FieldMin(filter *Bitmap, bitDepth uint) (min, count uint64, err error) {
consider := f.Row(uint64(bitDepth))
if filter != nil {
consider = consider.Intersect(filter)
}
// If there are no columns to consider, return early.
if consider.Count() == 0 {
return 0, 0, nil
}
for i := bitDepth; i > uint(0); i-- {
ii := i - 1 // allow for uint range: (bitdepth-1) to 0
row := f.Row(uint64(ii))
x := consider.Difference(row)
count = x.Count()
if count > 0 {
consider = x
} else {
min += (1 << ii)
if ii == 0 {
count = consider.Count()
}
}
}
return min, count, nil
}
// FieldMax returns the max of a given field as well as the number of columns involved.
// A bitmap can be passed in to optionally filter the computed columns.
func (f *Fragment) FieldMax(filter *Bitmap, bitDepth uint) (max, count uint64, err error) {
consider := f.Row(uint64(bitDepth))
if filter != nil {
consider = consider.Intersect(filter)
}
// If there are no columns to consider, return early.
if consider.Count() == 0 {
return 0, 0, nil
}
for i := bitDepth; i > uint(0); i-- {
ii := i - 1 // allow for uint range: (bitdepth-1) to 0
row := f.Row(uint64(ii))
x := row.Intersect(consider)
count = x.Count()
if count > 0 {
max += (1 << ii)
consider = x
} else if ii == 0 {
count = consider.Count()
}
}
return max, count, nil
}
// FieldRange returns bitmaps with a field value encoding matching the predicate.
func (f *Fragment) FieldRange(op pql.Token, bitDepth uint, predicate uint64) (*Bitmap, error) {
switch op {

View file

@ -256,6 +256,79 @@ func TestFragment_FieldSum(t *testing.T) {
})
}
// Ensure a fragment can find the max of field values.
func TestFragment_FieldMinMax(t *testing.T) {
const bitDepth = 16
f := test.MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "")
defer f.Close()
// Set values.
if _, err := f.SetFieldValue(1000, bitDepth, 382); err != nil {
t.Fatal(err)
} else if _, err := f.SetFieldValue(2000, bitDepth, 300); err != nil {
t.Fatal(err)
} else if _, err := f.SetFieldValue(3000, bitDepth, 2818); err != nil {
t.Fatal(err)
} else if _, err := f.SetFieldValue(4000, bitDepth, 300); err != nil {
t.Fatal(err)
} else if _, err := f.SetFieldValue(5000, bitDepth, 2818); err != nil {
t.Fatal(err)
} else if _, err := f.SetFieldValue(6000, bitDepth, 2817); err != nil {
t.Fatal(err)
} else if _, err := f.SetFieldValue(7000, bitDepth, 0); err != nil {
t.Fatal(err)
}
t.Run("Min", func(t *testing.T) {
tests := []struct {
filter *pilosa.Bitmap
exp uint64
cnt uint64
}{
{filter: nil, exp: 0, cnt: 1},
{filter: pilosa.NewBitmap(2000, 4000, 5000), exp: 300, cnt: 2},
{filter: pilosa.NewBitmap(2000, 4000), exp: 300, cnt: 2},
{filter: pilosa.NewBitmap(1), exp: 0, cnt: 0},
{filter: pilosa.NewBitmap(1000), exp: 382, cnt: 1},
{filter: pilosa.NewBitmap(7000), exp: 0, cnt: 1},
}
for i, test := range tests {
if min, cnt, err := f.FieldMin(test.filter, bitDepth); err != nil {
t.Fatal(err)
} else if min != test.exp {
t.Errorf("test %d expected min: %v, but got: %v", i, test.exp, min)
} else if cnt != test.cnt {
t.Errorf("test %d expected cnt: %v, but got: %v", i, test.cnt, cnt)
}
}
})
t.Run("Max", func(t *testing.T) {
tests := []struct {
filter *pilosa.Bitmap
exp uint64
cnt uint64
}{
{filter: nil, exp: 2818, cnt: 2},
{filter: pilosa.NewBitmap(2000, 4000, 5000), exp: 2818, cnt: 1},
{filter: pilosa.NewBitmap(2000, 4000), exp: 300, cnt: 2},
{filter: pilosa.NewBitmap(1), exp: 0, cnt: 0},
{filter: pilosa.NewBitmap(1000), exp: 382, cnt: 1},
{filter: pilosa.NewBitmap(7000), exp: 0, cnt: 1},
}
for i, test := range tests {
if max, cnt, err := f.FieldMax(test.filter, bitDepth); err != nil {
t.Fatal(err)
} else if max != test.exp {
t.Errorf("test %d expected max: %v, but got: %v", i, test.exp, max)
} else if cnt != test.cnt {
t.Errorf("test %d expected cnt: %v, but got: %v", i, test.cnt, cnt)
}
}
})
}
// Ensure a fragment query for matching fields.
func TestFragment_FieldRange(t *testing.T) {
const bitDepth = 16

View file

@ -763,6 +763,46 @@ func (f *Frame) FieldSum(filter *Bitmap, name string) (sum, count int64, err err
return int64(vsum) + (int64(vcount) * field.Min), int64(vcount), nil
}
// FieldMin returns the min for a field.
// An optional filtering bitmap can be provided.
func (f *Frame) FieldMin(filter *Bitmap, name string) (min, count int64, err error) {
field := f.Field(name)
if field == nil {
return 0, 0, ErrFieldNotFound
}
view := f.View(ViewFieldPrefix + name)
if view == nil {
return 0, 0, nil
}
vmin, vcount, err := view.FieldMin(filter, field.BitDepth())
if err != nil {
return 0, 0, err
}
return int64(vmin) + field.Min, int64(vcount), nil
}
// FieldMax returns the max for a field.
// An optional filtering bitmap can be provided.
func (f *Frame) FieldMax(filter *Bitmap, name string) (max, count int64, err error) {
field := f.Field(name)
if field == nil {
return 0, 0, ErrFieldNotFound
}
view := f.View(ViewFieldPrefix + name)
if view == nil {
return 0, 0, nil
}
vmax, vcount, err := view.FieldMax(filter, field.BitDepth())
if err != nil {
return 0, 0, err
}
return int64(vmax) + field.Min, int64(vcount), nil
}
func (f *Frame) FieldRange(name string, op pql.Token, predicate int64) (*Bitmap, error) {
// Retrieve and validate field.
field := f.Field(name)

43
view.go
View file

@ -357,6 +357,49 @@ func (v *View) FieldSum(filter *Bitmap, bitDepth uint) (sum, count uint64, err e
return sum, count, nil
}
// FieldMin returns the min and count of a field.
func (v *View) FieldMin(filter *Bitmap, bitDepth uint) (min, count uint64, err error) {
var minHasValue bool
for _, f := range v.Fragments() {
fmin, fcount, err := f.FieldMin(filter, bitDepth)
if err != nil {
return min, count, err
}
// Don't consider a min based on zero columns.
if fcount == 0 {
continue
}
if !minHasValue {
min = fmin
minHasValue = true
count += fcount
continue
}
if fmin < min {
min = fmin
count += fcount
}
}
return min, count, nil
}
// FieldMax returns the max and count of a field.
func (v *View) FieldMax(filter *Bitmap, bitDepth uint) (max, count uint64, err error) {
for _, f := range v.Fragments() {
fmax, fcount, err := f.FieldMax(filter, bitDepth)
if err != nil {
return max, count, err
}
if fcount > 0 && fmax > max {
max = fmax
count += fcount
}
}
return max, count, nil
}
// FieldRange returns bitmaps with a field value encoding matching the predicate.
func (v *View) FieldRange(op pql.Token, bitDepth uint, predicate uint64) (*Bitmap, error) {
bm := NewBitmap()