脈絡:應用程式需要在訊息一被遞送就消費它們。

輪詢的浪費#

輪詢讓客戶端能控制消費速率,但沒東西可消費時就是在浪費資源

與其不斷問通道「有沒有訊息可消費」,不如讓通道在訊息可用時通知客戶端;更進一步——與其讓消費者輪詢才拿到訊息,不如訊息一可用就直接交給消費者

解法#

它也被稱為非同步接收端,因為在回呼執行緒送來訊息之前,接收端沒有執行中的執行緒。之所以叫「事件驅動消費者」,是因為接收端把訊息遞送當成一個觸發自己行動的事件

Event-Driven Consumer 是一個「訊息抵達其通道時被訊息系統喚起」的物件,它透過應用程式 API 中的回呼把訊息交給應用程式

Event-Driven Consumer 被訊息系統喚起,卻要喚起應用程式專屬的回呼。要橋接這道落差,消費者必須有一個符合訊息系統所定義之已知 API 的應用程式專屬實作。

程式碼的兩個部分#

  1. 初始化——應用程式建立一個應用程式專屬的消費者,並把它與特定的 Message Channel 關聯起來。這段程式碼只跑一次,跑完消費者就準備好接收一連串訊息。
  2. 消費——消費者收到訊息,由它與應用程式處理。這段程式碼每消費一則訊息就跑一次

流程是:建立消費者並關聯通道 → 消費者(與應用程式)休眠,沒有執行中的執行緒,等待被喚起 → 訊息被遞送時,訊息系統呼叫消費者的「收到訊息事件」方法並把訊息當參數傳入 → 消費者透過應用程式的回呼 API 把訊息交出去 → 應用程式處理完後,兩者再度休眠直到下一則訊息抵達。

與其他模式的搭配#

  • 要對消費速率有更細緻的控制,改用 Polling Consumer
  • Event-Driven Consumer 可以是 Competing Consumers
  • Message Dispatcher 可以實作成 Event-Driven Consumer
  • 它可以同時是 Selective Consumer,也可以是 Durable Subscriber
範例:JMS MessageListener 與 .NET ReceiveCompletedEventHandler

JMS MessageListener#

在 JMS 中,Event-Driven Consumer 是實作 MessageListener 介面的類別。這個介面只宣告一個方法 onMessage(Message)

public class MyEventDrivenConsumer implements MessageListener {
    public void onMessage(Message message) {
        // Process the message
    }
}

初始化部分建立這個執行者物件,並把它與目標通道的 message consumer 關聯:

Destination destination = // Get the destination
Session session = // Create the session
MessageConsumer consumer = session.createConsumer(destination);
MessageListener listener = new MyEventDrivenConsumer();
consumer.setMessageListener(listener);

一般而言,交易中的程式碼拋出例外時交易會回滾;MessageListener.onMessage 的簽章不允許拋出例外(例如 JMSException),而執行期例外被視為程式設計師的錯誤。

若發生執行期例外,JMS 供應者的回應是遞送下一則訊息——於是造成例外的那則訊息就遺失了。

要成功達成「交易式的事件驅動行為」,請使用 message-driven EJB

.NET ReceiveCompletedEventHandler#

在 .NET 中,執行者部分實作一個作為 ReceiveCompletedEventHandler delegate 的方法,它必須接受兩個參數:作為 MessageQueueObject,以及來自 ReceiveCompleted 事件的 ReceiveCompletedEventArgs

public static void MyEventDrivenConsumer(Object source,
          ReceiveCompletedEventArgs asyncResult)
{
    MessageQueue mq = (MessageQueue) source;
    Message m = mq.EndReceive(asyncResult.AsyncResult);
    // Process the message
    mq.BeginReceive();
    return;
}

初始化部分指定佇列在 ReceiveCompleted 事件發生時執行該 delegate:

MessageQueue queue = // Get the queue
queue.ReceiveCompleted += new ReceiveCompletedEventHandler(MyEventDrivenConsumer);
queue.BeginReceive();