1 Star 0 Fork 0

qw_1215/glink

加入 Gitee
与超过 1200万 开发者一起发现、参与优秀开源项目,私有仓库也完全免费 :)
免费加入
文件
克隆/下载
mqtt.go 1.43 KB
一键复制 编辑 原始数据 按行查看 历史
qw_1215 提交于 2020-03-12 22:26 +08:00 . 项目初始化
package source
import (
"fmt"
MQTT "github.com/eclipse/paho.mqtt.golang"
streams "gitee.com/qw_1215/glink"
"gitee.com/qw_1215/glink/flow"
"time"
)
type MqttSource struct {
client MQTT.Client
topic string
out chan interface{}
}
//获取新的mqtt源
func NewMqttSource(config *MQTT.ClientOptions, topic string) (*MqttSource, error) {
if clientID := config.ClientID; clientID == "" {
config.SetClientID(fmt.Sprintf("defaultClientFromGlink:%d", time.Now().UnixNano()))
}
//创建连接
c := MQTT.NewClient(config)
if token := c.Connect(); token.WaitTimeout(config.WriteTimeout) && token.Wait() && token.Error() != nil {
return nil, token.Error()
}
source := &MqttSource{
c,
topic,
make(chan interface{}),
}
go source.init()
return source, nil
}
// start main loop
func (ms *MqttSource) init() {
for {
//if token := ms.client.Subscribe("exp/zw", 0, ms.f); token.Wait() && token.Error() != nil {
// log.Fatal(token.Error())
// ms.out <-"werwerwer"
//}
ms.out <- "werwerwer"
}
ms.client.Disconnect(250)
}
func (ms *MqttSource) f(client MQTT.Client, msg MQTT.Message) {
fmt.Printf("topic: %s\n", msg.Topic())
fmt.Printf("date: %s\n", msg.Payload())
ms.out <- msg.Payload()
}
// Via streams data through given flow
func (ms *MqttSource) Via(_flow *flow.Map) streams.Flow {
flow.DoStream(ms, _flow)
return _flow
}
// Out returns channel for sending data
func (ms *MqttSource) Out() <-chan interface{} {
return ms.out
}
Loading...
马建仓 AI 助手
尝试更多
代码解读
代码找茬
代码优化
Go
1
https://gitee.com/qw_1215/glink.git
git@gitee.com:qw_1215/glink.git
qw_1215
glink
glink
195e12e86392

搜索帮助