Created
October 28, 2015 16:03
-
-
Save mmirolim/c0921726bd79dba10c6a to your computer and use it in GitHub Desktop.
example usage for fallback with circuit breaker
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "errors" | |
| "log" | |
| "sync/atomic" | |
| "time" | |
| as "github.com/aerospike/aerospike-client-go" | |
| "github.com/rubyist/circuitbreaker" | |
| "gopkg.in/mgo.v2" | |
| ) | |
| var ( | |
| asErrCounter int64 = 0 | |
| mongoErrCounter int64 = 0 | |
| ac *as.Client | |
| ms *mgo.Session | |
| cbPanel *circuit.Panel | |
| ErrAsNotConnected = errors.New("aerospike client not connected") | |
| ErrMgoNotConnected = errors.New("mongo client not connected") | |
| ) | |
| func init() { | |
| var err error | |
| // define a client to connect to | |
| ac, err = as.NewClient("127.0.0.1", 3000) | |
| fatalOnError(err) | |
| // try to connect, set timeout for request | |
| ms, err = mgo.DialWithTimeout("127.0.0.1", time.Second) | |
| fatalOnError(err) | |
| } | |
| func main() { | |
| // create circuit panel for all required services | |
| cbPanel = healthCheckServices("as", "mongo") | |
| // lets check how it behaves | |
| time.Sleep(60 * time.Second) | |
| } | |
| // add circuit breaker to panel | |
| func healthCheckServices(asname, mongoname string) *circuit.Panel { | |
| cbpanel := circuit.NewPanel() | |
| // Creates a circuit breaker that will trip if the function fails 10 times | |
| cb := circuit.NewThresholdBreaker(10) | |
| cbpanel.Add(asname, cb) | |
| asCbEvents := cb.Subscribe() | |
| go func() { | |
| log.Println("start listening to aeCbEvents") | |
| for { | |
| e := <-asCbEvents | |
| // Monitor breaker events like BreakerTripped, BreakerReset, BreakerFail, BreakerReady | |
| switch e { | |
| case circuit.BreakerTripped: | |
| log.Println("aerospike breaker tripped") | |
| case circuit.BreakerReset: | |
| log.Println("aerospike breaker reset") | |
| case circuit.BreakerFail: | |
| log.Println("aerospike service as unavailable") | |
| atomic.AddInt64(&asErrCounter, 1) | |
| log.Println("aerospikeErrCounter", atomic.LoadInt64(&asErrCounter)) | |
| case circuit.BreakerReady: | |
| log.Println("aerospike breaker ready") | |
| atomic.AddInt64(&asErrCounter, -1) | |
| log.Println("aerospike ErrCounter after BreakerReady", atomic.LoadInt64(&asErrCounter)) | |
| } | |
| } | |
| }() | |
| // start pings to services | |
| go func() { | |
| cb, _ := cbpanel.Get(asname) | |
| for { | |
| // make ping every 100ms | |
| time.Sleep(100 * time.Millisecond) | |
| cb.Call(func() error { | |
| key, _ := as.NewKey("test", "ping", "key") | |
| if _, err := ac.Get(nil, key); err != nil { | |
| log.Println("aerospike putbins error", err) | |
| return ErrAsNotConnected | |
| } | |
| log.Println("aerospike ok") | |
| return nil | |
| }, 0) | |
| } | |
| }() | |
| // Creates a circuit breaker that will trip if the function fails 10 times | |
| cb = circuit.NewThresholdBreaker(10) | |
| cbpanel.Add(mongoname, cb) | |
| mCbEvents := cb.Subscribe() | |
| go func() { | |
| log.Println("start listening to mCbEvents") | |
| for { | |
| e := <-mCbEvents | |
| // Monitor breaker events like BreakerTripped, BreakerReset, BreakerFail, BreakerReady | |
| switch e { | |
| case circuit.BreakerTripped: | |
| log.Println("mongo breaker tripped") | |
| case circuit.BreakerReset: | |
| log.Println("mongo breaker reset") | |
| case circuit.BreakerFail: | |
| log.Println("mongo service as unavailable") | |
| atomic.AddInt64(&mongoErrCounter, 1) | |
| log.Println("mongoErrCounter", atomic.LoadInt64(&mongoErrCounter)) | |
| case circuit.BreakerReady: | |
| log.Println("mongo breaker ready") | |
| atomic.AddInt64(&mongoErrCounter, -1) | |
| log.Println("mongoErrCounter after BreakerReady", atomic.LoadInt64(&mongoErrCounter)) | |
| } | |
| } | |
| }() | |
| // start pings to services | |
| go func() { | |
| cb, _ := cbpanel.Get(mongoname) | |
| for { | |
| // make ping every 100ms | |
| time.Sleep(500 * time.Millisecond) | |
| cb.Call(func() error { | |
| ms.Refresh() | |
| if err := ms.Ping(); err != nil { | |
| log.Println("mongo ping err", err) | |
| return err | |
| } | |
| log.Println("mongo ok") | |
| return nil | |
| }, 0) | |
| } | |
| }() | |
| return cbpanel | |
| } | |
| func fatalOnError(err error) { | |
| if err != nil { | |
| log.Fatal(err) | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment