mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-11 07:11:02 +00:00
Merge pull request #4 from travisturner/new-rows-iterate-cleanup
New rows iterate cleanup
This commit is contained in:
commit
86d238edb0
5 changed files with 918 additions and 327 deletions
|
|
@ -363,9 +363,12 @@ func encodeQueryResponse(m *pilosa.QueryResponse) *internal.QueryResponse {
|
|||
case pilosa.RowIDs:
|
||||
pb.Results[i].Type = queryResultTypeRowIDs
|
||||
pb.Results[i].RowIDs = result
|
||||
case pilosa.GroupByCounts:
|
||||
pb.Results[i].Type = queryResultTypeGroupByCounts
|
||||
pb.Results[i].GroupByCounts = encodeGroupByCount(result)
|
||||
case []pilosa.GroupCount:
|
||||
pb.Results[i].Type = queryResultTypeGroupCounts
|
||||
pb.Results[i].GroupCounts = encodeGroupCounts(result)
|
||||
case pilosa.RowIdentifiers:
|
||||
pb.Results[i].Type = queryResultTypeRowIdentifiers
|
||||
pb.Results[i].RowIdentifiers = encodeRowIdentifiers(result)
|
||||
case nil:
|
||||
pb.Results[i].Type = queryResultTypeNil
|
||||
}
|
||||
|
|
@ -898,7 +901,6 @@ func decodeQueryResponse(pb *internal.QueryResponse, m *pilosa.QueryResponse) {
|
|||
}
|
||||
m.Results = make([]interface{}, len(pb.Results))
|
||||
decodeQueryResults(pb.Results, m.Results)
|
||||
|
||||
}
|
||||
|
||||
func decodeColumnAttrSets(pb []*internal.ColumnAttrSet, m []*pilosa.ColumnAttrSet) {
|
||||
|
|
@ -929,7 +931,8 @@ const (
|
|||
queryResultTypeUint64
|
||||
queryResultTypeBool
|
||||
queryResultTypeRowIDs
|
||||
queryResultTypeGroupByCounts
|
||||
queryResultTypeGroupCounts
|
||||
queryResultTypeRowIdentifiers
|
||||
)
|
||||
|
||||
func decodeQueryResult(pb *internal.QueryResult) interface{} {
|
||||
|
|
@ -946,8 +949,8 @@ func decodeQueryResult(pb *internal.QueryResult) interface{} {
|
|||
return pb.Changed
|
||||
case queryResultTypeNil:
|
||||
return nil
|
||||
case queryResultTypeGroupByCounts:
|
||||
return decodeGroupByCounts(pb.GroupByCounts)
|
||||
case queryResultTypeGroupCounts:
|
||||
return decodeGroupCounts(pb.GroupCounts)
|
||||
}
|
||||
panic(fmt.Sprintf("unknown type: %d", pb.Type))
|
||||
}
|
||||
|
|
@ -998,12 +1001,24 @@ func decodeAttr(attr *internal.Attr) (key string, value interface{}) {
|
|||
}
|
||||
}
|
||||
|
||||
func decodeGroupByCounts(a []*internal.GroupLine) pilosa.GroupByCounts {
|
||||
gbc := make(pilosa.GroupByCounts, 0)
|
||||
func decodeGroupCounts(a []*internal.GroupCount) []pilosa.GroupCount {
|
||||
other := make([]pilosa.GroupCount, len(a))
|
||||
for i := range a {
|
||||
gbc = append(gbc, pilosa.GroupLine{a[i].Groups, a[i].Total})
|
||||
other[i] = pilosa.GroupCount{
|
||||
decodeFieldRows(a[i].Group),
|
||||
a[i].Count,
|
||||
}
|
||||
}
|
||||
return gbc
|
||||
return other
|
||||
}
|
||||
|
||||
func decodeFieldRows(a []*internal.FieldRow) []pilosa.FieldRow {
|
||||
other := make([]pilosa.FieldRow, len(a))
|
||||
for i := range a {
|
||||
other[i].Field = a[i].Field
|
||||
other[i].RowID = a[i].RowID
|
||||
}
|
||||
return other
|
||||
}
|
||||
|
||||
func decodePairs(a []*internal.Pair) []pilosa.Pair {
|
||||
|
|
@ -1057,14 +1072,36 @@ func encodeRow(r *pilosa.Row) *internal.Row {
|
|||
}
|
||||
}
|
||||
|
||||
func encodeGroupByCount(counts pilosa.GroupByCounts) []*internal.GroupLine {
|
||||
result := make([]*internal.GroupLine, len(counts))
|
||||
func encodeRowIdentifiers(r pilosa.RowIdentifiers) *internal.RowIdentifiers {
|
||||
return &internal.RowIdentifiers{
|
||||
Rows: r.Rows,
|
||||
Keys: r.Keys,
|
||||
//Attrs: encodeAttrs(r.Attrs),
|
||||
}
|
||||
}
|
||||
|
||||
func encodeGroupCounts(counts []pilosa.GroupCount) []*internal.GroupCount {
|
||||
result := make([]*internal.GroupCount, len(counts))
|
||||
for i := range counts {
|
||||
result[i] = &internal.GroupLine{Groups: counts[i].Groups, Total: counts[i].Total}
|
||||
result[i] = &internal.GroupCount{
|
||||
Group: encodeFieldRows(counts[i].Group),
|
||||
Count: counts[i].Count,
|
||||
}
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
func encodeFieldRows(a []pilosa.FieldRow) []*internal.FieldRow {
|
||||
other := make([]*internal.FieldRow, len(a))
|
||||
for i := range a {
|
||||
other[i] = &internal.FieldRow{
|
||||
Field: a[i].Field,
|
||||
RowID: a[i].RowID,
|
||||
}
|
||||
}
|
||||
return other
|
||||
}
|
||||
|
||||
func encodePairs(a pilosa.Pairs) []*internal.Pair {
|
||||
other := make([]*internal.Pair, len(a))
|
||||
for i := range a {
|
||||
|
|
|
|||
239
executor.go
239
executor.go
|
|
@ -196,9 +196,9 @@ func (e *executor) executeCall(ctx context.Context, index string, c *pql.Call, s
|
|||
case "TopN":
|
||||
e.Holder.Stats.CountWithCustomTags(c.Name, 1, 1.0, []string{indexTag})
|
||||
return e.executeTopN(ctx, index, c, shards, opt)
|
||||
case "Rows":
|
||||
case "RowIDs":
|
||||
e.Holder.Stats.CountWithCustomTags(c.Name, 1, 1.0, []string{indexTag})
|
||||
return e.executeRows(ctx, index, c, shards, opt)
|
||||
return e.executeRowIDs(ctx, index, c, shards, opt)
|
||||
case "GroupBy":
|
||||
e.Holder.Stats.CountWithCustomTags(c.Name, 1, 1.0, []string{indexTag})
|
||||
return e.executeGroupBy(ctx, index, c, shards, opt)
|
||||
|
|
@ -722,9 +722,23 @@ func (e *executor) executeDifferenceShard(ctx context.Context, index string, c *
|
|||
return other, nil
|
||||
}
|
||||
|
||||
// RowIdentifiers is a return type for a list of
|
||||
// row ids or row keys. The names `Rows` and `Keys`
|
||||
// are meant to follow the same convention as the
|
||||
// Row query which returns `Columns` and `Keys`.
|
||||
// TODO: Rename this to something better. Anything.
|
||||
type RowIdentifiers struct {
|
||||
Rows []uint64 `json:"rows"`
|
||||
Keys []string `json:"keys,omitempty"`
|
||||
}
|
||||
|
||||
// RowIDs is a query return type for just uint64 row ids.
|
||||
// It should only be used internally (since RowIdentifiers
|
||||
// is the external return type), but it is exported because
|
||||
// the proto package needs access to it.
|
||||
type RowIDs []uint64
|
||||
|
||||
func (r RowIDs) Merge(other RowIDs) RowIDs {
|
||||
func (r RowIDs) merge(other RowIDs) RowIDs {
|
||||
i, j := 0, 0
|
||||
result := make(RowIDs, 0)
|
||||
for i < len(r) && j < len(other) {
|
||||
|
|
@ -751,22 +765,24 @@ func (r RowIDs) Merge(other RowIDs) RowIDs {
|
|||
}
|
||||
return result
|
||||
}
|
||||
func (e *executor) executeGroupBy(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *execOptions) (GroupByCounts, error) {
|
||||
|
||||
func (e *executor) executeGroupBy(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *execOptions) ([]GroupCount, error) {
|
||||
// Execute calls in bulk on each remote node and merge.
|
||||
mapFn := func(shard uint64) (interface{}, error) {
|
||||
return e.executeGroupByShard(ctx, index, c, shard)
|
||||
}
|
||||
// Merge returned results at coordinating node.
|
||||
reduceFn := func(prev, v interface{}) interface{} {
|
||||
other, _ := prev.(GroupByCounts)
|
||||
return other.Merge(v.(GroupByCounts))
|
||||
other, _ := prev.([]GroupCount)
|
||||
return mergeGroupCounts(other, v.([]GroupCount))
|
||||
}
|
||||
// Get full result set.
|
||||
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
results, _ := other.(GroupByCounts)
|
||||
results, _ := other.([]GroupCount)
|
||||
|
||||
// Apply offset.
|
||||
if offset, hasOffset, err := c.UintArg("offset"); err != nil {
|
||||
return nil, err
|
||||
|
|
@ -786,55 +802,53 @@ func (e *executor) executeGroupBy(ctx context.Context, index string, c *pql.Call
|
|||
return results, nil
|
||||
}
|
||||
|
||||
// FieldRow is used to distinguish rows in a group by result.
|
||||
type FieldRow struct {
|
||||
Field string `json:"field"`
|
||||
RowID uint64 `json:"rowID"`
|
||||
RowKey string `json:"rowKey,omitempty"`
|
||||
}
|
||||
|
||||
func (fr FieldRow) String() string {
|
||||
return fmt.Sprintf("%s.%d", fr.Field, fr.RowID)
|
||||
}
|
||||
|
||||
// TODO: we shouldn't need to string this
|
||||
func uniqueGroupString(fr []FieldRow) string {
|
||||
s := []string{}
|
||||
for _, f := range fr {
|
||||
s = append(s, f.String())
|
||||
}
|
||||
return strings.Join(s, "-")
|
||||
}
|
||||
|
||||
// gbi is a groupBy item.
|
||||
type gbi struct {
|
||||
row *Row
|
||||
fieldKey string
|
||||
rowID uint64
|
||||
fieldRow FieldRow
|
||||
}
|
||||
type GroupLine struct {
|
||||
Groups []string
|
||||
Total uint64
|
||||
}
|
||||
type GroupByCounts []GroupLine
|
||||
|
||||
func (gbc GroupByCounts) Merge(other GroupByCounts) GroupByCounts {
|
||||
m := make(map[string]struct {
|
||||
i int
|
||||
total uint64
|
||||
})
|
||||
for i := range gbc {
|
||||
m[strings.Join(gbc[i].Groups, "-")] = struct {
|
||||
i int
|
||||
total uint64
|
||||
}{total: gbc[i].Total, i: i}
|
||||
type GroupCount struct {
|
||||
Group []FieldRow `json:"group"`
|
||||
Count uint64 `json:"count"`
|
||||
}
|
||||
|
||||
func mergeGroupCounts(gc, other []GroupCount) []GroupCount {
|
||||
m := make(map[string]int)
|
||||
for i := range gc {
|
||||
m[uniqueGroupString(gc[i].Group)] = i
|
||||
}
|
||||
for i := range other {
|
||||
key := strings.Join(other[i].Groups, "-")
|
||||
o, found := m[key]
|
||||
if found {
|
||||
gbc[o.i].Total += other[i].Total
|
||||
if idx, found := m[uniqueGroupString(other[i].Group)]; found {
|
||||
gc[idx].Count += other[i].Count
|
||||
} else {
|
||||
gbc = append(gbc, other[i])
|
||||
gc = append(gc, other[i])
|
||||
}
|
||||
}
|
||||
return gbc
|
||||
return gc
|
||||
}
|
||||
func makeGroup(parts []gbi) GroupLine {
|
||||
var other *Row
|
||||
line := GroupLine{}
|
||||
for i, o := range parts {
|
||||
if i == 0 {
|
||||
other = o.row
|
||||
} else {
|
||||
other = other.Intersect(o.row)
|
||||
}
|
||||
line.Groups = append(line.Groups, o.fieldKey)
|
||||
}
|
||||
line.Total = other.Count()
|
||||
return line
|
||||
}
|
||||
func (e *executor) executeGroupByShard(ctx context.Context, index string, c *pql.Call, shard uint64) (GroupByCounts, error) {
|
||||
|
||||
func (e *executor) executeGroupByShard(ctx context.Context, index string, c *pql.Call, shard uint64) ([]GroupCount, error) {
|
||||
// Fetch index.
|
||||
idx := e.Holder.Index(index)
|
||||
if idx == nil {
|
||||
|
|
@ -866,8 +880,9 @@ func (e *executor) executeGroupByShard(ctx context.Context, index string, c *pql
|
|||
return nil, errors.Wrap(ErrFieldNotFound, fmt.Sprintf("executeGroupBy: %s", fieldDirective.(string)))
|
||||
}
|
||||
}
|
||||
results := make(GroupByCounts, 0)
|
||||
var work listOfGBILists
|
||||
|
||||
results := make([]GroupCount, 0)
|
||||
var work [][]gbi
|
||||
for _, fieldDirective := range fieldDirectives.([]interface{}) {
|
||||
fieldName := getFieldName(fieldDirective.(string))
|
||||
// Fetch fragment.
|
||||
|
|
@ -880,54 +895,51 @@ func (e *executor) executeGroupByShard(ctx context.Context, index string, c *pql
|
|||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
set := make(gbiList, 0)
|
||||
|
||||
set := make([]gbi, 0)
|
||||
for _, rowID := range frag.rowsWithFilter(filter) {
|
||||
set = append(set, gbi{
|
||||
row: frag.row(rowID),
|
||||
rowID: rowID,
|
||||
fieldKey: fmt.Sprintf("%s.%d", fieldName, rowID),
|
||||
row: frag.row(rowID),
|
||||
fieldRow: FieldRow{
|
||||
Field: fieldName,
|
||||
RowID: rowID,
|
||||
},
|
||||
})
|
||||
}
|
||||
work = append(work, set)
|
||||
}
|
||||
for _, group := range product(work) {
|
||||
group.gl.Total = group.row.Count()
|
||||
if group.gl.Total > 0 {
|
||||
results = append(results, group.gl)
|
||||
group.gCnt.Count = group.row.Count()
|
||||
if group.gCnt.Count > 0 {
|
||||
results = append(results, group.gCnt)
|
||||
}
|
||||
}
|
||||
return results, nil
|
||||
}
|
||||
|
||||
type gbiList []gbi
|
||||
type listOfGBILists []gbiList
|
||||
|
||||
// pi is a product process item.
|
||||
type pi struct {
|
||||
row *Row
|
||||
gl GroupLine
|
||||
// ppi is a product process item.
|
||||
type ppi struct {
|
||||
row *Row
|
||||
gCnt GroupCount
|
||||
}
|
||||
type piList []pi
|
||||
|
||||
// product generates the cartiesian product of the input
|
||||
// using tail recursion
|
||||
func product(input listOfGBILists) piList {
|
||||
res := make(piList, 0)
|
||||
if len(input) == 0 { //base return empty list
|
||||
res = append(res, pi{gl: GroupLine{Groups: make([]string, 0)}})
|
||||
} else {
|
||||
res = productHelper(input, res)
|
||||
// product generates the cartesian product of the input
|
||||
// using tail recursion.
|
||||
func product(input [][]gbi) []ppi {
|
||||
if len(input) == 0 { // base return empty list
|
||||
return []ppi{
|
||||
{gCnt: GroupCount{Group: make([]FieldRow, 0)}},
|
||||
}
|
||||
}
|
||||
return res
|
||||
}
|
||||
func productHelper(lists listOfGBILists, res piList) piList {
|
||||
head := lists[0] //take first element of the list
|
||||
tail := product(lists[1:]) //invoke product on remaining element
|
||||
|
||||
res := make([]ppi, 0)
|
||||
head := input[0] // take first element of the list
|
||||
tail := product(input[1:]) // invoke product on remaining element
|
||||
for h := range head { // for each head
|
||||
for t := range tail { //iterate over the tail
|
||||
s := pi{gl: GroupLine{Groups: make([]string, 0)}}
|
||||
s.gl.Groups = append([]string{head[h].fieldKey}, tail[t].gl.Groups...) //had to insert at the front to match input order
|
||||
if tail[t].row != nil { //first time around nothing to intersect
|
||||
for t := range tail { // iterate over the tail
|
||||
s := ppi{gCnt: GroupCount{Group: make([]FieldRow, 0)}}
|
||||
s.gCnt.Group = append([]FieldRow{head[h].fieldRow}, tail[t].gCnt.Group...) // had to insert at the front to match input order
|
||||
if tail[t].row != nil { // first time around nothing to intersect
|
||||
s.row = head[h].row.Intersect(tail[t].row)
|
||||
} else {
|
||||
s.row = head[h].row
|
||||
|
|
@ -937,15 +949,16 @@ func productHelper(lists listOfGBILists, res piList) piList {
|
|||
}
|
||||
return res
|
||||
}
|
||||
func (e *executor) executeRows(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *execOptions) (RowIDs, error) {
|
||||
|
||||
func (e *executor) executeRowIDs(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *execOptions) (RowIDs, error) {
|
||||
// Execute calls in bulk on each remote node and merge.
|
||||
mapFn := func(shard uint64) (interface{}, error) {
|
||||
return e.executeRowsShard(ctx, index, c, shard)
|
||||
return e.executeRowIDsShard(ctx, index, c, shard)
|
||||
}
|
||||
// Merge returned results at coordinating node.
|
||||
reduceFn := func(prev, v interface{}) interface{} {
|
||||
other, _ := prev.(RowIDs)
|
||||
return other.Merge(v.(RowIDs))
|
||||
return other.merge(v.(RowIDs))
|
||||
}
|
||||
// Get full result set.
|
||||
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
|
||||
|
|
@ -971,7 +984,8 @@ func (e *executor) executeRows(ctx context.Context, index string, c *pql.Call, s
|
|||
}
|
||||
return results, nil
|
||||
}
|
||||
func (e *executor) executeRowsShard(ctx context.Context, index string, c *pql.Call, shard uint64) (RowIDs, error) {
|
||||
|
||||
func (e *executor) executeRowIDsShard(ctx context.Context, index string, c *pql.Call, shard uint64) (RowIDs, error) {
|
||||
// Fetch index.
|
||||
idx := e.Holder.Index(index)
|
||||
if idx == nil {
|
||||
|
|
@ -980,7 +994,7 @@ func (e *executor) executeRowsShard(ctx context.Context, index string, c *pql.Ca
|
|||
// Fetch field name from argument.
|
||||
fieldName, ok := c.Args["field"].(string)
|
||||
if !ok {
|
||||
return nil, errors.New("Rows() argument required: field")
|
||||
return nil, errors.New("RowIDs() argument required: field")
|
||||
}
|
||||
// Fetch field.
|
||||
f := e.Holder.Field(index, fieldName)
|
||||
|
|
@ -2087,7 +2101,62 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
|
|||
return other, nil
|
||||
}
|
||||
}
|
||||
|
||||
case []GroupCount:
|
||||
other := make([]GroupCount, 0)
|
||||
for _, gl := range result {
|
||||
|
||||
group := make([]FieldRow, len(gl.Group))
|
||||
for i, g := range gl.Group {
|
||||
group[i] = g
|
||||
|
||||
// TODO: It may be useful to cache this field lookup.
|
||||
field := idx.Field(g.Field)
|
||||
if field == nil {
|
||||
return nil, ErrFieldNotFound
|
||||
}
|
||||
if field.keys() {
|
||||
key, err := e.TranslateStore.TranslateRowToString(index, g.Field, g.RowID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
group[i].RowKey = key
|
||||
}
|
||||
}
|
||||
|
||||
other = append(other, GroupCount{
|
||||
Group: group,
|
||||
Count: gl.Count,
|
||||
})
|
||||
}
|
||||
return other, nil
|
||||
|
||||
case RowIDs:
|
||||
other := RowIdentifiers{}
|
||||
|
||||
fieldName := callArgString(call, "field")
|
||||
if fieldName == "" {
|
||||
return nil, ErrFieldNotFound
|
||||
}
|
||||
|
||||
if field := idx.Field(fieldName); field == nil {
|
||||
return nil, ErrFieldNotFound
|
||||
} else if field.keys() {
|
||||
other.Keys = make([]string, len(result))
|
||||
for i, id := range result {
|
||||
key, err := e.TranslateStore.TranslateRowToString(index, fieldName, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
other.Keys[i] = key
|
||||
}
|
||||
} else {
|
||||
other.Rows = result
|
||||
}
|
||||
|
||||
return other, nil
|
||||
}
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
|
|
@ -2136,7 +2205,7 @@ func needsShards(calls []*pql.Call) bool {
|
|||
switch call.Name {
|
||||
case "Clear", "Set", "SetRowAttrs", "SetColumnAttrs":
|
||||
continue
|
||||
case "Count", "TopN", "Rows":
|
||||
case "Count", "TopN", "RowIDs":
|
||||
return true
|
||||
// default catches Bitmap calls
|
||||
default:
|
||||
|
|
|
|||
|
|
@ -1233,12 +1233,12 @@ Set(4500001, fn=4)
|
|||
}); err != nil {
|
||||
t.Fatalf("GroupBy querying: %v", err)
|
||||
} else {
|
||||
expected := pilosa.GroupByCounts{
|
||||
{Groups: []string{"f.10"}, Total: 4},
|
||||
{Groups: []string{"f.7"}, Total: 1},
|
||||
expected := []pilosa.GroupCount{
|
||||
{Group: []pilosa.FieldRow{{Field: "f", RowID: 10}}, Count: 4},
|
||||
{Group: []pilosa.FieldRow{{Field: "f", RowID: 7}}, Count: 1},
|
||||
}
|
||||
results := res.Results[0].(pilosa.GroupByCounts)
|
||||
checkGroupBy(expected, results, t)
|
||||
results := res.Results[0].([]pilosa.GroupCount)
|
||||
checkGroupBy(t, expected, results)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
@ -1639,7 +1639,7 @@ func benchmarkExistence(nn bool, b *testing.B) {
|
|||
func BenchmarkExecutor_Existence_True(b *testing.B) { benchmarkExistence(true, b) }
|
||||
func BenchmarkExecutor_Existence_False(b *testing.B) { benchmarkExistence(false, b) }
|
||||
|
||||
func TestExecutor_Execute_Rows(t *testing.T) {
|
||||
func TestExecutor_Execute_RowIDs(t *testing.T) {
|
||||
c := test.MustRunCluster(t, 1)
|
||||
defer c.Close()
|
||||
hldr := test.Holder{Holder: c[0].Server.Holder()}
|
||||
|
|
@ -1649,24 +1649,28 @@ func TestExecutor_Execute_Rows(t *testing.T) {
|
|||
hldr.SetBit("i", "general", 11, ShardWidth+2)
|
||||
hldr.SetBit("i", "general", 12, 2)
|
||||
hldr.SetBit("i", "general", 12, ShardWidth+2)
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Rows(field=general)`}); err != nil {
|
||||
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `RowIDs(field=general)`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if columns := res.Results[0].(pilosa.RowIDs); !reflect.DeepEqual(columns, pilosa.RowIDs{10, 11, 12}) {
|
||||
} else if columns := res.Results[0].(pilosa.RowIdentifiers); !reflect.DeepEqual(columns, pilosa.RowIdentifiers{Rows: []uint64{10, 11, 12}}) {
|
||||
t.Fatalf("unexpected columns: %+v", columns)
|
||||
}
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Rows(field=general, limit=2)`}); err != nil {
|
||||
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `RowIDs(field=general, limit=2)`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if columns := res.Results[0].(pilosa.RowIDs); !reflect.DeepEqual(columns, pilosa.RowIDs{10, 11}) {
|
||||
} else if columns := res.Results[0].(pilosa.RowIdentifiers); !reflect.DeepEqual(columns, pilosa.RowIdentifiers{Rows: []uint64{10, 11}}) {
|
||||
t.Fatalf("unexpected columns: %+v", columns)
|
||||
}
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Rows(field=general, offset=1,limit=2)`}); err != nil {
|
||||
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `RowIDs(field=general, offset=1,limit=2)`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if columns := res.Results[0].(pilosa.RowIDs); !reflect.DeepEqual(columns, pilosa.RowIDs{11, 12}) {
|
||||
} else if columns := res.Results[0].(pilosa.RowIdentifiers); !reflect.DeepEqual(columns, pilosa.RowIdentifiers{Rows: []uint64{11, 12}}) {
|
||||
t.Fatalf("unexpected columns: %+v", columns)
|
||||
}
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Rows(field=general, column=2)`}); err != nil {
|
||||
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `RowIDs(field=general, column=2)`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if columns := res.Results[0].(pilosa.RowIDs); !reflect.DeepEqual(columns, pilosa.RowIDs{11, 12}) {
|
||||
} else if columns := res.Results[0].(pilosa.RowIdentifiers); !reflect.DeepEqual(columns, pilosa.RowIdentifiers{Rows: []uint64{11, 12}}) {
|
||||
t.Fatalf("unexpected columns: %+v", columns)
|
||||
}
|
||||
}
|
||||
|
|
@ -1681,17 +1685,15 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
|
|||
hldr.SetBit("i", "general", 11, ShardWidth+2)
|
||||
hldr.SetBit("i", "general", 12, 2)
|
||||
hldr.SetBit("i", "general", 12, ShardWidth+2)
|
||||
hldr.SetBit("i", "sub", 10, 0)
|
||||
hldr.SetBit("i", "sub", 10, 1)
|
||||
hldr.SetBit("i", "sub", 10, 3)
|
||||
hldr.SetBit("i", "sub", 11, 2)
|
||||
hldr.SetBit("i", "sub", 11, 0)
|
||||
expected := pilosa.GroupByCounts{
|
||||
{Groups: []string{"general.10", "sub.11"}, Total: 1},
|
||||
{Groups: []string{"general.11", "sub.11"}, Total: 1},
|
||||
{Groups: []string{"general.12", "sub.11"}, Total: 1},
|
||||
{Groups: []string{"general.10", "sub.10"}, Total: 2},
|
||||
}
|
||||
|
||||
hldr.SetBit("i", "sub", 100, 0)
|
||||
hldr.SetBit("i", "sub", 100, 1)
|
||||
hldr.SetBit("i", "sub", 100, 3)
|
||||
hldr.SetBit("i", "sub", 100, ShardWidth+1)
|
||||
|
||||
hldr.SetBit("i", "sub", 110, 2)
|
||||
hldr.SetBit("i", "sub", 110, 0)
|
||||
|
||||
t.Run("No Field List Arguments", func(t *testing.T) {
|
||||
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy()`}); err != nil {
|
||||
if errors.Cause(err) != pilosa.ErrFieldsArgumentRequired {
|
||||
|
|
@ -1714,42 +1716,54 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
|
|||
}
|
||||
})
|
||||
t.Run("Basic", func(t *testing.T) {
|
||||
expected := []pilosa.GroupCount{
|
||||
{Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 110}}, Count: 1},
|
||||
{Group: []pilosa.FieldRow{{Field: "general", RowID: 11}, {Field: "sub", RowID: 110}}, Count: 1},
|
||||
{Group: []pilosa.FieldRow{{Field: "general", RowID: 12}, {Field: "sub", RowID: 110}}, Count: 1},
|
||||
{Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 100}}, Count: 3},
|
||||
}
|
||||
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy(fields=[general,sub])`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else {
|
||||
results := res.Results[0].(pilosa.GroupByCounts)
|
||||
checkGroupBy(expected, results, t)
|
||||
results := res.Results[0].([]pilosa.GroupCount)
|
||||
checkGroupBy(t, expected, results)
|
||||
}
|
||||
})
|
||||
expected = pilosa.GroupByCounts{
|
||||
{Groups: []string{"general.11"}, Total: 2},
|
||||
{Groups: []string{"general.12"}, Total: 2},
|
||||
}
|
||||
|
||||
t.Run("check field offset no limit", func(t *testing.T) {
|
||||
expected := []pilosa.GroupCount{
|
||||
{Group: []pilosa.FieldRow{{Field: "general", RowID: 11}}, Count: 2},
|
||||
{Group: []pilosa.FieldRow{{Field: "general", RowID: 12}}, Count: 2},
|
||||
}
|
||||
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy(fields=[general:11:])`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else {
|
||||
results := res.Results[0].(pilosa.GroupByCounts)
|
||||
checkGroupBy(expected, results, t)
|
||||
results := res.Results[0].([]pilosa.GroupCount)
|
||||
checkGroupBy(t, expected, results)
|
||||
}
|
||||
})
|
||||
expected = pilosa.GroupByCounts{
|
||||
{Groups: []string{"general.11"}, Total: 2},
|
||||
}
|
||||
|
||||
t.Run("check field offset limit", func(t *testing.T) {
|
||||
expected := []pilosa.GroupCount{
|
||||
{Group: []pilosa.FieldRow{{Field: "general", RowID: 11}}, Count: 2},
|
||||
}
|
||||
|
||||
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy(fields=[general:11:1])`}); err != nil {
|
||||
t.Fatal(err)
|
||||
} else {
|
||||
results := res.Results[0].(pilosa.GroupByCounts)
|
||||
checkGroupBy(expected, results, t)
|
||||
results := res.Results[0].([]pilosa.GroupCount)
|
||||
checkGroupBy(t, expected, results)
|
||||
}
|
||||
})
|
||||
}
|
||||
func checkGroupBy(expected, results pilosa.GroupByCounts, t *testing.T) {
|
||||
notIn := func(item pilosa.GroupLine, expected pilosa.GroupByCounts) bool {
|
||||
|
||||
func checkGroupBy(t *testing.T, expected, results []pilosa.GroupCount) {
|
||||
notIn := func(item pilosa.GroupCount, expected []pilosa.GroupCount) bool {
|
||||
for i := range expected {
|
||||
if item.Total == expected[i].Total {
|
||||
if reflect.DeepEqual(item.Groups, expected[i].Groups) {
|
||||
if item.Count == expected[i].Count {
|
||||
if reflect.DeepEqual(item.Group, expected[i].Group) {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
|
|
|||
File diff suppressed because it is too large
Load diff
|
|
@ -8,15 +8,26 @@ message Row {
|
|||
repeated Attr Attrs = 2;
|
||||
}
|
||||
|
||||
message RowIdentifiers {
|
||||
repeated uint64 Rows = 1;
|
||||
repeated string Keys = 2;
|
||||
//repeated Attr Attrs = 3;
|
||||
}
|
||||
|
||||
message Pair {
|
||||
uint64 ID = 1;
|
||||
string Key = 3;
|
||||
uint64 Count = 2;
|
||||
}
|
||||
|
||||
message GroupLine{
|
||||
repeated string Groups = 1;
|
||||
uint64 Total=2;
|
||||
message FieldRow{
|
||||
string Field = 1;
|
||||
uint64 RowID = 2;
|
||||
}
|
||||
|
||||
message GroupCount{
|
||||
repeated FieldRow Group = 1;
|
||||
uint64 Count = 2;
|
||||
}
|
||||
|
||||
message ValCount {
|
||||
|
|
@ -72,7 +83,8 @@ message QueryResult {
|
|||
bool Changed = 4;
|
||||
ValCount ValCount = 5;
|
||||
repeated uint64 RowIDs = 7;
|
||||
repeated GroupLine GroupByCounts = 8;
|
||||
repeated GroupCount GroupCounts = 8;
|
||||
RowIdentifiers RowIdentifiers = 9;
|
||||
}
|
||||
|
||||
message ImportRequest {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue