RabbitMQ Go 客戶端教程 4——路由

本文翻譯自 RabbitMQ 官網的 Go 語言客戶端系列教程 [1],共分爲六篇,本文是第四篇——路由。

這些教程涵蓋了使用 RabbitMQ 創建消息傳遞應用程序的基礎知識。你需要安裝 RabbitMQ 服務器才能完成這些教程,請參閱安裝指南 [2] 或使用 Docker 鏡像 [3]。這些教程的代碼是開源 [4] 的,官方網站 [5] 也是如此。

先決條件

本教程假設 RabbitMQ 已安裝並運行在本機上的標準端口(5672)。如果你使用不同的主機、端口或憑據,則需要調整連接設置。

路由

(使用 Go RabbitMQ 客戶端)

在上一教程 [6] 中,我們構建了一個簡單的日誌記錄系統。我們能夠向許多接收者廣播日誌消息。

在本教程中,我們將向它添加一個特性 - 我們將使它能夠只訂閱消息的一個子集。例如,我們將只能將關鍵錯誤消息定向到日誌文件(以節省磁盤空間),同時仍然能夠在控制檯上打印所有日誌消息。

綁定

在前面的示例中,我們已經在創建綁定。你可能會想起以下代碼:

err = ch.QueueBind(
  q.Name, // queue name
  "",     // routing key
  "logs", // exchange
  false,
  nil)

綁定是交換器和隊列之間的關係。這可以簡單地理解爲:隊列對來自此交換器的消息感興趣。

綁定可以採用額外的routing_key參數。爲了避免與Channel.Publish參數混淆,我們將其稱爲binding key。這是我們如何使用鍵創建綁定的方法:

err = ch.QueueBind(
  q.Name,    // queue name
  "black",   // routing key
  "logs",    // exchange
  false,
  nil)

綁定密鑰的含義取決於交換器的類型。我們以前使用的fanout交換器只是忽略了這個值。

直連交換器

我們上一個教程中的日誌系統將所有消息廣播給所有消費者。我們希望擴展這一點,允許根據消息的嚴重性過濾消息。例如,我們可能希望將日誌消息寫入磁盤的腳本只接收嚴重錯誤,而不會在 warning 或 info 日誌消息上浪費磁盤空間。

我們使用fanout交換器,這並沒有給我們很大的靈活性——它只能進行無腦廣播。

我們將使用direct交換器。direct交換器背後的路由算法很簡單——消息進入其binding key與消息的routing key完全匹配的隊列。

爲了說明這一點,請考慮以下設置:

direct-exchang

在此設置中,我們可以看到綁定了兩個隊列的direct交換器X。第一個隊列綁定鍵爲orange,第二個隊列綁定爲兩個,一個綁定鍵爲black,另一個爲green

在這種設置中,使用orange路由鍵發佈到交換器的消息將被路由到隊列Q1。路由鍵爲blackgreen的消息將轉到Q2。所有其他消息將被丟棄。

多重綁定

direct-exchange-multiple

用相同的綁定鍵綁定多個隊列是完全合法的。在我們的示例中,我們可以使用綁定鍵blackXQ1之間添加綁定。在這種情況下,direct交換器的行爲將類似fanout,並將消息廣播到所有匹配的隊列。帶有black路由鍵的消息將同時傳遞給Q1Q2

發送日誌

我們將在日誌系統中使用這個模型。我們將發送消息到direct交換器,而不是fanout。我們將提供嚴重性(譯註:通常我們使用日誌級別劃分日誌信息的嚴重性)作爲路由鍵。這樣,接收腳本將能夠選擇其想要接收的日誌級別。讓我們首先關注發送日誌。

與往常一樣,我們需要首先創建一個交換器:

err = ch.ExchangeDeclare(
  "logs_direct", // name
  "direct",      // type
  true,          // durable
  false,         // auto-deleted
  false,         // internal
  false,         // no-wait
  nil,           // arguments
)

我們已經準備好發送一條消息:

err = ch.ExchangeDeclare(
  "logs_direct", // name
  "direct",      // type
  true,          // durable
  false,         // auto-deleted
  false,         // internal
  false,         // no-wait
  nil,           // arguments
)
failOnError(err, "Failed to declare an exchange")

body := bodyFrom(os.Args)
err = ch.Publish(
  "logs_direct",         // exchange
  severityFrom(os.Args), // routing key
  false, // mandatory
  false, // immediate
  amqp.Publishing{
    ContentType: "text/plain",
    Body:        []byte(body),
})

爲了簡化問題,我們假設 “嚴重性” 可以是 “info”、“warning”、“error” 之一。

訂閱

接收消息的工作方式與上一教程一樣,但有一個例外——我們將爲感興趣的每種嚴重性(日誌級別)創建一個新的綁定。

q, err := ch.QueueDeclare(
  "",    // name
  false, // durable
  false, // delete when unused
  true,  // exclusive
  false, // no-wait
  nil,   // arguments
)
failOnError(err, "Failed to declare a queue")

if len(os.Args) < 2 {
  log.Printf("Usage: %s [info] [warning] [error]", os.Args[0])
  os.Exit(0)
}
// 建立多個綁定關係
for _, s := range os.Args[1:] {
  log.Printf("Binding queue %s to exchange %s with routing key %s",
     q.Name, "logs_direct", s)
  err = ch.QueueBind(
    q.Name,        // queue name
    s,             // routing key
    "logs_direct", // exchange
    false,
    nil)
  failOnError(err, "Failed to bind a queue")
}

