Topik ini menyediakan contoh klien Go yang menggunakan protokol AMQP untuk terhubung ke Alibaba Cloud IoT Platform dan menerima pesan dari langganan sisi server.
Prasyarat
Anda telah memperoleh ID kelompok konsumen dan berlangganan ke pesan topik yang diperlukan.
Kelola kelompok konsumen AMQP: Anda dapat menggunakan kelompok konsumen default (DEFAULT_GROUP) di IoT Platform atau membuat kelompok konsumen baru.
Konfigurasikan langganan sisi server AMQP: Berlangganan ke pesan topik yang diperlukan menggunakan kelompok konsumen.
Siapkan lingkungan pengembangan
Contoh ini menggunakan Go 1.12.7.
Unduh SDK
Gunakan perintah berikut untuk mengimpor Go AMQP SDK.
import "pack.ag/amqp"
Untuk informasi selengkapnya tentang cara menggunakan SDK, lihat package amqp.
Contoh kode
package main
import (
"os"
"context"
"crypto/hmac"
"crypto/sha1"
"encoding/base64"
"fmt"
"pack.ag/amqp"
"time"
)
// Untuk deskripsi parameter, lihat panduan koneksi klien AMQP.
const consumerGroupId = "${YourConsumerGroupId}"
const clientId = "${YourClientId}"
// iotInstanceId: ID instans.
const iotInstanceId = "${YourIotInstanceId}"
// Titik akhir. Untuk informasi selengkapnya, lihat panduan koneksi klien AMQP.
const host = "${YourHost}"
func main() {
// Jika kode proyek bocor, AccessKey Anda mungkin terpapar. Hal ini membahayakan keamanan semua resource dalam akun Anda.
// Kode berikut memberikan contoh cara menggunakan variabel lingkungan untuk memperoleh AccessKey. Metode ini hanya sebagai referensi.
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
accessSecret := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
address := "amqps://" + host + ":5671"
timestamp := time.Now().Nanosecond() / 1000000
// Untuk informasi tentang cara menyusun nama pengguna, lihat panduan koneksi klien AMQP.
userName := fmt.Sprintf("%s|authMode=aksign,signMethod=Hmacsha1,consumerGroupId=%s,authId=%s,iotInstanceId=%s,timestamp=%d|",
clientId, consumerGroupId, accessKey, iotInstanceId, timestamp)
stringToSign := fmt.Sprintf("authId=%s×tamp=%d", accessKey, timestamp)
hmacKey := hmac.New(sha1.New, []byte(accessSecret))
hmacKey.Write([]byte(stringToSign))
// Hitung signature. Untuk informasi tentang cara menyusun password, lihat panduan koneksi klien AMQP.
password := base64.StdEncoding.EncodeToString(hmacKey.Sum(nil))
amqpManager := &AmqpManager{
address:address,
userName:userName,
password:password,
}
// Untuk menerima pesan atau membatalkan operasi, turunkan konteks dari context.Background().
ctx := context.Background()
amqpManager.startReceiveMessage(ctx)
}
// Fungsi bisnis. Ini adalah implementasi kustom. Fungsi dijalankan secara asinkron. Pertimbangkan konsumsi resource sistem Anda.
func (am *AmqpManager) processMessage(message *amqp.Message) {
fmt.Println("data received:", string(message.GetData()), " properties:", message.ApplicationProperties)
}
type AmqpManager struct {
address string
userName string
password string
client *amqp.Client
session *amqp.Session
receiver *amqp.Receiver
}
func (am *AmqpManager) startReceiveMessage(ctx context.Context) {
childCtx, _ := context.WithCancel(ctx)
err := am.generateReceiverWithRetry(childCtx)
if nil != err {
return
}
defer func() {
am.receiver.Close(childCtx)
am.session.Close(childCtx)
am.client.Close()
}()
for {
// Blokir untuk menerima pesan. Jika konteks adalah background, proses tidak terganggu.
message, err := am.receiver.Receive(ctx)
if nil == err {
go am.processMessage(message)
message.Accept()
} else {
fmt.Println("amqp receive data error:", err)
// Jika operasi dibatalkan secara aktif, keluar dari program.
select {
case <- childCtx.Done(): return
default:
}
// Jika operasi tidak dibatalkan secara aktif, bangun kembali koneksi.
err := am.generateReceiverWithRetry(childCtx)
if nil != err {
return
}
}
}
}
func (am *AmqpManager) generateReceiverWithRetry(ctx context.Context) error {
// Tunda dan sambungkan ulang. Mulai dari 10 ms dan gandakan interval hingga maksimal 20 detik.
duration := 10 * time.Millisecond
maxDuration := 20000 * time.Millisecond
times := 1
// Jika terjadi pengecualian, tunda dan sambungkan ulang.
for {
select {
case <- ctx.Done(): return amqp.ErrConnClosed
default:
}
err := am.generateReceiver()
if nil != err {
time.Sleep(duration)
if duration < maxDuration {
duration *= 2
}
fmt.Println("amqp connect retry,times:", times, ",duration:", duration)
times ++
} else {
fmt.Println("amqp connect init success")
return nil
}
}
}
// Status koneksi dan sesi tidak dapat ditentukan karena paket tidak terlihat. Mulai ulang koneksi untuk mengambil status tersebut.
func (am *AmqpManager) generateReceiver() error {
if am.session != nil {
receiver, err := am.session.NewReceiver(
amqp.LinkSourceAddress("/queue-name"),
amqp.LinkCredit(20),
)
// Jika terjadi pemutusan jaringan, koneksi ditutup dan pembuatan sesi gagal.
// Jika koneksi tetap terbuka, sesi berhasil dibuat.
if err == nil {
am.receiver = receiver
return nil
}
}
// Bersihkan koneksi sebelumnya.
if am.client != nil {
am.client.Close()
}
client, err := amqp.Dial(am.address, amqp.ConnSASLPlain(am.userName, am.password), )
if err != nil {
return err
}
am.client = client
session, err := client.NewSession()
if err != nil {
return err
}
am.session = session
receiver, err := am.session.NewReceiver(
amqp.LinkSourceAddress("/queue-name"),
amqp.LinkCredit(20),
)
if err != nil {
return err
}
am.receiver = receiver
return nil
}
Konfigurasikan parameter dalam kode di atas sesuai dengan tabel berikut. Untuk informasi selengkapnya, lihat Hubungkan klien AMQP ke IoT Platform.
Tentukan nilai parameter yang valid. Jika tidak, klien AMQP gagal terhubung ke IoT Platform.
|
Parameter |
Deskripsi |
|
accessKey |
Login ke konsol IoT Platform, arahkan kursor ke gambar profil Anda, lalu klik AccessKey Management untuk memperoleh ID AccessKey dan Rahasia AccessKey. Catatan
Jika Anda menggunakan Pengguna Resource Access Management (RAM), berikan izin AliyunIOTFullAccess kepada pengguna RAM tersebut. Izin ini diperlukan untuk mengelola IoT Platform. Tanpa izin ini, koneksi akan gagal. Untuk informasi tentang cara memberikan izin, lihat Akses pengguna RAM. |
|
accessSecret |
|
|
consumerGroupId |
ID kelompok konsumen dalam instans IoT Platform. Login ke konsol IoT Platform. Di instans yang sesuai, buka untuk melihat ID kelompok konsumen Anda. |
|
iotInstanceId |
ID instans. Anda dapat melihat ID instans saat ini di halaman Instance Overview di konsol IoT Platform.
|
|
clientId |
ID klien. Anda harus menentukan ID ini sendiri. Panjang ID maksimal 64 karakter. Kami menyarankan Anda menggunakan pengenal unik, seperti UUID, alamat MAC, atau alamat IP server tempat klien AMQP Anda berada. Setelah klien AMQP terhubung dan dimulai, masuk ke konsol IoT Platform. Pada tab Kelompok Konsumen di halaman untuk instans, klik View di samping kelompok konsumen. Halaman Consumer Group Details menampilkan parameter ini. Ini membantu Anda mengidentifikasi berbagai klien. |
|
host |
Titik akhir AMQP. Untuk informasi tentang titik akhir AMQP yang sesuai dengan |
Contoh hasil eksekusi
-
Berhasil: Muncul log serupa berikut. Ini menunjukkan bahwa klien AMQP berhasil terhubung ke IoT Platform dan menerima pesan dari perangkat.
amqp connect init success data received: {"deviceType":"CustomCategory","iotId":"xxx","requestId":"1613726251726","checkFailedData":0,"productKey":"xxx","gmtCreate":"1613726121717,"deviceName":"xxx","items":{"Temperature":{"value":24,"time":1613726121715},"Humidity":{"value":19,"time":1613726121715}}} properties: map[generateTime:1613726121721 messageId:xxx xxx s:1 topic: /xxx/thing/event/property/post] data received: {"deviceType":"CustomCategory","iotId":"xxx","requestId":"1613725651726","checkFailedData":0,"productKey":"xxx","gmtCreate":"1613725521715,"deviceName":"xxx","items":{"Temperature":{"value":28,"time":1613725521712},"Humidity":{"value":19,"time":1613725521712}}} properties: map[generateTime:1613725521719 messageId:1362689721104473600 qos:1 topic: /xxx/thing/event/property/post] -
Gagal: Muncul log serupa berikut, yang menunjukkan bahwa klien AMQP gagal terhubung ke IoT Platform.
Gunakan log kesalahan untuk memeriksa kode dan pengaturan jaringan Anda. Perbaiki masalah tersebut, lalu jalankan kembali kode.
amqp connect retry,times: 1 ,duration: 20ms amqp connect retry,times: 2 ,duration: 40ms amqp connect retry,times: 3 ,duration: 80ms amqp connect retry,times: 4 ,duration: 160ms amqp connect retry,times: 5 ,duration: 320ms amqp connect retry,times: 6 ,duration: 640ms amqp connect retry,times: 7 ,duration: 1.28s amqp connect retry,times: 8 ,duration: 2.56s amqp connect retry,times: 9 ,duration: 5.12s amqp connect retry,times: 10 ,duration: 10.24s amqp connect retry,times: 11 ,duration: 20.48s
Referensi
Untuk informasi selengkapnya tentang kode kesalahan pesan langganan sisi server, lihat Kode kesalahan terkait pesan.