mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-12 07:41:02 +00:00
Merge pull request #360 from raskle/topn-multinode-tests
TopN tests across 3 nodes
This commit is contained in:
commit
d56a55a0d3
1 changed files with 138 additions and 0 deletions
138
client_test.go
138
client_test.go
|
|
@ -3,13 +3,151 @@ package pilosa_test
|
|||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"github.com/davecgh/go-spew/spew"
|
||||
"github.com/pilosa/pilosa"
|
||||
"github.com/pilosa/pilosa/internal"
|
||||
"github.com/pilosa/pilosa/pql"
|
||||
)
|
||||
|
||||
func createCluster(c *pilosa.Cluster) ([]*Server, []*Index) {
|
||||
numNodes := len(c.Nodes)
|
||||
idx := make([]*Index, numNodes)
|
||||
server := make([]*Server, numNodes)
|
||||
for i := 0; i < numNodes; i++ {
|
||||
idx[i] = MustOpenIndex()
|
||||
server[i] = NewServer()
|
||||
server[i].Handler.Host = server[i].Host()
|
||||
server[i].Handler.Cluster = c
|
||||
server[i].Handler.Cluster.Nodes[i].Host = server[i].Host()
|
||||
server[i].Handler.Index = idx[i].Index
|
||||
}
|
||||
return server, idx
|
||||
}
|
||||
|
||||
// Test distributed TopN Bitmap count across 3 nodes.
|
||||
func TestClient_MultiNode(t *testing.T) {
|
||||
cluster := NewCluster(3)
|
||||
s, idx := createCluster(cluster)
|
||||
|
||||
for i := 0; i < len(cluster.Nodes); i++ {
|
||||
defer idx[i].Close()
|
||||
defer s[i].Close()
|
||||
}
|
||||
|
||||
s[0].Handler.Executor.ExecuteFn = func(ctx context.Context, db string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) {
|
||||
e := pilosa.NewExecutor()
|
||||
e.Index = idx[0].Index
|
||||
e.Host = cluster.Nodes[0].Host
|
||||
e.Cluster = cluster
|
||||
return e.Execute(ctx, db, query, slices, opt)
|
||||
}
|
||||
s[1].Handler.Executor.ExecuteFn = func(ctx context.Context, db string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) {
|
||||
e := pilosa.NewExecutor()
|
||||
e.Index = idx[1].Index
|
||||
e.Host = cluster.Nodes[1].Host
|
||||
e.Cluster = cluster
|
||||
return e.Execute(ctx, db, query, slices, opt)
|
||||
}
|
||||
s[2].Handler.Executor.ExecuteFn = func(ctx context.Context, db string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) {
|
||||
e := pilosa.NewExecutor()
|
||||
e.Index = idx[2].Index
|
||||
e.Host = cluster.Nodes[2].Host
|
||||
e.Cluster = cluster
|
||||
return e.Execute(ctx, db, query, slices, opt)
|
||||
}
|
||||
|
||||
// Create a dispersed set of bitmaps across 3 nodes such that each individual node and slice width increment would reveal a different TopN.
|
||||
idx[0].MustCreateFragmentIfNotExists("d", "f", 0).MustSetBits(99, 1, 2, 3, 4)
|
||||
idx[0].MustCreateFragmentIfNotExists("d", "f", 0).MustSetBits(100, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
|
||||
idx[0].MustCreateFragmentIfNotExists("d", "f", 0).MustSetBits(98, 1, 2, 3, 4, 5, 6)
|
||||
idx[0].MustCreateFragmentIfNotExists("d", "f", 0).MustSetBits(1, 4)
|
||||
idx[0].MustCreateFragmentIfNotExists("d", "f", 0).MustSetBits(22, 1, 2, 3, 4, 5)
|
||||
|
||||
idx[1].MustCreateFragmentIfNotExists("d", "f", 10).MustSetBits(100, (SliceWidth*10)+10)
|
||||
idx[1].MustCreateFragmentIfNotExists("d", "f", 10).MustSetBits(4, (SliceWidth*10)+10, (SliceWidth*10)+11, (SliceWidth*10)+12)
|
||||
idx[1].MustCreateFragmentIfNotExists("d", "f", 10).MustSetBits(4, (SliceWidth*10)+10, (SliceWidth*10)+11, (SliceWidth*10)+12, (SliceWidth*10)+13, (SliceWidth*10)+14, (SliceWidth*10)+15)
|
||||
idx[1].MustCreateFragmentIfNotExists("d", "f", 10).MustSetBits(2, (SliceWidth*10)+1, (SliceWidth*10)+2, (SliceWidth*10)+3, (SliceWidth*10)+4)
|
||||
idx[1].MustCreateFragmentIfNotExists("d", "f", 10).MustSetBits(3, (SliceWidth*10)+1, (SliceWidth*10)+2, (SliceWidth*10)+3, (SliceWidth*10)+4, (SliceWidth*10)+5)
|
||||
idx[1].MustCreateFragmentIfNotExists("d", "f", 10).MustSetBits(22, (SliceWidth*10)+1, (SliceWidth*10)+2, (SliceWidth*10)+10)
|
||||
|
||||
idx[2].MustCreateFragmentIfNotExists("d", "f", 6).MustSetBits(24, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12, (SliceWidth*6)+13, (SliceWidth*6)+14)
|
||||
idx[2].MustCreateFragmentIfNotExists("d", "f", 6).MustSetBits(20, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12, (SliceWidth*6)+13)
|
||||
idx[2].MustCreateFragmentIfNotExists("d", "f", 6).MustSetBits(21, (SliceWidth*6)+10)
|
||||
idx[2].MustCreateFragmentIfNotExists("d", "f", 6).MustSetBits(100, (SliceWidth*6)+10)
|
||||
idx[2].MustCreateFragmentIfNotExists("d", "f", 6).MustSetBits(99, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12)
|
||||
idx[2].MustCreateFragmentIfNotExists("d", "f", 6).MustSetBits(98, (SliceWidth*6)+10, (SliceWidth*6)+11)
|
||||
idx[2].MustCreateFragmentIfNotExists("d", "f", 6).MustSetBits(22, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12)
|
||||
|
||||
// Connect to each node to compare results.
|
||||
client := make([]*Client, 3)
|
||||
client[0] = MustNewClient(s[0].Host())
|
||||
client[1] = MustNewClient(s[0].Host())
|
||||
client[2] = MustNewClient(s[0].Host())
|
||||
|
||||
topN := 4
|
||||
q := fmt.Sprintf(`TopN(frame="%s", n=%d)`, "f", topN)
|
||||
|
||||
result, err := client[0].ExecuteQuery(context.Background(), "d", q, true)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Check the results before every node has the correct max slice value.
|
||||
pairs := result.(internal.QueryResponse).Results[0].Pairs
|
||||
for _, pair := range pairs {
|
||||
if pair.Key == 22 && pair.Count != 5 {
|
||||
t.Fatalf("Invalid Cluster wide MaxSlice prevents accurate calculation of %s", pair)
|
||||
}
|
||||
}
|
||||
|
||||
// Set max slice to correct value.
|
||||
idx[0].DB("d").SetRemoteMaxSlice(10)
|
||||
idx[1].DB("d").SetRemoteMaxSlice(10)
|
||||
idx[2].DB("d").SetRemoteMaxSlice(10)
|
||||
|
||||
result, err = client[0].ExecuteQuery(context.Background(), "d", q, true)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Test must return exactly N results.
|
||||
if len(result.(internal.QueryResponse).Results[0].Pairs) != topN {
|
||||
t.Fatalf("unexpected number of TopN results: %s", spew.Sdump(result))
|
||||
}
|
||||
p := []*internal.Pair{
|
||||
{Key: 100, Count: 12},
|
||||
{Key: 22, Count: 11},
|
||||
{Key: 98, Count: 8},
|
||||
{Key: 99, Count: 7}}
|
||||
|
||||
// Valdidate the Top 4 result counts.
|
||||
if !reflect.DeepEqual(result.(internal.QueryResponse).Results[0].Pairs, p) {
|
||||
t.Fatalf("Invalid TopN result set: %s", spew.Sdump(result))
|
||||
}
|
||||
|
||||
result1, err := client[1].ExecuteQuery(context.Background(), "d", q, true)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result2, err := client[2].ExecuteQuery(context.Background(), "d", q, true)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Compare TopN results across all nodes in the cluster.
|
||||
if !reflect.DeepEqual(result, result1) {
|
||||
t.Fatalf("TopN result should be the same on node0 and node1: %s", spew.Sdump(result1))
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(result, result2) {
|
||||
t.Fatalf("TopN result should be the same on node0 and node2: %s", spew.Sdump(result2))
|
||||
}
|
||||
}
|
||||
|
||||
// Ensure client can bulk import data.
|
||||
func TestClient_Import(t *testing.T) {
|
||||
idx := MustOpenIndex()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue