From 16539bbab12ccb1743714cd59f0e27e63406c228 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 10 Apr 2018 11:16:27 -0500 Subject: [PATCH] WIP: implement Min/Max BSI queries --- client_test.go | 36 ++++++++- executor.go | 196 ++++++++++++++++++++++++++++++++++++++++++++++- fragment.go | 64 ++++++++++++++++ fragment_test.go | 73 ++++++++++++++++++ frame.go | 40 ++++++++++ view.go | 43 +++++++++++ 6 files changed, 445 insertions(+), 7 deletions(-) diff --git a/client_test.go b/client_test.go index cdbaba788..4db9156c5 100644 --- a/client_test.go +++ b/client_test.go @@ -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) } } diff --git a/executor.go b/executor.go index e74cd63ba..022e48dda 100644 --- a/executor.go +++ b/executor.go @@ -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, + } +} diff --git a/fragment.go b/fragment.go index 438a80ae7..f40c0c8d9 100644 --- a/fragment.go +++ b/fragment.go @@ -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 { diff --git a/fragment_test.go b/fragment_test.go index 6a4a2f377..6b49b5d75 100644 --- a/fragment_test.go +++ b/fragment_test.go @@ -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 diff --git a/frame.go b/frame.go index def41d1b1..a8213e264 100644 --- a/frame.go +++ b/frame.go @@ -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) diff --git a/view.go b/view.go index f775f1b0e..dbcc73b8d 100644 --- a/view.go +++ b/view.go @@ -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()