Merge pull request #229 from jaffee/168-output-format

168 output format
This commit is contained in:
Matthew Jaffee 2017-01-05 11:52:50 -06:00 committed by GitHub
commit b8632e1a88
12 changed files with 145 additions and 201 deletions

View file

@ -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`.

View file

@ -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 &parallelBenchmark{
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,
}
}

View file

@ -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
}

View file

@ -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)

View file

@ -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"]
}
]

View file

@ -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.

View file

@ -1,6 +1,6 @@
{
"CreatorArgs": ["-type", "local", "-serverN", "3", "-replicaN", "1"],
"Agents": { "Type": "local" },
"AgentHosts": [],
"Benchmarks": [
{
"Num": 3,

View file

@ -1,6 +1,6 @@
{
"CreatorArgs": ["-type", "local", "-serverN", "3", "-replicaN", "1"],
"Agents": { "Type": "local" },
"AgentHosts": [],
"Benchmarks": [
{
"Num": 3,

View file

@ -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"]
}
]
}

View file

@ -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"]
}
]
}

View file

@ -1,6 +1,6 @@
{
"CreatorArgs": ["-type", "local", "-serverN", "3", "-replicaN", "1"],
"Agents": { "Type": "local" },
"AgentHosts": [],
"Benchmarks": [
{
"Num": 1,

View file

@ -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.