diff --git a/README.md b/README.md index 02b2e2216..3b7827922 100644 --- a/README.md +++ b/README.md @@ -239,8 +239,7 @@ bspawn uses a json config format that has 5 top level items - an example is belo { "CreatorArgs": ["-type", "local", "-serverN", "1", "-replicaN", "1"], "PilosaHosts": ["localhost:19327"], - "Agents": { "Type": "local" }, - "AgentHosts": ["localhost"], + "AgentHosts": ["agent.example.com"], "Benchmarks": [ { "Num": 1, @@ -261,11 +260,8 @@ Specifies the pilosa cluster that should be created to run benchmarks against. F #### PilosaHosts If PilosaHosts is set, CreatorArgs will be ignored, and an existing pilosa cluster specified by the list of hosts will be used. -#### Agents -Agents specifies the host(s) that the benchmark should be run from. Currently only running from localhost is supported. - #### AgentHosts -If AgentHosts is specified, Agents is ignored, and the existing agents specified here are used. This is not yet implemented. +If AgentHosts is not empty, the agents specified here are used; if it is empty, agents will be run locally. #### Benchmarks Benchmarks is where the actual benchmarks to run are specified - each contains a `Num` which is the number of agents that should run that benchmark, and Args which specifies the benchmark. The benchmarks in the `Benchmarks` list will be run concurrently. For more information about Args, see the `pilosactl bagent -help`. diff --git a/bench/bench.go b/bench/bench.go index 7949ae062..0ffcc0a40 100644 --- a/bench/bench.go +++ b/bench/bench.go @@ -1,12 +1,6 @@ package bench -import ( - "context" - "fmt" - "strconv" - "sync" - "time" -) +import "context" // Benchmark is an interface to guide the creation of new pilosa benchmarks or // benchmark components. It defines 2 methods, Init, and Run. These are separate @@ -43,97 +37,3 @@ type Command interface { func agentizeNum(n, iterations, agentNum int) int { return n + (agentNum * iterations) } - -type parallelBenchmark struct { - benchmarkers []Benchmark -} - -// Init calls Init for each benchmark. If there are any errors, it will return a -// non-nil error value. -func (pb *parallelBenchmark) Init(hosts []string, agentNum int) error { - var g ErrGroup - for i, _ := range pb.benchmarkers { - b := pb.benchmarkers[i] - g.Go(func() error { - return b.Init(hosts, agentNum) - }) - } - return g.Wait() -} - -// Run runs the parallel benchmark and returns it's results in a nested map - the -// top level keys are the indices of each benchmark in the list of benchmarks, -// and the values are the results of each benchmark's Run method. -func (pb *parallelBenchmark) Run(ctx context.Context, agentNum int) map[string]interface{} { - wg := sync.WaitGroup{} - results := make(map[string]interface{}, len(pb.benchmarkers)) - resultsLock := sync.Mutex{} - for i, b := range pb.benchmarkers { - wg.Add(1) - go func(i int, b Benchmark) { - defer wg.Done() - ret := b.Run(ctx, agentNum) - resultsLock.Lock() - results[strconv.Itoa(i)] = ret - resultsLock.Unlock() - }(i, b) - } - wg.Wait() - return results -} - -// Parallel takes a variable number of Benchmarks and returns a Benchmark -// which combines them and will run them in parallel. -func Parallel(bs ...Benchmark) Benchmark { - return ¶llelBenchmark{ - benchmarkers: bs, - } -} - -type serialBenchmark struct { - benchmarkers []Benchmark -} - -// Init calls Init for each benchmark. If there are any errors, it will return a -// non-nil error value. -func (sb *serialBenchmark) Init(hosts []string, agentNum int) error { - errors := make([]error, len(sb.benchmarkers)) - hadErr := false - for i, b := range sb.benchmarkers { - errors[i] = b.Init(hosts, agentNum) - if errors[i] != nil { - hadErr = true - } - } - if hadErr { - return fmt.Errorf("Had errs in serialBenchmark.Init: %v", errors) - } - return nil -} - -// Run runs the serial benchmark and returns it's results in a nested map - the -// top level keys are the indices of each benchmark in the list of benchmarks, -// and the values are the results of each benchmark's Run method. -func (sb *serialBenchmark) Run(ctx context.Context, agentNum int) map[string]interface{} { - results := make(map[string]interface{}, len(sb.benchmarkers)) - runtimes := make(map[string]interface{}) - total_start := time.Now() - for i, b := range sb.benchmarkers { - start := time.Now() - ret := b.Run(ctx, agentNum) - end := time.Now() - results[strconv.Itoa(i)] = ret - runtimes[strconv.Itoa(i)] = end.Sub(start) - } - runtimes["total"] = time.Now().Sub(total_start) - results["runtimes"] = runtimes - return results -} - -// Serial takes a variable number of Benchmarks and returns a Benchmark -// which combines then and will run each serially. -func Serial(bs ...Benchmark) Benchmark { - return &serialBenchmark{ - benchmarkers: bs, - } -} diff --git a/bench/import.go b/bench/import.go index 1b564533b..2854ea458 100644 --- a/bench/import.go +++ b/bench/import.go @@ -7,7 +7,6 @@ import ( "io" "io/ioutil" "math/rand" - "time" "sort" @@ -48,7 +47,7 @@ The following arguments are available: -base-bitmap-id int bits being set will all be greater than this - -maximum-bitmap-id int + -max-bitmap-id int bits being set will all be less than this -base-profile-id int @@ -135,9 +134,8 @@ func (b *Import) Init(hosts []string, agentNum int) error { b.MinBitsPerMap, b.MaxBitsPerMap, b.Seed+int64(agentNum), b.RandomBitmapOrder) b.numbits = num // set b.Paths - f.Close() b.Paths = []string{f.Name()} - return nil + return f.Close() } // Run runs the Import benchmark @@ -145,13 +143,11 @@ func (b *Import) Run(ctx context.Context, agentNum int) map[string]interface{} { results := make(map[string]interface{}) results["numbits"] = b.numbits results["db"] = b.Database - start := time.Now() err := b.ImportCommand.Run(ctx) if err != nil { results["error"] = err.Error() } - results["time"] = time.Now().Sub(start) return results } diff --git a/bench/zipf.go b/bench/zipf.go index 84759dc0a..41ada1293 100644 --- a/bench/zipf.go +++ b/bench/zipf.go @@ -153,7 +153,11 @@ func (b *ZipfSetBits) Run(ctx context.Context, agentNum int) map[string]interfac query := fmt.Sprintf("SetBit(%d, 'frame.n', %d)", b.BaseBitmapID+int64(bitmapID), b.BaseProfileID+int64(profID)) start = time.Now() - b.cli.ExecuteQuery(ctx, b.DB, query, true) + _, err := b.cli.ExecuteQuery(ctx, b.DB, query, true) + if err != nil { + results["error"] = fmt.Sprintf("Error executing query in zipf: %v", err) + return results + } s.Add(time.Now().Sub(start)) } AddToResults(s, results) diff --git a/cmd/pilosactl/import.json b/cmd/pilosactl/import.json index 0f5ca0d8a..920353dbc 100644 --- a/cmd/pilosactl/import.json +++ b/cmd/pilosactl/import.json @@ -1,9 +1,10 @@ { "CreatorArgs": ["-type", "local", "-serverN", "1", "-replicaN", "1"], - "Agents": { "Type": "local" }, + "AgentHosts": [], "Benchmarks": [ { "Num": 1, + "Name": "import", "Args": ["import", "-max-bitmap-id", "100000", "-max-profile-id", "10000", "-max-bits-per-map", "100", "-seed", "0", "-agent-controls", "width", "import", "-max-bitmap-id", "100000", "-max-profile-id", "10000", "-max-bits-per-map", "100", "-seed", "0", "-agent-controls", "width", "-random-bitmap-order", "-db", "randoload"] } ] diff --git a/cmd/pilosactl/main.go b/cmd/pilosactl/main.go index d5b36153d..7bebc8941 100644 --- a/cmd/pilosactl/main.go +++ b/cmd/pilosactl/main.go @@ -25,8 +25,6 @@ import ( "time" "unsafe" - "golang.org/x/crypto/ssh" - "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/bench" "github.com/pilosa/pilosa/creator" @@ -34,6 +32,7 @@ import ( "github.com/pilosa/pilosa/roaring" "github.com/satori/go.uuid" + "golang.org/x/crypto/ssh" ) var ( @@ -991,7 +990,6 @@ type CreateCommand struct { LogFilePrefix string Hosts []string GoMaxProcs int - RunUUID string SSHUser string @@ -1022,7 +1020,6 @@ func (cmd *CreateCommand) ParseFlags(args []string) error { fs.IntVar(&cmd.GoMaxProcs, "gomaxprocs", 0, "") fs.StringVar(&hosts, "hosts", "", "") fs.StringVar(&cmd.SSHUser, "ssh-user", "", "") - fs.StringVar(&cmd.RunUUID, "run-uuid", "", "") if err := fs.Parse(args); err != nil { return err @@ -1072,7 +1069,6 @@ The following flags are allowed: type CreateOutput struct { Hosts []string `json:"hosts"` LogFiles []string `json:"log-files"` - RunUUID string `json:"run-uuid"` } // Run executes cluster creation. @@ -1119,8 +1115,7 @@ func (cmd *CreateCommand) Run(ctx context.Context) error { defer clus.Shutdown() output := &CreateOutput{ - RunUUID: cmd.RunUUID, - Hosts: clus.Hosts(), + Hosts: clus.Hosts(), } logReaders := clus.Logs() @@ -1168,8 +1163,6 @@ type BagentCommand struct { // Enable pretty printing of results, for human consumption. HumanReadable bool `json:"human-readable"` - RunUUID string `json:"-"` // ignoring this here because we add it to the top level of the output - // Slice of pilosa hosts to run the Benchmarks against. Hosts []string `json:"hosts"` @@ -1193,7 +1186,7 @@ func NewBagentCommand(stdin io.Reader, stdout, stderr io.Writer) *BagentCommand } // ParseFlags parses command line flags for the BagentCommand. First the command -// wide flags `hosts` and `agentNum` are parsed. The rest of the flags should be +// wide flags `hosts` and `agent-num` are parsed. The rest of the flags should be // a series of subcommands along with their flags. ParseFlags runs each // subcommand's `ConsumeFlags` method which parses the flags for that command // and returns the rest of the argument slice which should contain further @@ -1204,9 +1197,8 @@ func (cmd *BagentCommand) ParseFlags(args []string) error { var pilosaHosts string fs.StringVar(&pilosaHosts, "hosts", "localhost:15000", "") - fs.IntVar(&cmd.AgentNum, "agentNum", 0, "") + fs.IntVar(&cmd.AgentNum, "agent-num", 0, "") fs.BoolVar(&cmd.HumanReadable, "human", false, "") - fs.StringVar(&cmd.RunUUID, "run-uuid", "", "") if err := fs.Parse(args); err != nil { return err @@ -1267,7 +1259,7 @@ The following arguments are available: -hosts Comma separated list of host:port describing all hosts in the cluster. - -agentNum N + -agent-num N An integer differentiating this agent from others in the fleet. -human @@ -1286,17 +1278,14 @@ The following arguments are available: // Run executes the benchmark agent. func (cmd *BagentCommand) Run(ctx context.Context) error { - sbm := bench.Serial(cmd.Benchmarks...) + sbm := serial(cmd.Benchmarks...) err := sbm.Init(cmd.Hosts, cmd.AgentNum) if err != nil { return fmt.Errorf("in cmd.Run initialization: %v", err) } res := sbm.Run(ctx, cmd.AgentNum) - res["metadata"] = cmd - if cmd.RunUUID != "" { - res["run-uuid"] = cmd.RunUUID - } + res["agent-num"] = cmd.AgentNum enc := json.NewEncoder(cmd.Stdout) if cmd.HumanReadable { enc.SetIndent("", " ") @@ -1325,12 +1314,16 @@ type BspawnCommand struct { // should include everything that comes after `pilosactl create` CreatorArgs []string - // If AgentHosts is specified, Agents is ignored, and the existing - // agents specified here are used. + // List of hosts to run agents on. If this is empty, agents will be run + // locally. AgentHosts []string + // If this is true, build and copy pilosactl binary to agent hosts. CopyBinary bool + // Makes output human readable + Human bool + // Benchmarks is a slice of Spawns which specifies all of the bagent // commands to run. These will all be run in parallel, started on each // of the agents in a round robin fashion. @@ -1346,7 +1339,8 @@ type BspawnCommand struct { // Spawn represents a bagent command run in parallel across Num agents. The // bagent command can run multiple Benchmarks serially within itself. type Spawn struct { - Num int // number of agents to run + Num int `json:"num"` // number of agents to run + Name string `json:"name"` // Should describe what this Spawn does Args []string // everything that comes after `pilosactl bagent [arguments]` } @@ -1390,11 +1384,13 @@ pilosactl spawn configfile // Run executes the main program execution. func (cmd *BspawnCommand) Run(ctx context.Context) error { runUUID := uuid.NewV1() + output := make(map[string]interface{}) + output["run-uuid"] = runUUID.String() if len(cmd.PilosaHosts) == 0 { // must create cluster r, w := io.Pipe() createCmd := NewCreateCommand(cmd.Stdin, w, cmd.Stderr) - err := createCmd.ParseFlags(append(cmd.CreatorArgs, []string{"-run-uuid", runUUID.String()}...)) + err := createCmd.ParseFlags(cmd.CreatorArgs) if err != nil { return err } @@ -1411,20 +1407,23 @@ func (cmd *BspawnCommand) Run(ctx context.Context) error { return err } cmd.PilosaHosts = clus.Hosts - - // print createOutput to stdout - enc := json.NewEncoder(cmd.Stdout) - enc.SetIndent("", " ") - err = enc.Encode(clus) - if err != nil { - return err - } + output["cluster"] = clus } - if len(cmd.AgentHosts) > 0 { - return cmd.spawnRemote(ctx, runUUID) - } else { - return cmd.spawnLocal(ctx, runUUID) + if len(cmd.AgentHosts) == 0 { + cmd.AgentHosts = []string{"localhost"} } + output["agents"] = cmd.AgentHosts + res, err := cmd.spawnRemote(ctx) + if err != nil { + return err + } + output["results"] = res + enc := json.NewEncoder(cmd.Stdout) + if cmd.Human { + enc.SetIndent("", " ") + output = bench.Prettify(output) + } + return enc.Encode(output) } func copyBinary(fleet pilosactl.SSHFleet, pkg, goos, goarch string) error { @@ -1445,33 +1444,53 @@ func copyBinary(fleet pilosactl.SSHFleet, pkg, goos, goarch string) error { return wc.Close() } -func (cmd *BspawnCommand) spawnRemote(ctx context.Context, runUUID uuid.UUID) error { +func (cmd *BspawnCommand) spawnRemote(ctx context.Context) (map[string]interface{}, error) { agentIndex := 0 agentConnections, err := pilosactl.SSHClients(cmd.AgentHosts, cmd.SSHUser, "", cmd.Stderr) if err != nil { - return err + return nil, err } if cmd.CopyBinary { err = copyBinary(agentConnections, "github.com/pilosa/pilosa/cmd/pilosactl", "linux", "amd64") if err != nil { - return err + return nil, err } } sessions := make([]*ssh.Session, 0) + results := make(map[string]interface{}) + resLock := sync.Mutex{} + wg := sync.WaitGroup{} for _, sp := range cmd.Benchmarks { + results[sp.Name] = make(map[int]interface{}) for i := 0; i < sp.Num; i++ { sess, err := agentConnections[agentIndex].NewSession() if err != nil { - return err + return nil, err } sessions = append(sessions, sess) - sess.Stdout = cmd.Stdout - sess.Stderr = cmd.Stderr - err = sess.Start("PATH=.:$PATH pilosactl bagent -agentNum=" + strconv.Itoa(i) + " -hosts=" + strings.Join(cmd.PilosaHosts, ",") + " -run-uuid=" + runUUID.String() + " " + strings.Join(sp.Args, " ")) + stdout, err := sess.StdoutPipe() if err != nil { - return err + return nil, err + } + wg.Add(1) + go func(stdout io.Reader, name string, num int) { + defer wg.Done() + dec := json.NewDecoder(stdout) + var v interface{} + err := dec.Decode(&v) + if err != nil { + fmt.Fprintf(cmd.Stderr, "error decoding json: %v, spawn: %v", err, name) + } + resLock.Lock() + results[name].(map[int]interface{})[num] = v + resLock.Unlock() + }(stdout, sp.Name, i) + sess.Stderr = cmd.Stderr + err = sess.Start("PATH=.:$PATH pilosactl bagent -agent-num=" + strconv.Itoa(i) + " -hosts=" + strings.Join(cmd.PilosaHosts, ",") + " " + strings.Join(sp.Args, " ")) + if err != nil { + return nil, err } } } @@ -1479,41 +1498,62 @@ func (cmd *BspawnCommand) spawnRemote(ctx context.Context, runUUID uuid.UUID) er for _, sess := range sessions { err = sess.Wait() if err != nil { - return fmt.Errorf("error waiting for remote bagent: %v", err) + return nil, fmt.Errorf("error waiting for remote bagent: %v", err) } } + wg.Wait() + return results, nil +} + +type serialBenchmark struct { + benchmarkers []bench.Benchmark +} + +// Init calls Init for each benchmark. If there are any errors, it will return a +// non-nil error value. +func (sb *serialBenchmark) Init(hosts []string, agentNum int) error { + errors := make([]error, len(sb.benchmarkers)) + hadErr := false + for i, b := range sb.benchmarkers { + errors[i] = b.Init(hosts, agentNum) + if errors[i] != nil { + hadErr = true + } + } + if hadErr { + return fmt.Errorf("Had errs in serialBenchmark.Init: %v", errors) + } return nil } -func (cmd *BspawnCommand) spawnLocal(ctx context.Context, runUUID uuid.UUID) error { - agents := []*BagentCommand{} - for _, sp := range cmd.Benchmarks { - for i := 0; i < sp.Num; i++ { - agentCmd := NewBagentCommand(cmd.Stdin, cmd.Stdout, cmd.Stderr) - agents = append(agents, agentCmd) - err := agentCmd.ParseFlags(append([]string{"-agentNum", strconv.Itoa(i), "-hosts", strings.Join(cmd.PilosaHosts, ","), "-run-uuid", runUUID.String()}, sp.Args...)) - if err != nil { - return err - } - } - } - errors := make([]error, len(agents)) +// Run runs the serial benchmark and returns it's results in a nested map - the +// top level keys are the indices of each benchmark in the list of benchmarks, +// and the values are the results of each benchmark's Run method. +func (sb *serialBenchmark) Run(ctx context.Context, agentNum int) map[string]interface{} { + benchmarks := make([]map[string]interface{}, len(sb.benchmarkers)) + results := map[string]interface{}{"benchmarks": benchmarks} - wg := sync.WaitGroup{} - for i, agent := range agents { - wg.Add(1) - go func(i int, agent *BagentCommand) { - defer wg.Done() - errors[i] = agent.Run(ctx) - }(i, agent) - } - wg.Wait() - for _, err := range errors { - if err != nil { - return fmt.Errorf("%v", errors) + total_start := time.Now() + for i, b := range sb.benchmarkers { + start := time.Now() + output := b.Run(ctx, agentNum) + if _, ok := output["runtime"]; ok { + panic(fmt.Sprintf("Benchmark %v added 'runtime' to its results", b)) } + output["runtime"] = time.Now().Sub(start) + ret := map[string]interface{}{"output": output, "metadata": b} + benchmarks[i] = ret + } + results["total_runtime"] = time.Now().Sub(total_start) + return results +} + +// serial takes a variable number of Benchmarks and returns a Benchmark +// which combines then and will run each serially. +func serial(bs ...bench.Benchmark) bench.Benchmark { + return &serialBenchmark{ + benchmarkers: bs, } - return nil } // readCSVRow reads a bitmap/profile pair from a CSV row. diff --git a/cmd/pilosactl/randSetAndQuery.json b/cmd/pilosactl/randSetAndQuery.json index 714f3d035..2888dc81a 100644 --- a/cmd/pilosactl/randSetAndQuery.json +++ b/cmd/pilosactl/randSetAndQuery.json @@ -1,6 +1,6 @@ { "CreatorArgs": ["-type", "local", "-serverN", "3", "-replicaN", "1"], - "Agents": { "Type": "local" }, + "AgentHosts": [], "Benchmarks": [ { "Num": 3, diff --git a/cmd/pilosactl/randspawn.json b/cmd/pilosactl/randspawn.json index 4c3d52b59..e929d946e 100644 --- a/cmd/pilosactl/randspawn.json +++ b/cmd/pilosactl/randspawn.json @@ -1,6 +1,6 @@ { "CreatorArgs": ["-type", "local", "-serverN", "3", "-replicaN", "1"], - "Agents": { "Type": "local" }, + "AgentHosts": [], "Benchmarks": [ { "Num": 3, diff --git a/cmd/pilosactl/slice-height.json b/cmd/pilosactl/slice-height.json index de3025397..aa0ffab2e 100644 --- a/cmd/pilosactl/slice-height.json +++ b/cmd/pilosactl/slice-height.json @@ -1,10 +1,10 @@ { - "CreatorArgs": ["-type", "local", "-serverN", "1", "-replicaN", "1"], - "Agents": { "Type": "local" }, + "CreatorArgs": ["-hosts=localhost:19444,localhost:19445,localhost:19446", "-log-file-prefix", "asdfe"], + "AgentHosts": [], "Benchmarks": [ { "Num": 1, - "Args": ["slice-height", "-max-time", "1", "-max-bits-per-map", "100"] + "Args": ["-human", "slice-height", "-max-time", "1", "-max-bits-per-map", "100"] } ] } diff --git a/cmd/pilosactl/spawn.json b/cmd/pilosactl/spawn.json index b3daa8a58..deb5e355b 100644 --- a/cmd/pilosactl/spawn.json +++ b/cmd/pilosactl/spawn.json @@ -1,10 +1,17 @@ { - "CreatorArgs": ["-type", "local", "-serverN", "3", "-replicaN", "1"], - "Agents": { "Type": "local" }, + "CreatorArgs": ["-hosts=localhost:19444,localhost:19445,localhost:19446", "-log-file-prefix", "asdfe"], + "AgentHosts": [], + "Human": true, "Benchmarks": [ { "Num": 3, + "Name": "set-diags", "Args": ["diagonal-set-bits", "-iterations", "30000", "-client-type", "round_robin"] + }, + { + "Num": 2, + "Name": "rand-plus-zipf", + "Args": ["random-set-bits", "-iterations", "20000", "zipf-set-bits", "-iterations", "100"] } ] } diff --git a/cmd/pilosactl/zipfspawn.json b/cmd/pilosactl/zipfspawn.json index 6dda5a283..289322727 100644 --- a/cmd/pilosactl/zipfspawn.json +++ b/cmd/pilosactl/zipfspawn.json @@ -1,6 +1,6 @@ { "CreatorArgs": ["-type", "local", "-serverN", "3", "-replicaN", "1"], - "Agents": { "Type": "local" }, + "AgentHosts": [], "Benchmarks": [ { "Num": 1, diff --git a/pilosactl/import.go b/pilosactl/import.go index 34d243f1e..ff157a931 100644 --- a/pilosactl/import.go +++ b/pilosactl/import.go @@ -19,25 +19,25 @@ import ( // ImportCommand represents a command for bulk importing data. type ImportCommand struct { // Destination host and port. - Host string + Host string `json:"host"` // Name of the database & frame to import into. - Database string - Frame string + Database string `json:"db"` + Frame string `json:"frame"` // Filenames to import from. - Paths []string + Paths []string `json:"paths"` // Size of buffer used to chunk import. - BufferSize int + BufferSize int `json:"buffer-size"` // Reusable client. - Client *pilosa.Client + Client *pilosa.Client `json:"-"` // Standard input/output - Stdin io.Reader - Stdout io.Writer - Stderr io.Writer + Stdin io.Reader `json:"-"` + Stdout io.Writer `json:"-"` + Stderr io.Writer `json:"-"` } // NewImportCommand returns a new instance of ImportCommand.