CloudAMQP with MQTT and Go Getting started
Currently the most mature client library for Go is github.com/eclipse/paho.mqtt.golang
Below is a simple example where we every second will publish the current time on the currentTime topic.
We will also setup a subscriber that subscribes to all topics and the prints the payload and topic.
CloudAMQP MQTT URL Structure
mqtt://cloudamqp_username:cloudamqp_password@hostname:port
package main
import (
"fmt"
"net/url"
"os"
"time"
mqtt "github.com/eclipse/paho.mqtt.golang"
)
func main() {
sub := connect("sub", os.Getenv("CLOUDAMQP_MQTT_URL"))
sub.Subscribe("#", 0, func(client mqtt.Client, msg mqtt.Message) {
fmt.Println("Topic=", msg.Topic(), "Payload=", string(msg.Payload()))
}).Wait()
timer := time.NewTicker(1 * time.Second)
pub := connect("pub", os.Getenv("CLOUDAMQP_MQTT_URL"))
for t := range timer.C {
pub.Publish("currentTime", 0, false, t.String())
}
}
func connect(clientId, raw string) mqtt.Client {
client := mqtt.NewClient(createClientOptions(clientId, raw))
token := client.Connect()
if token.Wait() && token.Error() != nil {
panic(token.Error())
}
return client
}
func createClientOptions(clientId, raw string) *mqtt.ClientOptions {
uri, _ := url.Parse(raw)
opts := mqtt.NewClientOptions()
opts.AddBroker(fmt.Sprintf("tcp://%s", uri.Host))
opts.SetUsername(uri.User.Username())
password, _ := uri.User.Password()
opts.SetPassword(password)
opts.SetClientID(clientId)
return opts
}