go-zero/core/syncx/singleflight.go

82 lines
1.7 KiB
Go
Raw Normal View History

2020-07-26 17:09:05 +08:00
package syncx
import "sync"
type (
// SingleFlight lets the concurrent calls with the same key to share the call result.
2020-07-26 17:09:05 +08:00
// For example, A called F, before it's done, B called F. Then B would not execute F,
// and shared the result returned by F which called by A.
// The calls with the same key are dependent, concurrent calls share the returned values.
// A ------->calls F with key<------------------->returns val
// B --------------------->calls F with key------>returns val
SingleFlight interface {
Do(key string, fn func() (any, error)) (any, error)
DoEx(key string, fn func() (any, error)) (any, bool, error)
2020-07-26 17:09:05 +08:00
}
call struct {
wg sync.WaitGroup
val any
2020-07-26 17:09:05 +08:00
err error
}
flightGroup struct {
2020-07-26 17:09:05 +08:00
calls map[string]*call
lock sync.Mutex
}
)
// NewSingleFlight returns a SingleFlight.
func NewSingleFlight() SingleFlight {
return &flightGroup{
2020-07-26 17:09:05 +08:00
calls: make(map[string]*call),
}
}
func (g *flightGroup) Do(key string, fn func() (any, error)) (any, error) {
c, done := g.createCall(key)
2020-09-30 12:31:35 +08:00
if done {
2020-07-26 17:09:05 +08:00
return c.val, c.err
}
2020-09-30 12:31:35 +08:00
g.makeCall(c, key, fn)
2020-07-26 17:09:05 +08:00
return c.val, c.err
}
func (g *flightGroup) DoEx(key string, fn func() (any, error)) (val any, fresh bool, err error) {
c, done := g.createCall(key)
2020-09-30 12:31:35 +08:00
if done {
2020-07-26 17:09:05 +08:00
return c.val, false, c.err
}
2020-09-30 12:31:35 +08:00
g.makeCall(c, key, fn)
2020-07-26 17:09:05 +08:00
return c.val, true, c.err
}
func (g *flightGroup) createCall(key string) (c *call, done bool) {
2020-09-30 12:31:35 +08:00
g.lock.Lock()
if c, ok := g.calls[key]; ok {
g.lock.Unlock()
c.wg.Wait()
return c, true
}
c = new(call)
2020-07-26 17:09:05 +08:00
c.wg.Add(1)
g.calls[key] = c
g.lock.Unlock()
2020-09-30 12:31:35 +08:00
return c, false
}
func (g *flightGroup) makeCall(c *call, key string, fn func() (any, error)) {
2020-07-26 17:09:05 +08:00
defer func() {
g.lock.Lock()
delete(g.calls, key)
g.lock.Unlock()
c.wg.Done()
}()
c.val, c.err = fn()
}