From fb706ab883e87ecf257740bcc3aea26bb2939cb8 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Fri, 12 Oct 2018 18:47:41 -0500 Subject: [PATCH] add GroupBy Rows(limit) test and fix bug run all group by tests on two cluster sizes --- executor_test.go | 372 +++++++++++++++++++++++++++++------------------ fragment.go | 18 ++- 2 files changed, 237 insertions(+), 153 deletions(-) diff --git a/executor_test.go b/executor_test.go index 666eb104d..6b903acc2 100644 --- a/executor_test.go +++ b/executor_test.go @@ -1358,8 +1358,8 @@ Set(4500001, fn=4) t.Fatalf("GroupBy querying: %v", err) } else { expected := []pilosa.GroupCount{ - {Group: []pilosa.FieldRow{{Field: "f", RowID: 10}}, Count: 4}, {Group: []pilosa.FieldRow{{Field: "f", RowID: 7}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "f", RowID: 10}}, Count: 4}, } results := res.Results[0].([]pilosa.GroupCount) checkGroupBy(t, expected, results) @@ -2346,174 +2346,256 @@ func TestExecutor_Execute_Rows_Keys(t *testing.T) { } func TestExecutor_Execute_GroupBy(t *testing.T) { - c := test.MustRunCluster(t, 1) - defer c.Close() - c.CreateField(t, "i", pilosa.IndexOptions{}, "general") - c.CreateField(t, "i", pilosa.IndexOptions{}, "sub") - c.ImportBits(t, "i", "general", [][2]uint64{ - {10, 0}, - {10, 1}, - {10, ShardWidth + 1}, - {11, 2}, - {11, ShardWidth + 2}, - {12, 2}, - {12, ShardWidth + 2}, - }) - c.ImportBits(t, "i", "sub", [][2]uint64{ - {100, 0}, - {100, 1}, - {100, 3}, - {100, ShardWidth + 1}, + groupByTest := func(t *testing.T, clusterSize int) { + c := test.MustRunCluster(t, 1) + defer c.Close() + c.CreateField(t, "i", pilosa.IndexOptions{}, "general") + c.CreateField(t, "i", pilosa.IndexOptions{}, "sub") + c.ImportBits(t, "i", "general", [][2]uint64{ + {10, 0}, + {10, 1}, + {10, ShardWidth + 1}, + {11, 2}, + {11, ShardWidth + 2}, + {12, 2}, + {12, ShardWidth + 2}, + }) + c.ImportBits(t, "i", "sub", [][2]uint64{ + {100, 0}, + {100, 1}, + {100, 3}, + {100, ShardWidth + 1}, - {110, 2}, - {110, 0}, - }) + {110, 2}, + {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 !strings.Contains(err.Error(), "need at least one child call") { - t.Fatalf("unexpected error: \"%v\"", err) + 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 !strings.Contains(err.Error(), "need at least one child call") { + t.Fatalf("unexpected error: \"%v\"", err) + } } - } - }) + }) - t.Run("Unknown Field ", func(t *testing.T) { - if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy(Rows(field=missing))`}); err != nil { - if errors.Cause(err) != pilosa.ErrFieldNotFound { - t.Fatalf("unexpected error\n\"%s\" not returned instead \n\"%s\"", pilosa.ErrFieldNotFound, err) + t.Run("Unknown Field ", func(t *testing.T) { + if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `GroupBy(Rows(field=missing))`}); err != nil { + if errors.Cause(err) != pilosa.ErrFieldNotFound { + t.Fatalf("unexpected error\n\"%s\" not returned instead \n\"%s\"", pilosa.ErrFieldNotFound, err) + } } - } - }) + }) - t.Run("Basic", func(t *testing.T) { - expected := []pilosa.GroupCount{ - {Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 100}}, Count: 3}, - {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}, - } + t.Run("Basic", func(t *testing.T) { + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "sub", RowID: 100}}, Count: 3}, + {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}, + } - results := c.Query(t, "i", `GroupBy(Rows(field=general), Rows(field=sub))`).Results[0].([]pilosa.GroupCount) - checkGroupBy(t, expected, results) - }) + results := c.Query(t, "i", `GroupBy(Rows(field=general), Rows(field=sub))`).Results[0].([]pilosa.GroupCount) + checkGroupBy(t, expected, results) + }) - 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}, - } + 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}, + } - results := c.Query(t, "i", `GroupBy(Rows(field=general, previous=10))`).Results[0].([]pilosa.GroupCount) - checkGroupBy(t, expected, results) - }) + results := c.Query(t, "i", `GroupBy(Rows(field=general, previous=10))`).Results[0].([]pilosa.GroupCount) + checkGroupBy(t, expected, results) + }) - t.Run("check field offset limit", func(t *testing.T) { - expected := []pilosa.GroupCount{ - {Group: []pilosa.FieldRow{{Field: "general", RowID: 11}}, Count: 2}, - } + t.Run("check field offset limit", func(t *testing.T) { + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "general", RowID: 11}}, Count: 2}, + } - results := c.Query(t, "i", `GroupBy(Rows(field=general, previous=10), limit=1)`).Results[0].([]pilosa.GroupCount) - checkGroupBy(t, expected, results) + results := c.Query(t, "i", `GroupBy(Rows(field=general, previous=10), limit=1)`).Results[0].([]pilosa.GroupCount) + checkGroupBy(t, expected, results) - }) + }) - c.CreateField(t, "i", pilosa.IndexOptions{}, "a") - c.CreateField(t, "i", pilosa.IndexOptions{}, "b") - c.ImportBits(t, "i", "a", [][2]uint64{ - {0, 1}, - {1, ShardWidth + 1}, - }) - c.ImportBits(t, "i", "b", [][2]uint64{ - {0, ShardWidth + 1}, - {1, 1}, - }) + c.CreateField(t, "i", pilosa.IndexOptions{}, "a") + c.CreateField(t, "i", pilosa.IndexOptions{}, "b") + c.ImportBits(t, "i", "a", [][2]uint64{ + {0, 1}, + {1, ShardWidth + 1}, + }) + c.ImportBits(t, "i", "b", [][2]uint64{ + {0, ShardWidth + 1}, + {1, 1}, + }) - t.Run("tricky data", func(t *testing.T) { - expected := []pilosa.GroupCount{ - {Group: []pilosa.FieldRow{{Field: "a", RowID: 0}, {Field: "b", RowID: 1}}, Count: 1}, - } + t.Run("tricky data", func(t *testing.T) { + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "a", RowID: 0}, {Field: "b", RowID: 1}}, Count: 1}, + } - results := c.Query(t, "i", `GroupBy(Rows(field=a), Rows(field=b), limit=1)`).Results[0].([]pilosa.GroupCount) - checkGroupBy(t, expected, results) - }) + results := c.Query(t, "i", `GroupBy(Rows(field=a), Rows(field=b), limit=1)`).Results[0].([]pilosa.GroupCount) + checkGroupBy(t, expected, results) + }) - // set the same bits in a single shard in three fields - c.CreateField(t, "i", pilosa.IndexOptions{}, "wa") - c.CreateField(t, "i", pilosa.IndexOptions{}, "wb") - c.CreateField(t, "i", pilosa.IndexOptions{}, "wc") - c.ImportBits(t, "i", "wa", [][2]uint64{ - {0, 0}, {0, 1}, {0, 2}, // all - {1, 1}, // odds - {2, 0}, {2, 2}, // evens - {3, 3}, // no overlap - }) - c.ImportBits(t, "i", "wb", [][2]uint64{ - {0, 0}, {0, 1}, {0, 2}, - {1, 1}, - {2, 0}, {2, 2}, - {3, 3}, - }) - c.ImportBits(t, "i", "wc", [][2]uint64{ - {0, 0}, {0, 1}, {0, 2}, - {1, 1}, - {2, 0}, {2, 2}, - {3, 3}, - }) + // set the same bits in a single shard in three fields + c.CreateField(t, "i", pilosa.IndexOptions{}, "wa") + c.CreateField(t, "i", pilosa.IndexOptions{}, "wb") + c.CreateField(t, "i", pilosa.IndexOptions{}, "wc") + c.ImportBits(t, "i", "wa", [][2]uint64{ + {0, 0}, {0, 1}, {0, 2}, // all + {1, 1}, // odds + {2, 0}, {2, 2}, // evens + {3, 3}, // no overlap + }) + c.ImportBits(t, "i", "wb", [][2]uint64{ + {0, 0}, {0, 1}, {0, 2}, + {1, 1}, + {2, 0}, {2, 2}, + {3, 3}, + }) + c.ImportBits(t, "i", "wc", [][2]uint64{ + {0, 0}, {0, 1}, {0, 2}, + {1, 1}, + {2, 0}, {2, 2}, + {3, 3}, + }) - t.Run("test wrapping with previous", func(t *testing.T) { - results := c.Query(t, "i", `GroupBy(Rows(field=wa), Rows(field=wb), Rows(field=wc, previous=1), limit=3)`).Results[0].([]pilosa.GroupCount) - expected := []pilosa.GroupCount{ - {Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 0}, {Field: "wc", RowID: 2}}, Count: 2}, - {Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 1}, {Field: "wc", RowID: 0}}, Count: 1}, - {Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 1}, {Field: "wc", RowID: 1}}, Count: 1}, - } - checkGroupBy(t, expected, results) - }) + t.Run("test wrapping with previous", func(t *testing.T) { + results := c.Query(t, "i", `GroupBy(Rows(field=wa), Rows(field=wb), Rows(field=wc, previous=1), limit=3)`).Results[0].([]pilosa.GroupCount) + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 0}, {Field: "wc", RowID: 2}}, Count: 2}, + {Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 1}, {Field: "wc", RowID: 0}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "wa", RowID: 0}, {Field: "wb", RowID: 1}, {Field: "wc", RowID: 1}}, Count: 1}, + } + checkGroupBy(t, expected, results) + }) - t.Run("test wrapping multiple", func(t *testing.T) { - results := c.Query(t, "i", `GroupBy(Rows(field=wa), Rows(field=wb, previous=2), Rows(field=wc, previous=2), limit=1)`).Results[0].([]pilosa.GroupCount) - expected := []pilosa.GroupCount{ - {Group: []pilosa.FieldRow{{Field: "wa", RowID: 1}, {Field: "wb", RowID: 0}, {Field: "wc", RowID: 0}}, Count: 1}, - } - checkGroupBy(t, expected, results) - }) + t.Run("test wrapping multiple", func(t *testing.T) { + results := c.Query(t, "i", `GroupBy(Rows(field=wa), Rows(field=wb, previous=2), Rows(field=wc, previous=2), limit=1)`).Results[0].([]pilosa.GroupCount) + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "wa", RowID: 1}, {Field: "wb", RowID: 0}, {Field: "wc", RowID: 0}}, Count: 1}, + } + checkGroupBy(t, expected, results) + }) - // TODO test multiple shards with distinct results (different rows) and same - // rows to ensure ordering, limit behavior and correctness - c.CreateField(t, "i", pilosa.IndexOptions{}, "ma") - c.CreateField(t, "i", pilosa.IndexOptions{}, "mb") - c.CreateField(t, "i", pilosa.IndexOptions{}, "mc") - c.ImportBits(t, "i", "ma", [][2]uint64{ - {0, 0}, - {1, ShardWidth}, - {2, 0}, - {3, ShardWidth}, - }) - c.ImportBits(t, "i", "mb", [][2]uint64{ - {0, 0}, - {1, ShardWidth}, - {2, 0}, - {3, ShardWidth}, - }) - t.Run("distinct rows in different shards", func(t *testing.T) { - results := c.Query(t, "i", `GroupBy(Rows(field=ma), Rows(field=mb), limit=5)`).Results[0].([]pilosa.GroupCount) - expected := []pilosa.GroupCount{ - {Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 0}}, Count: 1}, - {Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 2}}, Count: 1}, - {Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 1}}, Count: 1}, - {Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 3}}, Count: 1}, - {Group: []pilosa.FieldRow{{Field: "ma", RowID: 2}, {Field: "mb", RowID: 0}}, Count: 1}, - } - checkGroupBy(t, expected, results) + // test multiple shards with distinct results (different rows) and same + // rows to ensure ordering, limit behavior and correctness + c.CreateField(t, "i", pilosa.IndexOptions{}, "ma") + c.CreateField(t, "i", pilosa.IndexOptions{}, "mb") + c.ImportBits(t, "i", "ma", [][2]uint64{ + {0, 0}, + {1, ShardWidth}, + {2, 0}, + {3, ShardWidth}, + }) + c.ImportBits(t, "i", "mb", [][2]uint64{ + {0, 0}, + {1, ShardWidth}, + {2, 0}, + {3, ShardWidth}, + }) + t.Run("distinct rows in different shards", func(t *testing.T) { + results := c.Query(t, "i", `GroupBy(Rows(field=ma), Rows(field=mb), limit=5)`).Results[0].([]pilosa.GroupCount) + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 0}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 2}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 1}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 3}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 2}, {Field: "mb", RowID: 0}}, Count: 1}, + } + checkGroupBy(t, expected, results) + }) - }) + t.Run("distinct rows in different shards with row limit", func(t *testing.T) { + results := c.Query(t, "i", `GroupBy(Rows(field=ma), Rows(field=mb, limit=2), limit=5)`).Results[0].([]pilosa.GroupCount) + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 0}, {Field: "mb", RowID: 0}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 1}, {Field: "mb", RowID: 1}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 2}, {Field: "mb", RowID: 0}}, Count: 1}, + {Group: []pilosa.FieldRow{{Field: "ma", RowID: 3}, {Field: "mb", RowID: 1}}, Count: 1}, + } + checkGroupBy(t, expected, results) + }) - // TODO test column queries to row call (also with multiple shards) + c.CreateField(t, "i", pilosa.IndexOptions{}, "na") + c.CreateField(t, "i", pilosa.IndexOptions{}, "nb") + c.ImportBits(t, "i", "na", [][2]uint64{ + {0, 0}, + {0, ShardWidth}, + {1, 0}, + {1, ShardWidth}, + }) + c.ImportBits(t, "i", "nb", [][2]uint64{ + {0, 0}, + {0, ShardWidth}, + {1, 0}, + {1, ShardWidth}, + }) + t.Run("same rows in different shards", func(t *testing.T) { + results := c.Query(t, "i", `GroupBy(Rows(field=na), Rows(field=nb))`).Results[0].([]pilosa.GroupCount) + expected := []pilosa.GroupCount{ + {Group: []pilosa.FieldRow{{Field: "na", RowID: 0}, {Field: "nb", RowID: 0}}, Count: 2}, + {Group: []pilosa.FieldRow{{Field: "na", RowID: 0}, {Field: "nb", RowID: 1}}, Count: 2}, + {Group: []pilosa.FieldRow{{Field: "na", RowID: 1}, {Field: "nb", RowID: 0}}, Count: 2}, + {Group: []pilosa.FieldRow{{Field: "na", RowID: 1}, {Field: "nb", RowID: 1}}, Count: 2}, + } + checkGroupBy(t, expected, results) - // TODO test limit query to Rows calls + }) - // TODO test paging over results using previous. + // TODO test column queries to row call (also with multiple shards) + // test paging over results using previous. set the same bits in three + // fields + c.CreateField(t, "i", pilosa.IndexOptions{}, "ppa") + c.CreateField(t, "i", pilosa.IndexOptions{}, "ppb") + c.CreateField(t, "i", pilosa.IndexOptions{}, "ppc") + c.ImportBits(t, "i", "ppa", [][2]uint64{ + {0, 0}, + {1, 0}, + {2, 0}, + {3, 0}, {3, 91000}, {3, ShardWidth}, {3, ShardWidth * 2}, {3, ShardWidth * 3}, + }) + c.ImportBits(t, "i", "ppb", [][2]uint64{ + {0, 0}, + {1, 0}, + {2, 0}, + {3, 0}, {3, 91000}, {3, ShardWidth}, {3, ShardWidth * 2}, {3, ShardWidth * 3}, + }) + c.ImportBits(t, "i", "ppc", [][2]uint64{ + {0, 0}, + {1, 0}, + {2, 0}, + {3, 0}, {3, 91000}, {3, ShardWidth}, {3, ShardWidth * 2}, {3, ShardWidth * 3}, + }) + + t.Run("test wrapping with previous", func(t *testing.T) { + totalResults := make([]pilosa.GroupCount, 0) + var results []pilosa.GroupCount + results = c.Query(t, "i", `GroupBy(Rows(field=ppa), Rows(field=ppb), Rows(field=ppc), limit=3)`).Results[0].([]pilosa.GroupCount) + totalResults = append(totalResults, results...) + for len(totalResults) < 64 { + lastGroup := results[len(results)-1].Group + query := fmt.Sprintf("GroupBy(Rows(field=ppa, previous=%d), Rows(field=ppb, previous=%d), Rows(field=ppc, previous=%d), limit=3)", lastGroup[0].RowID, lastGroup[1].RowID, lastGroup[2].RowID) + results = c.Query(t, "i", query).Results[0].([]pilosa.GroupCount) + totalResults = append(totalResults, results...) + } + + expected := make([]pilosa.GroupCount, 64) + for i := 0; i < 64; i++ { + expected[i] = pilosa.GroupCount{Group: []pilosa.FieldRow{{Field: "ppa", RowID: uint64(i / 16)}, {Field: "ppb", RowID: uint64((i % 16) / 4)}, {Field: "ppc", RowID: uint64(i % 4)}}, Count: 1} + } + expected[63].Count = 5 + + checkGroupBy(t, expected, totalResults) + }) + } + for size := range []int{1, 3} { + t.Run(fmt.Sprintf("%d_nodes", size), func(t *testing.T) { + groupByTest(t, size) + }) + } } func BenchmarkGroupBy(b *testing.B) { diff --git a/fragment.go b/fragment.go index cf38972ab..f70ddf482 100644 --- a/fragment.go +++ b/fragment.go @@ -2022,8 +2022,9 @@ func filterWithRows(rows []uint64) rowFilter { // included, there must be one container in that row where all filters return // true. For a row to be skipped, at least one filter must return false for each // container in that row (it need not be the same filter for each). Any filter -// returning done == true will cause processing to stop and the rows accumulated -// so far will be returned. +// returning done == true will cause processing to stop after all filters for +// this container have been processed. The rows accumulated up to this point +// (including this row if all filters passed) will be returned. func (f *fragment) rows(start uint64, filters ...rowFilter) []uint64 { startKey := rowToKey(start) i, _ := f.storage.Containers.Iterator(startKey) @@ -2043,13 +2044,11 @@ func (f *fragment) rows(start uint64, filters ...rowFilter) []uint64 { } // apply filters - addRow := true + addRow, done := true, false for _, filter := range filters { - var done bool - addRow, done = filter(vRow, key, c) - if done { - return rows - } + var d bool + addRow, d = filter(vRow, key, c) + done = done || d if !addRow { break } @@ -2058,6 +2057,9 @@ func (f *fragment) rows(start uint64, filters ...rowFilter) []uint64 { lastRow = vRow rows = append(rows, vRow) } + if done { + return rows + } } return rows }