package nsq
|
|
import (
|
"apsClient/conf"
|
"apsClient/constvar"
|
"apsClient/pkg/logx"
|
"apsClient/pkg/nsqclient"
|
"context"
|
"fmt"
|
)
|
|
func Consume(topic, channel string) (err error) {
|
c, err := nsqclient.NewNsqConsumer(context.Background(), topic, channel)
|
if err != nil {
|
logx.Errorf("NewNsqConsumer err:%v", err)
|
return
|
}
|
logx.Infof("Consume NewNsqConsumer topic:%v", topic)
|
var handler MsgHandler
|
switch topic {
|
case fmt.Sprintf(constvar.NsqTopicScheduleTask, conf.Conf.NsqConf.NodeId):
|
handler = new(ScheduleTask)
|
case fmt.Sprintf(constvar.NsqTopicSendPlcAddress, conf.Conf.NsqConf.NodeId):
|
handler = &PlcAddress{Topic: topic}
|
case fmt.Sprintf(constvar.NsqTopicProcessParamsResponse, conf.Conf.NsqConf.NodeId):
|
handler = &ProcessParams{Topic: topic}
|
}
|
c.AddHandler(handler.HandleMessage)
|
|
if len(conf.Conf.NsqConf.NsqlookupdAddr) > 0 {
|
if err = c.RunLookupd(conf.Conf.NsqConf.NsqlookupdAddr, 1); err != nil {
|
logx.Errorf("RunLookupd err:%v", err)
|
return
|
}
|
} else {
|
if err = c.Run(conf.Conf.NsqConf.NsqdAddr, 1); err != nil {
|
logx.Errorf("Run err:%v", err)
|
return
|
}
|
}
|
return
|
}
|