代码拉取完成,页面将自动刷新
package nsq
import (
"errors"
"fmt"
"net/http"
"sync"
_nsq "github.com/nsqio/go-nsq"
"gitee.com/pangxianfei/multiapp/kernel/zone"
)
type nsq struct {
producer *producer
conn string
consumerList map[hashTopicChannel]*consumer
Lock sync.RWMutex
connectedProducerList []*producer
connectedConsumerList []*consumer
}
func (n *nsq) Register(topicName string, channelName string) (err error) {
if !_nsq.IsValidTopicName(topicName) {
return errors.New("topic name is invalid")
}
if !_nsq.IsValidChannelName(channelName) {
return errors.New("channel name is invalid")
}
// default register to nsqd index 0
resp, mqErr := http.Post(fmt.Sprintf("%s/channel/create?topic=%s&channel=%s", n.nsqlookupdHttpConnectionArgsList()[0], topicName, channelName), "", nil)
if mqErr != nil {
return mqErr
}
if resp.StatusCode != 200 {
return errors.New(fmt.Sprintf("response error, code: %d", resp.StatusCode))
}
return nil
}
func (n *nsq) Unregister(topicName string, channelName string) (err error) { return }
func (n *nsq) Push(topicName string, channelName string, delay zone.Duration, body []byte) (err error) {
// for SupportBroadCasting queue driver channelName is not used
if delay > 0 {
return n.producer.p.DeferredPublish(topicName, delay, body)
}
return n.producer.p.Publish(topicName, body)
}
func (n *nsq) Pop(topicName string, channelName string, handler func(hash string, body []byte) (handlerErr error), maxInFlight int) (err error) {
htc, mqErr := n.consumerConnect(topicName, channelName)
if err != nil {
return mqErr
}
n.consumer(htc).c.AddHandler(_nsq.HandlerFunc(func(message *_nsq.Message) error {
return handler(string(message.ID[:]), message.Body)
}))
if mqErr := n.consumer(htc).c.ConnectToNSQDs(n.nsqdTcpConnectionArgsList()); err != nil {
return mqErr
}
n.consumer(htc).c.ChangeMaxInFlight(maxInFlight)
n.addConnectedConsumer(n.consumer(htc))
return nil
}
func (n *nsq) SupportBroadCasting() bool {
return true
}
func (n *nsq) Close() (err error) {
// stop producer
for _, p := range n.connectedProducerList {
p.p.Stop()
}
// stop consumer
for _, c := range n.connectedConsumerList {
c.c.Stop()
}
return nil
}
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。