1
0
mirror of https://github.com/go-micro/go-micro.git synced 2024-12-24 10:07:04 +02:00

Add the notion of a node selector for routing

This commit is contained in:
Asim 2015-12-07 21:09:10 +00:00
parent 050fa26afb
commit 4e6b9347d9
4 changed files with 106 additions and 2 deletions

View File

@ -8,10 +8,18 @@ import (
"github.com/micro/go-micro/registry"
)
// NodeSelector is used to retrieve a node to which a request
// should be routed. It takes a list of services and selects
// a single node. If a node cannot be selected it should return
// an error. A list of services is provided as a service may
// have 1 or more versions.
type NodeSelector func(service []*registry.Service) (*registry.Node, error)
func init() {
rand.Seed(time.Now().UnixNano())
}
// Built in random hashed node selector
func nodeSelector(service []*registry.Service) (*registry.Node, error) {
if len(service) == 0 {
return nil, errors.NotFound("go.micro.client", "Service not found")

View File

@ -14,6 +14,7 @@ type options struct {
registry registry.Registry
transport transport.Transport
wrappers []Wrapper
selector NodeSelector
}
// Broker to be used for pub/sub
@ -51,6 +52,13 @@ func Transport(t transport.Transport) Option {
}
}
// Selector is used to select a node to route a request to
func Selector(s NodeSelector) Option {
return func(o *options) {
o.selector = s
}
}
// Adds a Wrapper to a list of options passed into the client
func Wrap(w Wrapper) Option {
return func(o *options) {

View File

@ -43,6 +43,10 @@ func newRpcClient(opt ...Option) Client {
opts.broker = broker.DefaultBroker
}
if opts.selector == nil {
opts.selector = nodeSelector
}
rc := &rpcClient{
once: once,
opts: opts,
@ -146,7 +150,7 @@ func (r *rpcClient) Call(ctx context.Context, request Request, response interfac
return errors.InternalServerError("go.micro.client", err.Error())
}
node, err := nodeSelector(service)
node, err := r.opts.selector(service)
if err != nil {
return err
}
@ -169,7 +173,7 @@ func (r *rpcClient) Stream(ctx context.Context, request Request, responseChan in
return nil, errors.InternalServerError("go.micro.client", err.Error())
}
node, err := nodeSelector(service)
node, err := r.opts.selector(service)
if err != nil {
return nil, err
}

View File

@ -0,0 +1,84 @@
package main
import (
"fmt"
"math/rand"
"time"
"github.com/micro/go-micro/client"
"github.com/micro/go-micro/cmd"
c "github.com/micro/go-micro/context"
"github.com/micro/go-micro/errors"
example "github.com/micro/go-micro/examples/server/proto/example"
"github.com/micro/go-micro/registry"
"golang.org/x/net/context"
)
func init() {
rand.Seed(time.Now().Unix())
}
// A random node selector
func randomSelector(s []*registry.Service) (*registry.Node, error) {
if len(s) == 0 {
return nil, errors.NotFound("go.micro.client", "Service not found")
}
i := rand.Int()
j := i % len(s)
if len(s[j].Nodes) == 0 {
return nil, errors.NotFound("go.micro.client", "Service not found")
}
n := i % len(s[j].Nodes)
return s[j].Nodes[n], nil
}
// Wraps the node selector so that it will log what node was selected
func wrapSelector(fn client.NodeSelector) client.NodeSelector {
return func(s []*registry.Service) (*registry.Node, error) {
n, err := fn(s)
if err != nil {
return nil, err
}
fmt.Printf("Selected node %v\n", n)
return n, nil
}
}
func call(i int) {
// Create new request to service go.micro.srv.example, method Example.Call
req := client.NewRequest("go.micro.srv.example", "Example.Call", &example.Request{
Name: "John",
})
// create context with metadata
ctx := c.WithMetadata(context.Background(), map[string]string{
"X-User-Id": "john",
"X-From-Id": "script",
})
rsp := &example.Response{}
// Call service
if err := client.Call(ctx, req, rsp); err != nil {
fmt.Println("call err: ", err, rsp)
return
}
fmt.Println("Call:", i, "rsp:", rsp.Msg)
}
func main() {
cmd.Init()
client.DefaultClient = client.NewClient(
client.Selector(wrapSelector(randomSelector)),
)
fmt.Println("\n--- Call example ---\n")
for i := 0; i < 10; i++ {
call(i)
}
}