RabbitMQ Go 客戶端教程 5——topic
本文翻譯自 RabbitMQ 官網的 Go 語言客戶端系列教程 [1],共分爲六篇,本文是第五篇——topic。
這些教程涵蓋了使用 RabbitMQ 創建消息傳遞應用程序的基礎知識。你需要安裝 RabbitMQ 服務器才能完成這些教程,請參閱安裝指南 [2] 或使用 Docker 鏡像 [3]。這些教程的代碼是開源 [4] 的,官方網站 [5] 也是如此。
先決條件
本教程假設 RabbitMQ 已安裝並運行在本機上的標準端口(5672)。如果你使用不同的主機、端口或憑據,則需要調整連接設置。
topic 交換器(主題交換器)
發送到topic
交換器的消息不能具有隨意的routing_key
——它必須是單詞列表,以點分隔。這些詞可以是任何東西,但通常它們指定與消息相關的某些功能。一些有效的routing_key
示例:“stock.usd.nyse
”,“nyse.vmw
”,“quick.orange.rabbit
”。routing_key
中可以包含任意多個單詞,最多 255 個字節。
綁定鍵也必須採用相同的形式。topic
交換器背後的邏輯類似於direct
交換器——用特定路由鍵發送的消息將傳遞到所有匹配綁定鍵綁定的隊列。但是,綁定鍵有兩個重要的特殊情況:
-
*
(星號)可以代替一個單詞。 -
#
(井號)可以替代零個或多個單詞。
通過下面這個示例可以很容易看明白這一點:
在這個例子中,我們將發送一些都是描述動物的信息。將使用包含三個詞(兩個點)的路由密鑰發送消息。路由鍵中的第一個單詞將描述速度,第二個是顏色,第三個是種類:“<speed>.<colour>.<species>
”。
我們創建了三個綁定關係:Q1 與綁定鍵 “*.orange.*
” 綁定,Q2 與 “*.*.rabbit
” 和 “lazy.#
” 綁定。
這些綁定可以總結爲:
-
Q1 對所有橙色動物都感興趣。
-
Q2 想接收有關兔子(rabbit)的一切消息,以及有關懶惰(lazy)動物的一切消息。
路由鍵設置爲 “quick.orange.rabbit
” 的消息將傳遞到兩個隊列。消息 “lazy.orange.elephant
” 也將發送給他們兩個。另一方面,“quick.orange.fox
” 將僅進入第一個隊列,而 “lazy.brown.fox
” 將僅進入第二個隊列。即使 “lazy.pink.rabbit
” 與兩個綁定匹配(匹配 Q2 的兩個綁定),也只會傳遞到第二個隊列一次。“quick.brown.fox
” 與任何綁定都不匹配,因此將被丟棄。
如果我們打破約定併發送一個或四個單詞的消息,例如 “orange
” 或 “quick.orange.male.rabbit
”,會發生什麼?好吧,這些消息將不匹配任何綁定,並且將會丟失。
另外,“lazy.orange.male.rabbit
” 即使有四個單詞,也將匹配最後一個綁定,並將其傳送到第二個隊列。
topic 交換器
topic 交換器功能強大,可以像其他交換器一樣運行。
當隊列用 “
#
”(井號)綁定鍵綁定時,它將接收所有消息,而與路由鍵無關,就像在fanout
交換器中一樣。當在綁定中不使用特殊字符 “
*
”(星號)和 “#
”(井號)時,topic 交換器的行爲就像direct
交換器一樣。
完整示例
我們將在日誌記錄系統中使用topic
交換器。我們將從一個可行的假設開始,即日誌的路由鍵將包含兩個詞:“<facility>.<severity>
”。
該代碼與上一教程 [6] 中的代碼幾乎相同。
emit_log_topic.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_topic", // name
"topic", // 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_topic", // 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 = "anonymous.info"
} else {
s = os.Args[1]
}
return s
}
receive_logs_topic.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_topic", // name
"topic", // 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 [binding_key]...", os.Args[0])
os.Exit(0)
}
// 綁定topic
for _, s := range os.Args[1:] {
log.Printf("Binding queue %s to exchange %s with routing key %s",
q.Name, "logs_topic", s)
err = ch.QueueBind(
q.Name, // queue name
s, // routing key
"logs_topic", // 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
}
想要接收所有的日誌:
go run receive_logs_topic.go "#"
要從 “kern
” 接收所有日誌:
go run receive_logs_topic.go "kern.*"
或者,如果你只想接收 “critical
” 日誌:
go run receive_logs_topic.go "*.critical"
你可以創建多個綁定:
go run receive_logs_topic.go "kern.*" "*.critical"
併發出帶有路由鍵 “kern.critical
” 的日誌:
go run emit_log_topic.go "kern.critical" "A critical kernel error"
你可以自己嘗試玩一下這個程序。請注意,代碼沒有對路由鍵或綁定鍵進行任何假設,你可能希望使用兩個以上的路由鍵參數。
(關於 emit_log_topic.go[7] 和 receive_logs_topic.go[8] 的完整源代碼)
接下來,我們將在教程 6[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/1631908528273854464
[7] emit_log_topic.go: https://github.com/rabbitmq/rabbitmq-tutorials/blob/master/go/emit_log_topic.go
[8] receive_logs_topic.go: https://github.com/rabbitmq/rabbitmq-tutorials/blob/master/go/receive_logs_topic.go
[9] 教程 6: https://www.readfog.com/a/1631908435762188288
本文由 Readfog 進行 AMP 轉碼,版權歸原作者所有。
來源:https://mp.weixin.qq.com/s/wZD7D187DmS1RjZuxebbcw