From 7a7071a7ad57665884888cf1a59039343502cca5 Mon Sep 17 00:00:00 2001 From: Jeffrey Wilcke Date: Fri, 19 Feb 2016 20:22:24 +0100 Subject: [PATCH] balancer: implemented generic load balancer Balancer is a generic approach to load balancing given any type of tasks. Tasks are packages that do 1. perform an expensive tasks and 2. yield a good/bad result (i.e. an error). When tasks are given to the balancer the balancer will selects one of the least loaded workers and passes on the work and adds it to its work queue. The load balancer can be used in all sorts of places and should generally help with optimisations in places where parallel processing is useful. The following tasks could be paralellised: * The transaction pool's validation process whenever it's given a bunch of transactions. * Block & uncle pow verification * Transaction sender public key derivation --- balancer/balancer.go | 139 ++++++++++++++++++++++++++++++++++++++ balancer/balancer_test.go | 86 +++++++++++++++++++++++ 2 files changed, 225 insertions(+) create mode 100644 balancer/balancer.go create mode 100644 balancer/balancer_test.go diff --git a/balancer/balancer.go b/balancer/balancer.go new file mode 100644 index 0000000000..8ea39aabc3 --- /dev/null +++ b/balancer/balancer.go @@ -0,0 +1,139 @@ +package balancer + +import "container/heap" + +// Task repsents a single batch of work offered to a worker. +type Task struct { + fn func() error // work function + c chan error // return channel +} + +// NewTask returns a new task and sets the proper fields. +func NewTask(fn func() error, c chan error) Task { + return Task{ + fn: fn, + c: c, + } +} + +// Worker is a worker that will take one it's assigned tasks +// and execute it +type Worker struct { + id int // worker id + tasks chan Task // tasks to do (buffered) + pending int // count of pending work + index int // index in the heap +} + +// work will take the oldest task and execute the function and +// yield the result back in to the return error channel. +func (w *Worker) work(done chan *Worker) { + for { + task := <-w.tasks // get task... + task.c <- task.fn() // ...execute the task + done <- w // we're done + } +} + +// Pool is a pool of workers and implements containers.Heap +type Pool []*Worker + +func (p Pool) Len() int { return len(p) } +func (p Pool) Less(i, j int) bool { return p[i].pending < p[j].pending } +func (p Pool) Swap(i, j int) { p[i], p[j] = p[j], p[i] } +func (p *Pool) Push(x interface{}) { + w := x.(*Worker) // cast the worker + w.index = len(*p) // assign the new index + + *p = append(*p, x.(*Worker)) +} + +func (p *Pool) Pop() interface{} { + old := *p + n := len(old) + x := old[n-1] + *p = old[0 : n-1] + return x +} + +// Balancer is responsible for balancing any given tasks +// to the pool of workers. The workers are managed by the +// balancer and will try to make sure that the workers are +// equally balanced in "work to complete". +type Balancer struct { + pool Pool + done chan *Worker + + work chan Task +} + +// New returns a new load balancer +func New(poolSize int) *Balancer { + balancer := &Balancer{ + done: make(chan *Worker), + pool: make(Pool, poolSize), + work: make(chan Task), + } + + operations := make(chan struct{}, poolSize) + defer close(operations) + + // fill the pool with the given pool size + for i := 0; i < poolSize; i++ { + // create new worker + balancer.pool[i] = &Worker{id: i, tasks: make(chan Task, 10)} + // spawn worker process + go func(i int) { + operations <- struct{}{} + balancer.pool[i].work(balancer.done) + }(i) + } + // spawn own balancer task + go balancer.balance(balancer.work) + + // wait for workers to be operations + for i := 0; i < poolSize; i++ { + <-operations + } + + return balancer +} + +// Push pushes the given tasks in to the work channel. +func (b *Balancer) Push(work Task) { + go func() { b.work <- work }() +} + +func (b *Balancer) balance(work chan Task) { + for { + select { + case task := <-work: // get task + b.dispatch(task) // dispatch the tasks + case w := <-b.done: // worker is done + b.completed(w) // handle worker + } + } +} + +// dispatch dispatches the tasks to the least loaded worker. +func (b *Balancer) dispatch(task Task) { + // Take least loaded worker + w := heap.Pop(&b.pool).(*Worker) + // send it a task + w.tasks <- task + // add to its queue + w.pending++ + // put it back in the heap + heap.Push(&b.pool, w) +} + +// completed handles the worker and puts it back in the pool +// based on it's load. +func (b *Balancer) completed(w *Worker) { + // reduce one task + w.pending-- + // remove it from the heap + heap.Remove(&b.pool, w.index) + // put it back in place + heap.Push(&b.pool, w) +} diff --git a/balancer/balancer_test.go b/balancer/balancer_test.go new file mode 100644 index 0000000000..f645db01c0 --- /dev/null +++ b/balancer/balancer_test.go @@ -0,0 +1,86 @@ +package balancer + +import ( + "math/big" + "testing" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/crypto" +) + +func makeTxs(size int, b *testing.B) []types.Transactions { + key, _ := crypto.GenerateKey() + + batches := make([]types.Transactions, b.N) + for i := 0; i < b.N; i++ { + txs := make(types.Transactions, size) + for j := range txs { + var err error + txs[j], err = types.NewTransaction(0, common.Address{}, new(big.Int), new(big.Int), new(big.Int), nil).SignECDSA(key) + if err != nil { + b.Fatal(err) + } + } + batches[i] = txs + } + return batches +} + +const benchTxSize = 1024 + +func BenchmarkTxsRaw(b *testing.B) { + batches := makeTxs(benchTxSize, b) + b.ResetTimer() + + for i := 0; i < b.N; i++ { + txs := batches[i] + b.StartTimer() + + for _, tx := range txs { + if _, err := tx.From(); err != nil { + b.Fatal(err) + } + } + } +} + +func BenchmarkTxsLoadBalancer(b *testing.B) { + balancer := New(4) + batches := makeTxs(benchTxSize, b) + b.ResetTimer() + + for i := 0; i < b.N; i++ { + txs := batches[i] + + var ( + size = 8 + batchsize = len(txs) / size + ch = make(chan error, size) + ) + + for j := 0; j < size; j++ { + j := j + task := Task{ + fn: func() error { + for _, tx := range txs[j*size : j*size+batchsize] { + if _, err := tx.From(); err != nil { + return err + } + } + return nil + }, + c: ch, + } + + balancer.Push(task) + } + + for j := 0; j < size; j++ { + if err := <-ch; err != nil { + b.Error(err) + } + } + close(ch) + } +}