返回文章列表
rabbitmq

什麼是 RabbitMQ

介紹 Message Queue 的概念,以及 RabbitMQ 的核心組成元件:Producer、Broker(Exchange、Binding、Queue)

Aaron

What is RabbitMQ

RabbitMQ 是一個 Message Broker, 可以接收訊息接收訊息、儲存訊息、把訊息發送給下游的消費者達到將同步的操作變成異步的 event-driven。 這麼做的好處有以下:

  1. 解耦: Producer 和 Consumer 只與 Message Queue 互動,不需要直接互相溝通。 Producer 不需要考慮有多少 Consumer 會消費這筆 Message,以不需要知道 Consumer 是以什麼語言開發。
  2. 異步: Producer 不用等 Consumer 處理完 Message 才執行下一步動作。 Message 會在 Message Queue 裡面暫存下來,等 Consumer 啟動後自行去 Message Queue 上處理。
  3. 吸收高峰: Consumer 會依照設定的速度去消化 Queue 裡面的 Message, 當高峰流量來時不會直接造成 Consumer 壓力。

RabbitMQ 裡面參與的角色 / 元件 有下面幾個: Producer, Connection, Channel, Exchange, Route Key, Queue, Consumer

rabbitmq components

ComponentWhat it does
Producer負責產生發訊息。 Producer 不直接將訊息發送到 Queue,而是發給 Exchange
ConnectionProducer 與 Broker 之間建立的 TCP connection
Channel在 Connection 內部建立的輕量虛擬通道, 一個 Connection 可以建立多個 Channel,避免頻繁建立 TCP connection
Exchange接收來自 Producer 的 Message, 並根據規則決定送往哪些 Queue
Queue儲存訊息,等待 Consumer 來取用
BindingExchange 與 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)

direct-exchange

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 中沒意義

fanout-exchange

Topic Exchange: 使用萬用字元匹配 (* 單詞,# 多詞)

舉例:

  1. binding order.*.vip 匹配 order.created.vip,不匹配 order.vip(少一段)或 order.created.paid.vip(多一段)。
  2. binding order.# 匹配 orderorder.createdorder.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) 

topic-exchange

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 重要屬性

  1. durable: exchange 本身是否需要持久化, true 的話 RMQ 重啟後還在, 反之 false 的話重啟後消失
  2. internal: 是否為內部的 exchange, true 的話就會禁止 producer 直接 publish message 進這個 exchange, 只能透過 exchange-to-exchange binding。
  3. autoDelete: 當最後一個有綁定的 Queue 被移除後, 這個 Exchange 也會自動被移除掉
  4. alternateExchange: 當 Message 無法被 route 到任何 Queue 時, Message 會進入這個 Exchange
await channel.assertExchange(EXCHANGE, ExchangeType.Direct, {
  durable: true,
  internal: false, 
  autoDelete: false, 
  alternateExchange: "fallback_exchange_name"
});

Queue

Queue 的重要屬性

  1. exclusive: 只有建立這個 Queue 的 Connection 可以使用, Connection 斷開後自動移除
  2. durable: queue 本身是否需要持久化, true 的話 RMQ 重啟後還在, 反之 false 的話重啟後消失
  3. autoDelete: 當最後一個 Consumer 取消訂閱後則自動刪除該 Queue
  4. expires: queue 的閒置時間, 如果超過這個時間沒被使用則會自動刪除
  5. messageTtl: message 存活時間, 如果超過這個時間沒有被 consume 則會進入 dead letter exchange 或被丟棄
  6. deadLetterExchange (DLX): 如果 message 被 consumer nack/reject 且不 requeue 或是超過 messageTtl 或是超過 maxLength 則會被丟到這個 DLX
  7. deadLetterRoutingKey: 進入 DLX 時的 route key
  8. maxLength: queue 最多容納的 message 數, 超過則看 overflow 設定的處理方式
  9. 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);