1 Star 0 Fork 0

qw_1215/glink

加入 Gitee
与超过 1200万 开发者一起发现、参与优秀开源项目,私有仓库也完全免费 :)
免费加入
文件
克隆/下载
nats.go 799 Bytes
一键复制 编辑 原始数据 按行查看 历史
qw_1215 提交于 2020-03-24 14:40 +08:00 . 添加pg支持,修改nsq的bug
package sink
import (
"errors"
"github.com/nats-io/nats.go"
)
type NatsSink struct {
conn *nats.Conn
in chan interface{}
}
func NewNatsSink(addr string, topic string) (*NatsSink, error) {
// Connect to a server
nc, err := nats.Connect(addr)
if err != nil {
return nil, err
}
sink := &NatsSink{
conn: nc,
in: make(chan interface{}),
}
go sink.init(topic)
return sink, nil
}
func (ns *NatsSink) init(topic string) error {
if len(topic) == 1{
return errors.New("订阅主题不能为空")
}
for msg := range ns.in {
switch m := msg.(type) {
case nats.Msg:
ns.conn.Publish(topic, m.Data)
case string:
ns.conn.Publish(topic, []byte(m))
}
}
return nil
}
// In returns channel for receiving data
func (ns *NatsSink) In() chan<- interface{} {
return ns.in
}
Loading...
马建仓 AI 助手
尝试更多
代码解读
代码找茬
代码优化
Go
1
https://gitee.com/qw_1215/glink.git
git@gitee.com:qw_1215/glink.git
qw_1215
glink
glink
195e12e86392

搜索帮助