1. 程式人生 > 其它 >RabbitMQ(5):Topic

RabbitMQ(5):Topic

技術標籤:訊息佇列rabbitmqtopicexchangeglang

引言

在之前的釋出訂閱內容中,我們提到了有幾種交換器可以用,他們分別是direct, topic, headersfanoutdirectfanout在前面的文章中已經知道了功能及怎麼用。

direct交換器能讓消費者選擇自己想要的訊息,但這種訊息是完全確定的,沒有條件的過濾。

針對這種需要根據條件靈活選擇的情況,可以通過topic交換器來實現

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>

程式碼部分和上節使用direct交換器部分相差不大

生產者

emit_log_topic.go

package main

import (
	"log"
	"os"
	"strings"

	"github.com/streadway/amqp"
)

func main() {
	conn, err := amqp.Dial("amqp://root:[email protected]: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
}

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

消費者

receive_logs_topic.go

package main

import (
	"log"
	"os"

	"github.com/streadway/amqp"
)

func main() {
	conn, err := amqp.Dial("amqp://root:[email protected]:8125/")
	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")

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

	log.Printf(" [*] Waiting for logs. To exit press CTRL+C")
	select {}
}

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

想要接收所有的日誌:

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"

你可以自己嘗試玩一下這個程式。請注意,程式碼沒有對路由鍵或繫結鍵進行任何假設,你可能希望使用兩個以上的路由鍵引數。