1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
|
package main
import (
"context"
"github.com/go-kit/kit/sd/etcdv3"
"time"
"github.com/go-kit/kit/sd"
"github.com/go-kit/kit/log"
"github.com/go-kit/kit/endpoint"
"io"
"github.com/go-kit/kit/sd/lb"
"fmt"
"google.golang.org/grpc"
"github.com/afex/hystrix-go/hystrix"
"github.com/go-kit/kit/circuitbreaker"
opzipkin "github.com/openzipkin/zipkin-go"
"github.com/go-kit/kit/tracing/zipkin"
grpctransport "github.com/go-kit/kit/transport/grpc"
"grpc-test/pb"
"github.com/openzipkin/zipkin-go/reporter/http"
)
func main() {
commandName := "my-endpoint"
hystrix.ConfigureCommand(commandName, hystrix.CommandConfig{
Timeout: 1000 * 30,
ErrorPercentThreshold: 1,
SleepWindow: 10000,
MaxConcurrentRequests: 1000,
RequestVolumeThreshold: 5,
})
breakerMw := circuitbreaker.Hystrix(commandName)
var (
//注册中心地址
etcdServer = "127.0.0.1:2379"
//监听的服务前缀
prefix = "/services/book/"
ctx = context.Background()
)
options := etcdv3.ClientOptions{
DialTimeout: time.Second * 3,
DialKeepAlive: time.Second * 3,
}
//连接注册中心
client, err := etcdv3.NewClient(ctx, []string{etcdServer}, options)
if err != nil {
panic(err)
}
logger := log.NewNopLogger()
//创建实例管理器, 此管理器会Watch监听etc中prefix的目录变化更新缓存的服务实例数据
instancer, err := etcdv3.NewInstancer(client, prefix, logger)
if err != nil {
panic(err)
}
//创建端点管理器, 此管理器根据Factory和监听的到实例创建endPoint并订阅instancer的变化动态更新Factory创建的endPoint
endpointer := sd.NewEndpointer(instancer, reqFactory, logger)
//创建负载均衡器
balancer := lb.NewRoundRobin(endpointer)
/**
我们可以通过负载均衡器直接获取请求的endPoint,发起请求
reqEndPoint,_ := balancer.Endpoint()
*/
/**
也可以通过retry定义尝试次数进行请求
*/
reqEndPoint := lb.Retry(3, 100*time.Second, balancer)
//增加熔断中间件
reqEndPoint = breakerMw(reqEndPoint)
//现在我们可以通过 endPoint 发起请求了
req := struct{}{}
for i := 1; i <= 1; i++ {
if _, err = reqEndPoint(ctx, req); err != nil {
fmt.Println(err)
}
}
}
//通过传入的 实例地址 创建对应的请求endPoint
func reqFactory(instanceAddr string) (endpoint.Endpoint, io.Closer, error) {
return func(ctx context.Context, request interface{}) (interface{}, error) {
fmt.Println("请求服务: ", instanceAddr, "当前时间: ", time.Now().Format("2006-01-02 15:04:05.99"))
conn, err := grpc.Dial(instanceAddr, grpc.WithInsecure())
if err != nil {
fmt.Println(err)
panic("connect error")
}
reporter := http.NewReporter("http://localhost:9411/api/v2/spans")
defer reporter.Close()
zkTracer, err := opzipkin.NewTracer(reporter)
zkClientTrace := zipkin.GRPCClientTrace(zkTracer)
bookInfoRequest := grpctransport.NewClient(
conn,
"BookService",
"GetBookInfo",
func(_ context.Context, in interface{}) (interface{}, error) { return nil, nil },
func(_ context.Context, out interface{}) (interface{}, error) {
return out, nil
},
book.BookInfo{},
zkClientTrace,
).Endpoint()
bookListRequest := grpctransport.NewClient(
conn,
"BookService",
"GetBookList",
func(_ context.Context, in interface{}) (interface{}, error) { return nil, nil },
func(_ context.Context, out interface{}) (interface{}, error) {
return out, nil
},
book.BookList{},
zkClientTrace,
).Endpoint()
parentSpan := zkTracer.StartSpan("bookCaller")
defer parentSpan.Flush()
ctx = opzipkin.NewContext(ctx, parentSpan)
infoRet ,_:= bookInfoRequest(ctx, request)
bi := infoRet.(*book.BookInfo)
fmt.Println("获取书籍详情")
fmt.Println("bookId: 1", " => ", "bookName:", bi.BookName)
listRet,_ := bookListRequest(ctx, request)
bl := listRet.(*book.BookList)
fmt.Println("获取书籍列表")
for _,b := range bl.BookList {
fmt.Println("bookId:", b.BookId, " => ", "bookName:", b.BookName)
}
return nil,nil
},nil,nil
}
|