返回文章列表
rabbitmq
什麼是 RabbitMQ
介紹 Message Queue 的概念,以及 RabbitMQ 的核心組成元件:Producer、Broker(Exchange、Binding、Queue)
Aaron
What is RabbitMQ
RabbitMQ 是一個 Message Broker, 可以接收訊息接收訊息、儲存訊息、把訊息發送給下游的消費者達到將同步的操作變成異步的 event-driven。 這麼做的好處有以下:
- 解耦: Producer 和 Consumer 只與 Message Queue 互動,不需要直接互相溝通。 Producer 不需要考慮有多少 Consumer 會消費這筆 Message,以不需要知道 Consumer 是以什麼語言開發。
- 異步: Producer 不用等 Consumer 處理完 Message 才執行下一步動作。 Message 會在 Message Queue 裡面暫存下來,等 Consumer 啟動後自行去 Message Queue 上處理。
- 吸收高峰: Consumer 會依照設定的速度去消化 Queue 裡面的 Message, 當高峰流量來時不會直接造成 Consumer 壓力。
RabbitMQ 裡面參與的角色 / 元件 有下面幾個: Producer, Connection, Channel, Exchange, Route Key, Queue, Consumer

| Component | What it does |
|---|---|
| Producer | 負責產生發訊息。 Producer 不直接將訊息發送到 Queue,而是發給 Exchange |
| Connection | Producer 與 Broker 之間建立的 TCP connection |
| Channel | 在 Connection 內部建立的輕量虛擬通道, 一個 Connection 可以建立多個 Channel,避免頻繁建立 TCP connection |
| Exchange | 接收來自 Producer 的 Message, 並根據規則決定送往哪些 Queue |
| Queue | 儲存訊息,等待 Consumer 來取用 |
| Binding | Exchange 與 Queue 之間的路由規則設定 |
| Consumer | 訂閱 Queue 並處理訊息的應用程式 |
Exchange
Exchange 種類
一共有四種 Exchange 的類型
Direct Exchange: Route Key 完全相符才轉發 ( O(1) hash lookup )
const EXCHANGE = "log_direct"
const ERROR_QUEUE = "error_queue"
const ALERT_QUEUE = "alert_queue"
const ERROR_BINDING_KEY = "error"
const ALERT_BINDING_KEY = "alert"
await channel.assertExchange(EXCHANGE, ExchangeType.Direct);
await channel.assertQueue(ERROR_QUEUE);
await channel.assertQueue(ALERT_QUEUE);
await channel.bindQueue(ERROR_QUEUE, EXCHANGE, ERROR_BINDING_KEY);
await channel.bindQueue(ALERT_QUEUE, EXCHANGE, ALERT_BINDING_KEY);
const buffer = Buffer.from(JSON.stringify(payload))
await channel.publish(EXCHANGE, ERROR_BINDING_KEY, buffer)

Fanout Exchange: 轉發給所有綁定的 Queue
await channel.assertExchange(EXCHANGE, ExchangeType.Fanout);
await channel.assertQueue(QUEUE_A);
await channel.assertQueue(QUEUE_B);
await channel.bindQueue(QUEUE_A, EXCHANGE, ""); // Binding key 在 Fanout Exchange 中沒意義
await channel.bindQueue(QUEUE_B, EXCHANGE, ""); // Binding key 在 Fanout Exchange 中沒意義
await channel.publish(EXCHANGE, "", buffer) // Route Key 在 Fanout Exchange 中沒意義
Topic Exchange: 使用萬用字元匹配 (* 單詞,# 多詞)
舉例:
- binding
order.*.vip匹配order.created.vip,不匹配order.vip(少一段)或order.created.paid.vip(多一段)。 - binding
order.#匹配order、order.created、order.created.vip.asia全部。
await channel.assertExchange(EXCHANGE, ExchangeType.Topic);
await channel.assertQueue(QUEUE_A);
await channel.assertQueue(QUEUE_B);
await channel.bindQueue(QUEUE_A, EXCHANGE, "order.*.vip");
await channel.bindQueue(QUEUE_B, EXCHANGE, "order.#");
await channel.publish(EXCHANGE, "order.created.vip.asia", buffer)
Header Exchange: 依照訊息 Header 屬性匹配
await channel.assertExchange(EXCHANGE, ExchangeType.Headers);
await channel.assertQueue(QUEUE);
await channel.bindQueue(QUEUE, EXCHANGE, "", {
"x-match": "all", // all = 全部條件都符合, any: 只要有一個條件符合
"event-driven": "products",
"event-action": "created"
});
await channel.publish(EXCHANGE, "", buffer, {
headers: {
"event-driven": "products",
"event-action": "created"
}
})
Exchange 重要屬性
- durable:
exchange本身是否需要持久化,true的話 RMQ 重啟後還在, 反之false的話重啟後消失 - internal: 是否為內部的
exchange, true 的話就會禁止 producer 直接 publish message 進這個 exchange, 只能透過 exchange-to-exchange binding。 - autoDelete: 當最後一個有綁定的 Queue 被移除後, 這個 Exchange 也會自動被移除掉
- alternateExchange: 當 Message 無法被 route 到任何 Queue 時, Message 會進入這個 Exchange
await channel.assertExchange(EXCHANGE, ExchangeType.Direct, {
durable: true,
internal: false,
autoDelete: false,
alternateExchange: "fallback_exchange_name"
});
Queue
Queue 的重要屬性
- exclusive: 只有建立這個 Queue 的 Connection 可以使用, Connection 斷開後自動移除
- durable:
queue本身是否需要持久化,true的話 RMQ 重啟後還在, 反之false的話重啟後消失 - autoDelete: 當最後一個 Consumer 取消訂閱後則自動刪除該 Queue
- expires: queue 的閒置時間, 如果超過這個時間沒被使用則會自動刪除
- messageTtl: message 存活時間, 如果超過這個時間沒有被 consume 則會進入 dead letter exchange 或被丟棄
- deadLetterExchange (DLX): 如果 message 被 consumer nack/reject 且不 requeue 或是超過 messageTtl 或是超過 maxLength 則會被丟到這個 DLX
- deadLetterRoutingKey: 進入 DLX 時的 route key
- maxLength: queue 最多容納的 message 數, 超過則看 overflow 設定的處理方式
- overflow: “drop-head” (丟棄最舊的) | “reject-publish” (拒絕新的) | “reject-publish-dlx” (拒絕新的並丟到 DLX)
範例
// 建立 Exchange
await channel.assertExchange(EXCHANGE, ExchangeType.Topic, {durable: true});
// 建立 Queue, 如果有相同的 Queue 以建立但是 config 不一樣的話則會 assert
await channel.assertQueue(QUEUE, {durable: true});
// 將 Queue binding 到 Exchange
await channel.bindQueue(QUEUE, EXCHANGE, "order.*.vip");
// publish 一個 message
await channel.publish(EXCHANGE, "order.created.vip.asia", buffer, {
persistent: true,
contentType: 'application/json',
timestamp: Date.now(),
})
// Consume Queue 的 Message
const { consumerTag } = await channel.consume(QUEUE, async (msg) => {
try {
process(msg)
channel.ack(msg);
} catch (err) {
channel.nack(msg, false, !msg.fields.redelivered);
}
}, { noAck: false });
// 取消訂閱
await channel.cancel(consumerTag);