2 Star 1 Fork 0

法马智慧/fmgo

加入 Gitee
与超过 1200万 开发者一起发现、参与优秀开源项目,私有仓库也完全免费 :)
免费加入
克隆/下载
natsx.go 2.63 KB
一键复制 编辑 原始数据 按行查看 历史
Asdybing 提交于 2022-03-11 14:29 . no commit message
package natsx
import (
"encoding/json"
"fmt"
"gitee.com/fmpt/fmgo/logx"
"github.com/nats-io/nats.go"
)
func NewNatsClient(addr string) *NatsClient {
var client NatsClient
client.address = fmt.Sprintf("nats://%s", addr)
client.topics = make(map[string]HFunc)
return &client
}
//重新连接 处理
func (c *NatsClient) onReconnected(conn *nats.Conn) {
c.conn = conn
c.subscribe()
}
//失去连接 处理
func (c *NatsClient) onDisconnected(conn *nats.Conn, err error) {
}
func (c *NatsClient) subscribe() {
for topic, handlerFunc := range c.topics {
// 必须用副本
t, h := topic, handlerFunc
if _, err := c.conn.Subscribe(t, func(msg *nats.Msg) {
h(msg.Subject, msg.Reply, msg.Data)
}); err != nil {
logx.Error("Subscribe Topic error : ", err)
continue
}
}
}
// Subscribe 监听主题 必须在 Start 后
func (c *NatsClient) Subscribe(topic string, h HFunc) {
if _, err := c.conn.Subscribe(topic, func(msg *nats.Msg) {
h(msg.Subject, msg.Reply, msg.Data)
}); err != nil {
logx.Error("Subscribe Topic error : ", err)
}
}
func (c *NatsClient) AddTopic(topic string, hFunc HFunc) {
c.topics[topic] = hFunc
}
func (c *NatsClient) AddTopics(topics map[string]HFunc) {
c.topics = topics
}
func (c *NatsClient) Push(topic string, param *NatsMessage) error {
data, err := json.Marshal(param)
if err != nil {
return err
}
if err = c.conn.Publish(topic, data); err != nil {
return err
}
return nil
}
func (c *NatsClient) Push2(topic string, param interface{}) error {
data, err := json.Marshal(param)
if err != nil {
return err
}
if err = c.conn.Publish(topic, data); err != nil {
return err
}
return nil
}
func (c *NatsClient) Response(topic string, param *NatsMessageRsp) error {
data, err := json.Marshal(param)
if err != nil {
return err
}
if err := c.conn.Publish(topic, data); err != nil {
return err
}
return nil
}
func (c *NatsClient) Request(topic string, param *NatsMessage) (*NatsMessageRsp, error) {
data, err := json.Marshal(param)
if err != nil {
return nil, err
}
msg, err := c.conn.Request(topic, data, natsRequestTimeout)
if err != nil {
return nil, err
}
var res NatsMessageRsp
err = json.Unmarshal(msg.Data, &res)
if err != nil {
return nil, err
}
return &res, nil
}
func (c *NatsClient) Start() {
for {
var err error
c.conn, err = nats.Connect(c.address, nats.MaxReconnects(natsReconnectCount), nats.Timeout(natsConnTimeout), nats.ReconnectWait(natsReconnectWaitTime),
nats.DisconnectErrHandler(c.onDisconnected), nats.ReconnectHandler(c.onReconnected))
if err != nil {
logx.Error("Connect to nats svc error : ", err)
continue
}
go c.subscribe()
break
}
}
马建仓 AI 助手
尝试更多
代码解读
代码找茬
代码优化
Go
1
https://gitee.com/fmpt/fmgo.git
git@gitee.com:fmpt/fmgo.git
fmpt
fmgo
fmgo
v1.3.5

搜索帮助