producer.go 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139
  1. package nsqclient
  2. import (
  3. "fmt"
  4. "time"
  5. nsq "github.com/nsqio/go-nsq"
  6. )
  7. type Producer interface {
  8. Publish(topic string, body []byte) error
  9. MultiPublish(topic string, body [][]byte) error
  10. DeferredPublish(topic string, delay time.Duration, body []byte) error
  11. }
  12. var _ Producer = (*producer)(nil)
  13. type producer struct {
  14. pool Pool
  15. }
  16. var (
  17. // name pool producer
  18. nsqList = make(map[string]Pool)
  19. )
  20. type Config struct {
  21. Addr string `toml:"addr" json:"addr"`
  22. InitSize int `toml:"init_size" json:"init_size"`
  23. MaxSize int `toml:"max_size" json:"max_size"`
  24. }
  25. func CreateProducerPool(configs map[string]Config) {
  26. for name, conf := range configs {
  27. n, err := newProducerPool(conf.Addr, conf.InitSize, conf.MaxSize)
  28. if err == nil {
  29. nsqList[name] = n
  30. // 支持ip:port寻址
  31. nsqList[conf.Addr] = n
  32. }
  33. }
  34. }
  35. func DestroyProducerPool() {
  36. for _, p := range nsqList {
  37. p.Close()
  38. }
  39. }
  40. func GetProducer(key ...string) (*producer, error) {
  41. k := "default"
  42. if len(key) > 0 {
  43. k = key[0]
  44. }
  45. if n, ok := nsqList[k]; ok {
  46. return &producer{n}, nil
  47. }
  48. return nil, fmt.Errorf("GetProducer can't get producer")
  49. }
  50. // CreateNSQProducer create nsq producer
  51. func newProducer(addr string, options ...func(*nsq.Config)) (*nsq.Producer, error) {
  52. cfg := nsq.NewConfig()
  53. for _, option := range options {
  54. option(cfg)
  55. }
  56. producer, err := nsq.NewProducer(addr, cfg)
  57. if err != nil {
  58. return nil, err
  59. }
  60. // producer.SetLogger(log.New(os.Stderr, "", log.Flags()), nsq.LogLevelError)
  61. return producer, nil
  62. }
  63. // CreateNSQProducerPool create a nwq producer pool
  64. func newProducerPool(addr string, initSize, maxSize int, options ...func(*nsq.Config)) (Pool, error) {
  65. factory := func() (*nsq.Producer, error) {
  66. // TODO 这里应该执行ping方法来确定连接是正常的否则不应该创建conn
  67. return newProducer(addr, options...)
  68. }
  69. nsqPool, err := NewChannelPool(initSize, maxSize, factory)
  70. if err != nil {
  71. return nil, err
  72. }
  73. return nsqPool, nil
  74. }
  75. func NewProducer(addr string) (*producer, error) {
  76. CreateProducerPool(map[string]Config{"default": {addr, 1, 1}})
  77. return GetProducer()
  78. }
  79. func retry(num int, fn func() error) error {
  80. var err error
  81. for i := 0; i < num; i++ {
  82. err = fn()
  83. if err == nil {
  84. break
  85. }
  86. }
  87. return err
  88. }
  89. func (p *producer) Publish(topic string, body []byte) error {
  90. nsq, err := p.pool.Get()
  91. if err != nil {
  92. return err
  93. }
  94. defer nsq.Close()
  95. return retry(2, func() error {
  96. return nsq.Publish(topic, body)
  97. })
  98. }
  99. func (p *producer) MultiPublish(topic string, body [][]byte) error {
  100. nsq, err := p.pool.Get()
  101. if err != nil {
  102. return err
  103. }
  104. defer nsq.Close()
  105. return retry(2, func() error {
  106. return nsq.MultiPublish(topic, body)
  107. })
  108. }
  109. func (p *producer) DeferredPublish(topic string, delay time.Duration, body []byte) error {
  110. nsq, err := p.pool.Get()
  111. if err != nil {
  112. return err
  113. }
  114. defer nsq.Close()
  115. return retry(2, func() error {
  116. return nsq.DeferredPublish(topic, delay, body)
  117. })
  118. }