inprocess.go 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143
  1. package broker
  2. import (
  3. "context"
  4. "fmt"
  5. "sync"
  6. "sync/atomic"
  7. "time"
  8. "git.nspix.com/golang/micro/log"
  9. )
  10. type InProcBus struct {
  11. ctx context.Context
  12. once sync.Once
  13. eventChan chan *Event
  14. subscribers sync.Map
  15. numOfWorker int32
  16. cwChan chan struct{}
  17. }
  18. func (bus *InProcBus) worker() {
  19. for {
  20. select {
  21. case e, ok := <-bus.eventChan:
  22. if ok {
  23. _ = bus.Dispatch(e)
  24. }
  25. case <-bus.cwChan:
  26. return
  27. case <-bus.ctx.Done():
  28. return
  29. }
  30. }
  31. }
  32. func (bus *InProcBus) resizePool(num int) {
  33. needPoolSize := int32(float64(num) / 1.5)
  34. if needPoolSize <= 1 {
  35. needPoolSize = 1
  36. }
  37. if needPoolSize > 20 {
  38. needPoolSize = 20
  39. }
  40. for {
  41. if needPoolSize == atomic.LoadInt32(&bus.numOfWorker) {
  42. break
  43. } else if needPoolSize < atomic.LoadInt32(&bus.numOfWorker) {
  44. bus.cwChan <- struct{}{}
  45. } else {
  46. go bus.worker()
  47. atomic.AddInt32(&bus.numOfWorker, 1)
  48. }
  49. }
  50. }
  51. func (bus *InProcBus) eventLoop() {
  52. ticker := time.NewTicker(time.Second * 5)
  53. defer ticker.Stop()
  54. for {
  55. select {
  56. case <-ticker.C:
  57. bus.resizePool(len(bus.eventChan))
  58. case <-bus.ctx.Done():
  59. return
  60. }
  61. }
  62. }
  63. //Subscribers 获取所有订阅者的快照
  64. func (bus *InProcBus) Subscribers() []Subscriber {
  65. vs := make([]Subscriber, 0)
  66. bus.subscribers.Range(func(key, value interface{}) bool {
  67. vs = append(vs, value.(Subscriber))
  68. return true
  69. })
  70. return vs
  71. }
  72. // Publish 发布一个事件
  73. func (bus *InProcBus) Publish(e *Event) {
  74. select {
  75. case bus.eventChan <- e:
  76. default:
  77. log.Warnf("event queue is full, event %s@%s has been discarded", e.Name, e.Namespace)
  78. }
  79. return
  80. }
  81. // Dispatch 分配事件,关心返回值
  82. func (bus *InProcBus) Dispatch(e *Event) (err error) {
  83. return bus.DispatchCtx(bus.ctx, e)
  84. }
  85. //DispatchCtx 分配一个事件
  86. func (bus *InProcBus) DispatchCtx(ctx context.Context, e *Event) (err error) {
  87. bus.subscribers.Range(func(key, value interface{}) bool {
  88. sub := value.(Subscriber)
  89. if sub.HasTopic(e.Name) {
  90. if err = sub.Process(ctx, e); err != nil {
  91. log.Warnf("subscriber %s process event %s error: %s", sub.ID(), e.Name, err.Error())
  92. }
  93. }
  94. return true
  95. })
  96. releaseEvent(e)
  97. return
  98. }
  99. //Subscribe 订阅一个事件
  100. func (bus *InProcBus) Subscribe(sub Subscriber) (err error) {
  101. if _, ok := bus.subscribers.Load(sub.ID()); ok {
  102. err = fmt.Errorf("subscriber %s already exists", sub.ID())
  103. return
  104. }
  105. if err = sub.OnAttach(); err != nil {
  106. return
  107. }
  108. bus.subscribers.Store(sub.ID(), sub)
  109. return
  110. }
  111. //UnSubscribe 取消一个订阅
  112. func (bus *InProcBus) UnSubscribe(sub Subscriber) (err error) {
  113. bus.subscribers.Delete(sub.ID())
  114. return
  115. }
  116. func (bus *InProcBus) WithContext(ctx context.Context) {
  117. bus.ctx = ctx
  118. }
  119. func NewInPrcBus(ctx context.Context) *InProcBus {
  120. bus := &InProcBus{
  121. ctx: ctx,
  122. eventChan: make(chan *Event, 1024),
  123. }
  124. bus.once.Do(func() {
  125. atomic.AddInt32(&bus.numOfWorker, 1)
  126. go bus.worker()
  127. })
  128. return bus
  129. }