try grpc lb interface
This commit is contained in:
@@ -1,7 +1,8 @@
|
|||||||
package balancer
|
package roundrobin
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
@@ -11,16 +12,16 @@ import (
|
|||||||
"google.golang.org/grpc/resolver"
|
"google.golang.org/grpc/resolver"
|
||||||
)
|
)
|
||||||
|
|
||||||
const Name = "roundrobin"
|
const Name = "zero_rr"
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
balancer.Register(newBuilder())
|
balancer.Register(newRoundRobinBuilder())
|
||||||
}
|
}
|
||||||
|
|
||||||
type roundRobinPickerBuilder struct {
|
type roundRobinPickerBuilder struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
func newBuilder() balancer.Builder {
|
func newRoundRobinBuilder() balancer.Builder {
|
||||||
return base.NewBalancerBuilder(Name, new(roundRobinPickerBuilder))
|
return base.NewBalancerBuilder(Name, new(roundRobinPickerBuilder))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -48,6 +49,7 @@ type roundRobinPicker struct {
|
|||||||
|
|
||||||
func (p *roundRobinPicker) Pick(ctx context.Context, info balancer.PickInfo) (
|
func (p *roundRobinPicker) Pick(ctx context.Context, info balancer.PickInfo) (
|
||||||
conn balancer.SubConn, done func(balancer.DoneInfo), err error) {
|
conn balancer.SubConn, done func(balancer.DoneInfo), err error) {
|
||||||
|
fmt.Println(p.conns)
|
||||||
p.lock.Lock()
|
p.lock.Lock()
|
||||||
defer p.lock.Unlock()
|
defer p.lock.Unlock()
|
||||||
|
|
||||||
36
rpcx/internal/discovclient.go
Normal file
36
rpcx/internal/discovclient.go
Normal file
@@ -0,0 +1,36 @@
|
|||||||
|
package internal
|
||||||
|
|
||||||
|
import (
|
||||||
|
"zero/core/discov"
|
||||||
|
"zero/rpcx/internal/balancer/roundrobin"
|
||||||
|
"zero/rpcx/internal/resolver"
|
||||||
|
|
||||||
|
"google.golang.org/grpc"
|
||||||
|
"google.golang.org/grpc/connectivity"
|
||||||
|
)
|
||||||
|
|
||||||
|
type DiscovClient struct {
|
||||||
|
conn *grpc.ClientConn
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewDiscovClient(etcd discov.EtcdConf, opts ...ClientOption) (*DiscovClient, error) {
|
||||||
|
resolver.RegisterResolver(etcd)
|
||||||
|
opts = append(opts, WithDialOption(grpc.WithBalancerName(roundrobin.Name)))
|
||||||
|
conn, err := dial("discov:///", opts...)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return &DiscovClient{
|
||||||
|
conn: conn,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *DiscovClient) Next() (*grpc.ClientConn, bool) {
|
||||||
|
state := c.conn.GetState()
|
||||||
|
if state == connectivity.Ready {
|
||||||
|
return c.conn, true
|
||||||
|
} else {
|
||||||
|
return nil, false
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -6,15 +6,19 @@ import (
|
|||||||
"google.golang.org/grpc/resolver"
|
"google.golang.org/grpc/resolver"
|
||||||
)
|
)
|
||||||
|
|
||||||
type discovResolver struct {
|
const discovScheme = "discov"
|
||||||
scheme string
|
|
||||||
etcd discov.EtcdConf
|
type discovBuilder struct {
|
||||||
cc resolver.ClientConn
|
etcd discov.EtcdConf
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *discovResolver) Build(target resolver.Target, cc resolver.ClientConn, opts resolver.BuildOptions) (
|
func (b *discovBuilder) Scheme() string {
|
||||||
|
return discovScheme
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b *discovBuilder) Build(target resolver.Target, cc resolver.ClientConn, opts resolver.BuildOptions) (
|
||||||
resolver.Resolver, error) {
|
resolver.Resolver, error) {
|
||||||
sub, err := discov.NewSubscriber(r.etcd.Hosts, r.etcd.Key)
|
sub, err := discov.NewSubscriber(b.etcd.Hosts, b.etcd.Key)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -27,12 +31,18 @@ func (r *discovResolver) Build(target resolver.Target, cc resolver.ClientConn, o
|
|||||||
Addr: val,
|
Addr: val,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
r.cc.UpdateState(resolver.State{
|
cc.UpdateState(resolver.State{
|
||||||
Addresses: addrs,
|
Addresses: addrs,
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
return r, nil
|
return &discovResolver{
|
||||||
|
cc: cc,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type discovResolver struct {
|
||||||
|
cc resolver.ClientConn
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *discovResolver) Close() {
|
func (r *discovResolver) Close() {
|
||||||
@@ -41,13 +51,8 @@ func (r *discovResolver) Close() {
|
|||||||
func (r *discovResolver) ResolveNow(options resolver.ResolveNowOptions) {
|
func (r *discovResolver) ResolveNow(options resolver.ResolveNowOptions) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *discovResolver) Scheme() string {
|
func RegisterResolver(etcd discov.EtcdConf) {
|
||||||
return r.scheme
|
resolver.Register(&discovBuilder{
|
||||||
}
|
etcd: etcd,
|
||||||
|
|
||||||
func RegisterResolver(scheme string, etcd discov.EtcdConf) {
|
|
||||||
resolver.Register(&discovResolver{
|
|
||||||
scheme: scheme,
|
|
||||||
etcd: etcd,
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
17
rpcx/lb/main.go
Normal file
17
rpcx/lb/main.go
Normal file
@@ -0,0 +1,17 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"zero/core/discov"
|
||||||
|
"zero/core/lang"
|
||||||
|
"zero/rpcx/internal"
|
||||||
|
)
|
||||||
|
|
||||||
|
func main() {
|
||||||
|
cli, err := internal.NewDiscovClient(discov.EtcdConf{
|
||||||
|
Hosts: []string{"localhost:2379"},
|
||||||
|
Key: "rpcx",
|
||||||
|
})
|
||||||
|
lang.Must(err)
|
||||||
|
|
||||||
|
cli.Next()
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user