import(..."github.com/cloudwego/kitex/client""github.com/cloudwego/kitex/pkg/circuitbreak""github.com/cloudwego/kitex/pkg/rpcinfo")// GenServiceCBKeyFunc returns a key which determines the granularity of the CBSuitefuncGenServiceCBKeyFunc(rirpcinfo.RPCInfo)string{// circuitbreak.RPCInfo2Key returns "$fromServiceName/$toServiceName/$method"returncircuitbreak.RPCInfo2Key(ri)}funcmain(){// build a new CBSuite withcbs:=circuitbreak.NewCBSuite(GenServiceCBKeyFunc)varopts[]client.Option// add to the client optionsopts=append(opts,client.WithCircuitBreaker(cbs))// init clientcli,err:=echoservice.NewClient(targetService,opts...)// update circuit breaker config for a certain key (should be consistent with GenServiceCBKeyFunc)// this can be called at any time, and will take effect for following requestscbs.UpdateServiceCBConfig("fromServiceName/toServiceName/method",circuitbreak.CBConfig{Enable:true,ErrRate:0.3,// requests will be blocked if error rate >= 30%MinSample:200,// this config takes effect if sampled requests are more than `MinSample`})// send requests with the client above...}
// Func is the definition for fallback func, which can do fallback both for error and resp.// Notice !! The args and result are not the real rpc req and resp, are respectively XXArgs and XXXResult of generated code.// setup eg: client.WithFallback(fallback.NewFallbackPolicy(yourFunc))typeFuncfunc(ctxcontext.Context,argsutils.KitexArgs,resultutils.KitexResult,errerror)(fbErrerror)// use democlient.WithFallback(fallback.NewFallbackPolicy(func(ctxcontext.Context,argsutils.KitexArgs,resultutils.KitexResult,errerror)(fbErrerror){// your fallback logic...result.SetSuccess(yourFallbackResult)return}))
// RealReqRespFunc is the definition for fallback func with real rpc req as param, and must return the real rpc resp.// setup eg: client.WithFallback(fallback.NewFallbackPolicy(fallback.UnwrapHelper(yourRealReqRespFunc)))typeRealReqRespFuncfunc(ctxcontext.Context,req,respinterface{},errerror)(fbRespinterface{},fbErrerror)// use democlient.WithFallback(fallback.NewFallbackPolicy(fallback.UnwrapHelper(func(ctxcontext.Context,req,respinterface{},errerror)(fbRespinterface{},fbErrerror){// your fallback logic...returnfbResp,fbErr}))
// 方法1:XXXArgs/XXXResult as paramsfallback.NewFallbackPolicy(func(ctxcontext.Context,argsutils.KitexArgs,resultutils.KitexResult,errerror)(fbErrerror){// your fallback logic...result.SetSuccess(yourFallbackResult)return})// 方法2:real rpc req/resp as paramsfallback.NewFallbackPolicy(fallback.UnwrapHelper(func(ctxcontext.Context,req,respinterface{},errerror)(fbRespinterface{},fbErrerror){// your fallback logic...return})
只对 Error(包括业务 Error) 进行 Fallback
1
2
3
4
5
6
7
8
9
10
11
12
// 1: XXXArgs/XXXResult as paramsfallback.ErrorFallback(func(ctxcontext.Context,argsutils.KitexArgs,resultutils.KitexResult,errerror)(fbErrerror){// your fallback logic...result.SetSuccess(yourFallbackResult)return})// 2: real rpc req/resp as paramsfallback.ErrorFallback(fallback.UnwrapHelper(func(ctxcontext.Context,req,respinterface{},errerror)(fbRespinterface{},fbErrerror){// your fallback logic...return})
只对超时和熔断 Error 进行 Fallback
1
2
3
4
5
6
7
8
9
10
11
12
// 1: XXXArgs/XXXResult as paramsfallback.TimeoutAndCBFallback(func(ctxcontext.Context,argsutils.KitexArgs,resultutils.KitexResult,errerror)(fbErrerror){// your fallback logic...result.SetSuccess(yourFallbackResult)return})// 2: real rpc req/resp as paramsfallback.TimeoutAndCBFallback(fallback.UnwrapHelper(func(ctxcontext.Context,req,respinterface{},errerror)(fbRespinterface{},fbErrerror){// your fallback logic...return}))
import("context""time""github.com/cloudwego/kitex/pkg/limiter""github.com/cloudwego/kitex/pkg/rpcinfo""github.com/cloudwego/kitex/server")typeqpsLimiterstruct{}func(l*qpsLimiter)Acquire(ctxcontext.Context)bool{ri:=rpcinfo.GetRPCInfo(ctx)md:=ri.From().Method()returnacquire(md)// return true to allow this request}func(l*qpsLimiter)Status(ctxcontext.Context)(max,currentint,intervaltime.Duration){// max: the maximum number of requests allowed in the interval;// current: the remaining number of requests allowed in the interval;return}typeconnectionLimiterstruct{}func(l*connectionLimiter)Acquire(ctxcontext.Context)bool{ri:=rpcinfo.GetRPCInfo(ctx)addr:=ri.From().Address()returnacquire(addr)// return true to allow this connection}func(l*connectionLimiter)Release(ctxcontext.Context){ri:=rpcinfo.GetRPCInfo(ctx)addr:=ri.From().Address()returnrelease(addr)// release occupied resource by the connection, only called after the release is successful.}func(l*connectionLimiter)Status(ctxcontext.Context)(limit,occupiedint){// limit: the maximum number of connections allowed.// occupied: the number of existing connections.return}funcmain(){myQPSLimiter:=&qpsLimiter{}myConnectionLimiter:=&connectionLimiter{}svr:=xxxservice.NewServer(handler,server.WithQPSLimiter(myQPSLimiter),server.WithConnectionLimiter(myConnectionLimiter))svr.Run()}
// LimitReporter is the interface define to report(metric or print log) when limit happentypeLimitReporterinterface{ConnOverloadReport()QPSOverloadReport()}