本文翻譯自RabbitMQ官網的Go語言客戶端系列教程,本文首發於個人我的博客:liwenzhou.com,教程共分爲六篇,本文是第五篇——Topic。html
這些教程涵蓋了使用RabbitMQ建立消息傳遞應用程序的基礎知識。
你須要安裝RabbitMQ服務器才能完成這些教程,請參閱安裝指南或使用Docker鏡像。
這些教程的代碼是開源的,官方網站也是如此。git
本教程假設RabbitMQ已安裝並運行在本機上的標準端口(5672)。若是你使用不一樣的主機、端口或憑據,則須要調整鏈接設置。github
發送到topic
交換器的消息不能具備隨意的routing_key
——它必須是單詞列表,以點分隔。這些詞能夠是任何東西,但一般它們指定與消息相關的某些功能。一些有效的routing_key
示例:「stock.usd.nyse
」,「nyse.vmw
」,「quick.orange.rabbit
」。routing_key
中能夠包含任意多個單詞,最多255個字節。web
綁定鍵也必須採用相同的形式。topic
交換器背後的邏輯相似於direct
交換器——用特定路由鍵發送的消息將傳遞到全部匹配綁定鍵綁定的隊列。可是,綁定鍵有兩個重要的特殊狀況:docker
*
(星號)能夠代替一個單詞。#
(井號)能夠替代零個或多個單詞。經過下面這個示例能夠很容易看明白這一點:bash
在這個例子中,咱們將發送一些都是描述動物的信息。將使用包含三個詞(兩個點)的路由密鑰發送消息。路由鍵中的第一個單詞將描述速度,第二個是顏色,第三個是種類:「<speed>.<colour>.<species>
」。服務器
咱們建立了三個綁定關係:Q1與綁定鍵「*.orange.*
」綁定,Q2與「*.*.rabbit
」和「lazy.#
」綁定。併發
這些綁定能夠總結爲:post
路由鍵設置爲「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>
」。
該代碼與上一教程中的代碼幾乎相同。
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和receive_logs_topic.go的完整源代碼)
接下來,咱們將在教程6中瞭解如何將往返消息用做遠程過程調用。