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
}