2 Star 0 Fork 0

TeamsHub/backend-gopkg

加入 Gitee
与超过 1200万 开发者一起发现、参与优秀开源项目,私有仓库也完全免费 :)
免费加入
克隆/下载
queue.go 2.02 KB
一键复制 编辑 原始数据 按行查看 历史
HCY 提交于 2024-05-10 13:07 . edit pkg
package queue
import (
myconfig "gitee.com/wuzheng0709/backend-gopkg/infrastructure/config"
"gitee.com/wuzheng0709/backend-gopkg/infrastructure/pkg/gin/log"
"github.com/nsqio/go-nsq"
"os"
"os/signal"
"strings"
"sync"
"syscall"
"time"
)
var producer *nsq.Producer
var addrNsqLookups []string
var logLevel nsq.LogLevel
var consumers []*nsq.Consumer
var config *nsq.Config
func Init(addrNsq string, addrNsqLookup string, maxInFlight int) {
log.Info(" nsq 链接中。。。")
if addrNsq == "" {
addrNsq = myconfig.C.Nsq.Address
}
addrNsqLookups = strings.Split(addrNsqLookup, ",")
config = nsq.NewConfig()
config.MaxInFlight = maxInFlight
p, err := nsq.NewProducer(addrNsq, config)
if err != nil {
log.Fatal(err)
return
}
logLevel = nsq.LogLevelWarning
p.SetLogger(log.NsqLogger(), logLevel)
producer = p
if err = p.Ping(); err != nil {
log.Fatal(err)
return
}
log.Info(" nsq 链接成功 ")
go func() {
ch := make(chan os.Signal)
signal.Notify(ch, syscall.SIGINT, syscall.SIGTERM)
<-ch
gracefulStop()
}()
}
func Publish(topic string, body []byte, delay ...time.Duration) (err error) {
log.Info(" nsq 链接成功 ", len(delay), topic, body)
if len(delay) == 0 {
err = producer.Publish(topic, body)
} else {
err = producer.DeferredPublish(topic, delay[0], body)
}
return
}
func Subscribe(topic string, channel string, handler nsq.Handler) (err error) {
c, err := nsq.NewConsumer(topic, channel, config)
if err != nil {
log.Fatal(err)
return
}
c.AddHandler(handler)
c.SetLogger(log.NsqLogger(), logLevel)
err = c.ConnectToNSQLookupds(addrNsqLookups)
if err != nil {
log.Fatal(err)
return
}
consumers = append(consumers, c)
return
}
func gracefulStop() {
producer.Stop()
var wg sync.WaitGroup
for _, c := range consumers {
wg.Add(1)
go func(c *nsq.Consumer) {
c.Stop()
// disconnect from all lookupd
for _, addr := range addrNsqLookups {
err := c.DisconnectFromNSQLookupd(addr)
if err != nil {
log.Error(err)
}
}
wg.Done()
}(c)
}
wg.Wait()
}
马建仓 AI 助手
尝试更多
代码解读
代码找茬
代码优化
1
https://gitee.com/wuzheng0709/backend-gopkg.git
git@gitee.com:wuzheng0709/backend-gopkg.git
wuzheng0709
backend-gopkg
backend-gopkg
v1.4.11

搜索帮助