queue.go 5.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251
  1. package queue
  2. import (
  3. "runtime"
  4. "sync"
  5. "sync/atomic"
  6. "time"
  7. "git.i2edu.net/i2/go-zero/core/logx"
  8. "git.i2edu.net/i2/go-zero/core/rescue"
  9. "git.i2edu.net/i2/go-zero/core/stat"
  10. "git.i2edu.net/i2/go-zero/core/threading"
  11. "git.i2edu.net/i2/go-zero/core/timex"
  12. )
  13. const queueName = "queue"
  14. type (
  15. // A Queue is a message queue.
  16. Queue struct {
  17. name string
  18. metrics *stat.Metrics
  19. producerFactory ProducerFactory
  20. producerRoutineGroup *threading.RoutineGroup
  21. consumerFactory ConsumerFactory
  22. consumerRoutineGroup *threading.RoutineGroup
  23. producerCount int
  24. consumerCount int
  25. active int32
  26. channel chan string
  27. quit chan struct{}
  28. listeners []Listener
  29. eventLock sync.Mutex
  30. eventChannels []chan interface{}
  31. }
  32. // A Listener interface represents a listener that can be notified with queue events.
  33. Listener interface {
  34. OnPause()
  35. OnResume()
  36. }
  37. // A Poller interface wraps the method Poll.
  38. Poller interface {
  39. Name() string
  40. Poll() string
  41. }
  42. // A Pusher interface wraps the method Push.
  43. Pusher interface {
  44. Name() string
  45. Push(string) error
  46. }
  47. )
  48. // NewQueue returns a Queue.
  49. func NewQueue(producerFactory ProducerFactory, consumerFactory ConsumerFactory) *Queue {
  50. q := &Queue{
  51. metrics: stat.NewMetrics(queueName),
  52. producerFactory: producerFactory,
  53. producerRoutineGroup: threading.NewRoutineGroup(),
  54. consumerFactory: consumerFactory,
  55. consumerRoutineGroup: threading.NewRoutineGroup(),
  56. producerCount: runtime.NumCPU(),
  57. consumerCount: runtime.NumCPU() << 1,
  58. channel: make(chan string),
  59. quit: make(chan struct{}),
  60. }
  61. q.SetName(queueName)
  62. return q
  63. }
  64. // AddListener adds a listener to q.
  65. func (q *Queue) AddListener(listener Listener) {
  66. q.listeners = append(q.listeners, listener)
  67. }
  68. // Broadcast broadcasts message to all event channels.
  69. func (q *Queue) Broadcast(message interface{}) {
  70. go func() {
  71. q.eventLock.Lock()
  72. defer q.eventLock.Unlock()
  73. for _, channel := range q.eventChannels {
  74. channel <- message
  75. }
  76. }()
  77. }
  78. // SetName sets the name of q.
  79. func (q *Queue) SetName(name string) {
  80. q.name = name
  81. q.metrics.SetName(name)
  82. }
  83. // SetNumConsumer sets the number of consumers.
  84. func (q *Queue) SetNumConsumer(count int) {
  85. q.consumerCount = count
  86. }
  87. // SetNumProducer sets the number of producers.
  88. func (q *Queue) SetNumProducer(count int) {
  89. q.producerCount = count
  90. }
  91. // Start starts q.
  92. func (q *Queue) Start() {
  93. q.startProducers(q.producerCount)
  94. q.startConsumers(q.consumerCount)
  95. q.producerRoutineGroup.Wait()
  96. close(q.channel)
  97. q.consumerRoutineGroup.Wait()
  98. }
  99. // Stop stops q.
  100. func (q *Queue) Stop() {
  101. close(q.quit)
  102. }
  103. func (q *Queue) consume(eventChan chan interface{}) {
  104. var consumer Consumer
  105. for {
  106. var err error
  107. if consumer, err = q.consumerFactory(); err != nil {
  108. logx.Errorf("Error on creating consumer: %v", err)
  109. time.Sleep(time.Second)
  110. } else {
  111. break
  112. }
  113. }
  114. for {
  115. select {
  116. case message, ok := <-q.channel:
  117. if ok {
  118. q.consumeOne(consumer, message)
  119. } else {
  120. logx.Info("Task channel was closed, quitting consumer...")
  121. return
  122. }
  123. case event := <-eventChan:
  124. consumer.OnEvent(event)
  125. }
  126. }
  127. }
  128. func (q *Queue) consumeOne(consumer Consumer, message string) {
  129. threading.RunSafe(func() {
  130. startTime := timex.Now()
  131. defer func() {
  132. duration := timex.Since(startTime)
  133. q.metrics.Add(stat.Task{
  134. Duration: duration,
  135. })
  136. logx.WithDuration(duration).Infof("%s", message)
  137. }()
  138. if err := consumer.Consume(message); err != nil {
  139. logx.Errorf("Error occurred while consuming %v: %v", message, err)
  140. }
  141. })
  142. }
  143. func (q *Queue) pause() {
  144. for _, listener := range q.listeners {
  145. listener.OnPause()
  146. }
  147. }
  148. func (q *Queue) produce() {
  149. var producer Producer
  150. for {
  151. var err error
  152. if producer, err = q.producerFactory(); err != nil {
  153. logx.Errorf("Error on creating producer: %v", err)
  154. time.Sleep(time.Second)
  155. } else {
  156. break
  157. }
  158. }
  159. atomic.AddInt32(&q.active, 1)
  160. producer.AddListener(routineListener{
  161. queue: q,
  162. })
  163. for {
  164. select {
  165. case <-q.quit:
  166. logx.Info("Quitting producer")
  167. return
  168. default:
  169. if v, ok := q.produceOne(producer); ok {
  170. q.channel <- v
  171. }
  172. }
  173. }
  174. }
  175. func (q *Queue) produceOne(producer Producer) (string, bool) {
  176. // avoid panic quit the producer, just log it and continue
  177. defer rescue.Recover()
  178. return producer.Produce()
  179. }
  180. func (q *Queue) resume() {
  181. for _, listener := range q.listeners {
  182. listener.OnResume()
  183. }
  184. }
  185. func (q *Queue) startConsumers(number int) {
  186. for i := 0; i < number; i++ {
  187. eventChan := make(chan interface{})
  188. q.eventLock.Lock()
  189. q.eventChannels = append(q.eventChannels, eventChan)
  190. q.eventLock.Unlock()
  191. q.consumerRoutineGroup.Run(func() {
  192. q.consume(eventChan)
  193. })
  194. }
  195. }
  196. func (q *Queue) startProducers(number int) {
  197. for i := 0; i < number; i++ {
  198. q.producerRoutineGroup.Run(func() {
  199. q.produce()
  200. })
  201. }
  202. }
  203. type routineListener struct {
  204. queue *Queue
  205. }
  206. func (rl routineListener) OnProducerPause() {
  207. if atomic.AddInt32(&rl.queue.active, -1) <= 0 {
  208. rl.queue.pause()
  209. }
  210. }
  211. func (rl routineListener) OnProducerResume() {
  212. if atomic.AddInt32(&rl.queue.active, 1) == 1 {
  213. rl.queue.resume()
  214. }
  215. }