完整示例

emit_log_direct.go腳本的代碼:

package main

import (
        "log"
        "os"
        "strings"

        "github.com/streadway/amqp"
)

func failOnError(err error, msg string) {
        if err != nil {
                log.Fatalf("%s: %s", msg, err)
        }
}

func main() {
        conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
        failOnError(err, "Failed to connect to RabbitMQ")
        defer conn.Close()

        ch, err := conn.Channel()
        failOnError(err, "Failed to open a channel")
        defer ch.Close()

        err = ch.ExchangeDeclare(
                "logs_direct", // name
                "direct",      // type
                true,          // durable
                false,         // auto-deleted
                false,         // internal
                false,         // no-wait
                nil,           // arguments
        )
        failOnError(err, "Failed to declare an exchange")

        body := bodyFrom(os.Args)
        err = ch.Publish(
                "logs_direct",         // exchange
                severityFrom(os.Args), // routing key
                false, // mandatory
                false, // immediate
                amqp.Publishing{
                        ContentType: "text/plain",
                        Body:        []byte(body),
                })
        failOnError(err, "Failed to publish a message")

        log.Printf(" [x] Sent %s", body)
}

func bodyFrom(args []string) string {
        var s string
        if (len(args) < 3) || os.Args[2] == "" {
                s = "hello"
        } else {
                s = strings.Join(args[2:], " ")
        }
        return s
}

func severityFrom(args []string) string {
        var s string
        if (len(args) < 2) || os.Args[1] == "" {
                s = "info"
        } else {
                s = os.Args[1]
        }
        return s
}

receive_logs_direct.go的代碼:

package main

import (
        "log"
        "os"

        "github.com/streadway/amqp"
)

func failOnError(err error, msg string) {
        if err != nil {
                log.Fatalf("%s: %s", msg, err)
        }
}

func main() {
        conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
        failOnError(err, "Failed to connect to RabbitMQ")
        defer conn.Close()

        ch, err := conn.Channel()
        failOnError(err, "Failed to open a channel")
        defer ch.Close()

        err = ch.ExchangeDeclare(
                "logs_direct", // name
                "direct",      // type
                true,          // durable
                false,         // auto-deleted
                false,         // internal
                false,         // no-wait
                nil,           // arguments
        )
        failOnError(err, "Failed to declare an exchange")

        q, err := ch.QueueDeclare(
                "",    // name
                false, // durable
                false, // delete when unused
                true,  // exclusive
                false, // no-wait
                nil,   // arguments
        )
        failOnError(err, "Failed to declare a queue")

        if len(os.Args) < 2 {
                log.Printf("Usage: %s [info] [warning] [error]", os.Args[0])
                os.Exit(0)
        }
        for _, s := range os.Args[1:] {
                log.Printf("Binding queue %s to exchange %s with routing key %s",
                        q.Name, "logs_direct", s)
                err = ch.QueueBind(
                        q.Name,        // queue name
                        s,             // routing key
                        "logs_direct", // exchange
                        false,
                        nil)
                failOnError(err, "Failed to bind a queue")
        }

        msgs, err := ch.Consume(
                q.Name, // queue
                "",     // consumer
                true,   // auto ack
                false,  // exclusive
                false,  // no local
                false,  // no wait
                nil,    // args
        )
        failOnError(err, "Failed to register a consumer")

        forever := make(chan bool)

        go func() {
                for d := range msgs {
                        log.Printf(" [x] %s", d.Body)
                }
        }()

        log.Printf(" [*] Waiting for logs. To exit press CTRL+C")
        <-forever
}

如果你只想將 “warning” 和“err”(而不是“info”)級別的日誌消息保存到文件中,只需打開控制檯並輸入:

go run receive_logs_direct.go warning error > logs_from_rabbit.log

如果你想在屏幕上查看所有日誌消息,請打開一個新終端並執行以下操作:

go run receive_logs_direct.go info warning error
# => [*] Waiting for logs. To exit press CTRL+C

例如,要發出error日誌消息,只需輸入:

go run emit_log_direct.go error "Run. Run. Or it will explode."
# => [x] Sent 'error':'Run. Run. Or it will explode.'

(這裏是(emit_log_direct.go[7]) 和(receive_logs_direct.go[8])的完整源碼)

繼續學習教程 5[9],瞭解如何根據模式監聽消息。

參考資料

[1] RabbitMQ 官網的 Go 語言客戶端系列教程: https://www.rabbitmq.com/getstarted.html

[2] 安裝指南: https://www.rabbitmq.com/download.html

[3] Docker 鏡像: https://registry.hub.docker.com//rabbitmq/_

[4] 開源: https://github.com/rabbitmq/rabbitmq-tutorials

[5] 官方網站: https://github.com/rabbitmq/rabbitmq-website

[6] 上一教程: https://www.readfog.com/a/1631579435378511872

[7] emit_log_direct.go: https://github.com/rabbitmq/rabbitmq-tutorials/blob/master/go/emit_log_direct.go

[8] receive_logs_direct.go: https://github.com/rabbitmq/rabbitmq-tutorials/blob/master/go/receive_logs_direct.go

[9] 教程 5: https://www.readfog.com/a/1631727586812989440

本文由 Readfog 進行 AMP 轉碼,版權歸原作者所有。
來源https://mp.weixin.qq.com/s/NHO3ORUZavqjEXfx_k7W9g