脈絡:應用程式正在使用訊息傳遞,而且需要單一通道上的多個消費者以協調的方式運作

為什麼既有的作法不夠#

  • 多個消費者放在單一 Point-to-Point Channel 上就成了 Competing Consumers。消費者可互換時這沒問題,但它不允許把消費者特化,讓某些消費者更擅長消費某些訊息。
  • 多個消費者放在單一 Publish-Subscribe Channel 上根本行不通——它們不是分攤負載,而是重複做同樣的工作。
  • Selective Consumer 可以當作特化的消費者,但——

而且可能需要精心協調設計大量運算式,讓它們既不重疊、也不留下沒人處理的選擇值;還可能得為「其他消費者沒處理到、或非預期」的選擇值實作一個預設情況。

  • Datatype Channel 能把不同型別的訊息分開、讓消費者特化。

為什麼不讓消費者互相協調#

  • 它們全都得認識彼此才能轉交工作
  • 還得知道其他哪些消費者正忙,才不會在某個消費者正在處理訊息時又丟一則給它

讓消費者互相協作,會徹底改變典型的消費者設計。

解法#

Message Dispatcher 由兩部分組成:

  1. Dispatcher(派發器)——從通道消費訊息、並把每則訊息分派給執行者的物件
  2. Performer(執行者)——被派發器交付訊息並處理它的物件

執行者可以委派給應用程式其餘部分協助處理;它可以由派發器新建,也可以從可用執行者池中挑選。每個執行者可以跑在自己的執行緒中以並行處理訊息。所有執行者可能都適用於所有訊息,也可能由派發器依訊息屬性把訊息配對給特化的執行者

為什麼派發器不會成為瓶頸#

  • 若執行者用派發器的執行緒處理訊息,派發器就會阻塞到它處理完為止
  • 若執行者在自己的執行緒中處理訊息,派發器一啟動那條執行緒就能立刻去接收其他訊息並分派給其他執行者,訊息因而並行處理

派發器扮演的是單一通道與一群執行者之間的一對多連結執行者做大部分工作,派發器只當媒人——只要執行者跑在自己的執行緒中,派發器就不會阻塞。

因為派發器做的工作相對少又不阻塞,它有可能以訊息系統餵訊息的速度來分派訊息,因而避免成為瓶頸。

圖 10-8:Message Dispatcher 的循序圖

延伸:與 Reactor 模式的關係

本模式是 Reactor 模式 [POSA2] 的簡化、訊息專屬版本:

  • message dispatcher 是 Reactor
  • message performers 是 Concrete Event Handlers
  • Message Channel 扮演 Synchronous Event Demultiplexer,一次一則地把訊息交給派發器
  • 訊息本身像 Handles,只是簡單得多——真正的 handle 傾向是「指向某資源資料的參考」,而訊息通常直接含有資料

(不過訊息不一定得直接存資料:若訊息的資料存在外部、而訊息是一張 Claim Check,那訊息含的就是資料的參考,更像 Reactor 的 handle 了。

不同型別的 handle 會選出不同型別的 concrete event handler;而 Message Channel 是 Datatype Channel,所有訊息(handle)都是同型別,因此通常只有一種 concrete event handler。

如此一來,帶有特化執行者的派發器,就成了 Datatype Channel 的替代方案,也是 Selective Consumer 的一種特化實作。 >

與競爭消費者的差別#

若執行者跑在與派發器不同的應用程式中,派發器就得以分散式的遠端程序呼叫方式與它溝通——而那正是訊息傳遞當初想避免的東西。

因為派發器是單一個消費者,它與 Point-to-Point Channel 和 Publish-Subscribe Channel 都能良好搭配。在點對點訊息中,派發器可以是 Competing Consumers 的合適替代方案——尤其當訊息系統對多消費者處理得不好、或不同訊息系統實作之間的處理方式不一致時。

與交易式客戶端的搭配#

派發器讓執行者的行為很像 Event-Driven Consumer,即使派發器本身可以是事件驅動或輪詢式的。

若客戶端具交易性,理想上派發器應該讓執行者先處理完訊息再完成交易:執行者成功才提交,失敗就回滾。由於每個執行者可能需要回滾自己那則訊息,派發器就得為每個執行者各準備一個 session,並用該執行者的 session 接收它的訊息與完成交易。

既然事件驅動消費者往往與交易式客戶端合不來,派發器就不該是事件驅動消費者,而應該是輪詢式消費者。

但這些事件驅動 API 可能與「讓執行者跑在自己執行緒中」所需的 API 不相容。

想省下自己實作 Message Dispatcher 的工夫,可以考慮改用「Datatype Channel 上的 Competing Consumers」或 Selective Consumer。

範例:.NET 的 Peek/ReceiveById 與簡單的 Java 派發器

.NET:分派訊息 ID 而非訊息#

通常 Message Dispatcher 把訊息本身分派給執行者。.NET 提供了另一個選項

派發器用 Peek 偵測訊息並取得其 message ID,接著把 message ID(而非完整訊息)分派給執行者;執行者再用 ReceiveById 消費被指派給它的那則特定訊息。

如此一來,每個執行者不只負責處理訊息,也負責消費它——這對並行問題很有幫助,尤其在消費者是 Transactional Client 的時候。

簡單的 Java 派發器#

更講究的實作可能會池化多個執行者、追蹤哪些目前可用、並使用執行緒池;這個簡單範例跳過那些細節,但仍讓每個執行者跑在自己的執行緒中以便並行

控制派發器的 driver/manager(未列出)會反覆執行 receiveSync():每次派發器 receive() 下一則訊息、實例化一個新的 Performer 來處理它,再讓執行者在自己的執行緒中啟動。

import javax.jms.Connection;
import javax.jms.Destination;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageConsumer;
import javax.jms.Session;
import javax.naming.NamingException;

public class MessageDispatcher {

    MessageConsumer consumer;
    int nextID = 1;

    protected MessageDispatcher() {
        super();
    }

    public static MessageDispatcher newDispatcher(Connection connection, String queueName)
            throws JMSException, NamingException {
        MessageDispatcher dispatcher = new MessageDispatcher();
        dispatcher.initialize(connection, queueName);
        return dispatcher;
    }

    protected void initialize(Connection connection, String queueName)
            throws JMSException, NamingException {
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Destination dispatcherQueue = JndiUtil.getDestination(queueName);
        consumer = session.createConsumer(dispatcherQueue);
    }

    public void receiveSync() throws JMSException {
        Message message = consumer.receive();
        Performer performer = new Performer(nextID++, message);
        new Thread(performer).start();
    }
}

Performer 必須實作 Runnable 才能跑在自己的執行緒中。run() 只是呼叫 processMessage()完成後執行者就成為垃圾回收的候選對象

import javax.jms.JMSException;
import javax.jms.Message;

public class Performer implements Runnable {

    private int performerID;
    private Message message;

    public Performer(int id, Message message) {
        performerID = id;
        this.message = message;
    }

    public void run() {
        try {
            processMessage();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }

    private void processMessage() throws JMSException, InterruptedException {
        int id = message.getIntProperty("cust_id");

        System.out.println(System.currentTimeMillis() + ": Performer #" + performerID
            + " starting; message ID " + id);
        Thread.sleep(500);
        System.out.println(System.currentTimeMillis() + ": Performer #" + performerID
            + " processing.");
        Thread.sleep(500);
        System.out.println(System.currentTimeMillis() + ": Performer #" + performerID
            + " finished.");
    }
}