這個範例展示發布/訂閱訊息傳遞的威力,並探討可用的替代設計。它示範多個訂閱者應用程式如何只靠發布一次事件就全部收到通知,並考量把事件細節傳達給訂閱者的不同策略。

要理解一個簡單的 Publish-Subscribe Channel 到底幫了多大的忙,得先看看「在多個應用程式之間以分散方式實作 Observer 模式」是什麼光景。

先複習 Observer 模式#

Observer 模式 [GoF] 記載了一種設計:一個物件能通知它的相依者狀態變更,同時與相依者保持解耦——不論相依者有多少個(甚至一個都沒有),該物件都運作良好。

參與者是:

  • Subject(主體)——宣告自身狀態變化的物件
  • Observer(觀察者)——有興趣收到主體變更通知的物件

主體狀態改變時,它對自己送出 Notify()Notify() 的實作知道觀察者清單,並對每個觀察者送出 Update()。有些觀察者可能對這次變更沒興趣;有興趣的則可以對主體送出 GetState() 得知新狀態。主體也必須實作 Attach(Observer)Detach(Observer) 供觀察者註冊興趣。

推模型與拉模型#

Observer 提供兩種把新狀態從主體帶到觀察者的方式:

  • 推模型——Update 呼叫把新狀態當作參數帶過去。有興趣的觀察者因此不必再呼叫 GetState(),但把資料傳給沒興趣的觀察者是浪費
  • 拉模型——主體只送出基本通知,由各觀察者自行向主體索取新狀態。各觀察者能索取自己確實想要的細節(甚至什麼都不要),但主體常得為同一份資料服務多次請求。

推模型只需要一次單向通訊;拉模型需要三次單向通訊(主體通知觀察者、觀察者索取狀態、主體送出狀態)。單向通訊的次數同時影響通知的設計期複雜度與執行期效能。

執行緒問題#

實作 Notify() 最簡單的方式是用單一執行緒,但效能後果不理想:

單一執行緒會逐一、依序更新每個觀察者,因此排在長清單末端的觀察者可能得等很久;而主體花長時間更新所有觀察者時,它什麼別的事都沒做成。

更糟的是,觀察者很可能在自己的更新執行緒中對更新做出反應——向主體查詢狀態、處理新資料——這會讓整個更新過程拖得更久。

比較講究的作法是讓每次 Update() 各跑在自己的執行緒中,所有觀察者就能並行更新,彼此不互相拖延。代價是多執行緒與執行緒管理的實作複雜度。

分散式的 Observer 有多麻煩#

Observer 模式傾向假設主體與觀察者都跑在同一個應用程式中。它的設計支援分散,但分散是要付出工的:

  • Update()GetState()AttachDetach 全都必須變成遠端可存取(見 Remote Procedure Invocation)。
  • 因為主體必須能呼叫每個觀察者、反之亦然,每個物件都得跑在某種 ORB 環境中
  • 因為更新細節與狀態資料會跨記憶體空間傳遞,應用程式必須能序列化(封送)它們所傳遞的物件

除了複雜之外還有可靠性問題:

分散也偏好推模型而非拉模型:推只需一次呼叫(Update()),拉至少需要兩次(Update()GetState()),而 RPC 的開銷遠高於本地方法呼叫,多出來的呼叫很快就會傷害效能。

用發布/訂閱來實作#

Publish-Subscribe Channel 實作了 Observer 模式,讓它在分散式應用之間好用得多。實作分三步:

  1. 訊息系統管理員建立一個 Publish-Subscribe Channel(在 Java 應用中表現為一個 JMS Topic
  2. 扮演主體的應用程式建立一個 TopicPublisher(一種 MessageProducer)在通道上送訊息
  3. 每個扮演觀察者的應用程式各建立一個 TopicSubscriber(一種 MessageConsumer)在通道上收訊息——這相當於 Observer 模式中的 Attach(Observer)

於是主體與觀察者透過通道建立了連結。往後主體每有變更要宣告,就送一則訊息;通道確保每個觀察者都收到這則訊息的一份副本

程式碼:推模型的 SubjectGateway 與 ObserverGateway

SubjectGateway#

public class SubjectGateway {

    public static final String UPDATE_TOPIC_NAME = "jms/Update";
    private Connection connection;
    private Session session;
    private MessageProducer updateProducer;

    public static SubjectGateway newGateway() throws JMSException, NamingException {
        SubjectGateway gateway = new SubjectGateway();
        gateway.initialize();
        return gateway;
    }

    protected void initialize() throws JMSException, NamingException {
        ConnectionFactory connectionFactory = JndiUtil.getQueueConnectionFactory();
        connection = connectionFactory.createConnection();
        session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Destination updateTopic = JndiUtil.getDestination(UPDATE_TOPIC_NAME);
        updateProducer = session.createProducer(updateTopic);

        connection.start();
    }

    public void notify(String state) throws JMSException {
        TextMessage message = session.createTextMessage(state);
        updateProducer.send(message);
    }

    public void release() throws JMSException {
        if (connection != null) {
            connection.stop();
            connection.close();
        }
    }
}

SubjectGateway 是主體(未列出)與訊息系統之間的 Messaging Gateway。主體建立這個 gateway 並用它廣播通知——本質上,主體的 Notify() 方法就實作成呼叫 SubjectGateway.notify(String),由 gateway 在更新通道上送一則訊息來宣告變更。

ObserverGateway#

public class ObserverGateway implements MessageListener {

    public static final String UPDATE_TOPIC_NAME = "jms/Update";
    private Observer observer;
    private Connection connection;
    private MessageConsumer updateConsumer;

    public static ObserverGateway newGateway(Observer observer)
            throws JMSException, NamingException {
        ObserverGateway gateway = new ObserverGateway();
        gateway.initialize(observer);
        return gateway;
    }

    protected void initialize(Observer observer) throws JMSException, NamingException {
        this.observer = observer;
        ConnectionFactory connectionFactory = JndiUtil.getQueueConnectionFactory();
        connection = connectionFactory.createConnection();
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Destination updateTopic = JndiUtil.getDestination(UPDATE_TOPIC_NAME);
        updateConsumer = session.createConsumer(updateTopic);
        updateConsumer.setMessageListener(this);
    }

    public void onMessage(Message message) {
        try {
            TextMessage textMsg = (TextMessage) message; // assume cast always works
            String newState = textMsg.getText();
            update(newState);
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }

    public void attach() throws JMSException {
        connection.start();
    }

    public void detach() throws JMSException {
        if (connection != null) {
            connection.stop();
            connection.close();
        }
    }

    private void update(String newState) throws JMSException {
        observer.update(newState);
    }
}

ObserverGateway 是觀察者與訊息系統之間的另一個 Messaging Gateway。觀察者建立 gateway,再用 attach() 啟動 Connection(相當於 Observer 模式中的 Attach(Observer))。這個 gateway 是 Event-Driven Consumer,因此實作 MessageListener 介面所要求的 onMessage

這兩個類別實作的是推模型版本的 Observer:訊息的存在告訴觀察者「有變更發生了」,訊息的內容則告訴觀察者主體的新狀態是什麼——新狀態是被過去的。

發布/訂閱相對於 RPC 的優勢#

  • 簡化通知——主體的 Notify() 實作變得極其簡單,程式碼只要在通道上送一則訊息;同樣地 Observer.Update() 只要收一則訊息。
  • 簡化 Attach/Detach——觀察者不是去 attach/detach 主體,而是訂閱/取消訂閱通道。主體根本不必實作 Attach(Observer)Detach(Observer)(雖然觀察者可能實作這兩個方法來封裝訂閱行為)。
  • 簡化並行執行緒——主體只需要一條執行緒就能並行更新所有觀察者,因為是通道在並行遞送通知訊息,而每個觀察者在自己的執行緒中處理更新——一個觀察者在自己更新執行緒裡做什麼,都不影響其他觀察者。
  • 簡化遠端存取——主體與觀察者都不必實作任何遠端方法,也不必跑在 ORB 中。它們只要存取訊息系統,分散的事由訊息系統處理。
  • 提高可靠性——因為通道使用訊息傳遞,通知會排隊等到觀察者能處理為止,這也讓觀察者能對通知做節流。想收到自己離線期間發出的通知,觀察者應該讓自己成為 Durable Subscriber

有一件事發布/訂閱沒有改變:序列化。不論用 RPC 還是訊息傳遞實作 Observer,狀態資料都得從主體的記憶體空間送到各觀察者的記憶體空間,因此都必須序列化(封送)

若說發布/訂閱有什麼缺點,那就是它需要訊息傳遞——主體與觀察者應用程式必須都能存取一套共享的訊息系統,並被實作成該系統的客戶端。不過,把應用程式變成訊息客戶端並不比走 RPC 途徑困難,多半還更容易

推模型 vs. 拉模型(在訊息傳遞下)#

發布/訂閱的另一個潛在缺點是:拉模型比推模型複雜。分散式應用之間多出來的通訊會顯著傷害效能,而且用訊息傳遞比用 RPC 更複雜

  • Update() 兩者都是單向通訊——一次回傳 void 的 RPC,或一則從主體到觀察者的 Event Message。
  • 棘手的是觀察者要查詢主體狀態時:GetState()雙向通訊——一次「索取並回傳狀態」的 RPC,或者一組 Request-Reply(一則 Command Message 索取狀態、一則獨立的 Document Message 回傳它)。

Request-Reply 之所以更麻煩,不只因為它需要一對訊息,更因為它需要一對通道

  • get-state-request 通道——從觀察者到主體
  • get-state-reply 通道——從主體回到觀察者

所有觀察者可以共用同一個請求通道,但它們多半各自需要自己的回覆通道——每個觀察者需要收到的不是隨便一個回應,而是對它自己那個請求的回應,而最容易確保這件事的方式就是各有各的回覆通道。(另一種選擇是共用單一回覆通道,靠 Correlation Identifier 判斷哪個回覆給哪個觀察者,但一觀察者一通道實作起來容易得多。)

通道爆炸與 TemporaryQueue#

一觀察者一回覆通道會導致通道爆炸。數量大或許還管得住,但訊息系統管理員無法預先知道該建幾個靜態通道——需要用到這些通道的觀察者數量在執行期是動態變化的。而且就算通道夠多,每個觀察者又怎麼知道該用哪一個?

JMS 的 TemporaryQueue 正是為此而生:觀察者可以建立一個專供自己使用的臨時佇列,把它指定為請求中的 Return Address,然後在那個佇列上等回覆。

但頻繁建立新佇列可能沒效率(視訊息系統實作而定),而且臨時佇列無法持久化(因此不能搭配 Guaranteed Delivery)。

程式碼:拉模型的 PullSubjectGateway 與 PullObserverGateway

PullSubjectGateway#

public class PullSubjectGateway {

    public static final String UPDATE_TOPIC_NAME = "jms/Update";
    private PullSubject subject;
    private Connection connection;
    private Session session;
    private MessageProducer updateProducer;

    protected void initialize(PullSubject subject) throws JMSException, NamingException {
        this.subject = subject;

        ConnectionFactory connectionFactory = JndiUtil.getQueueConnectionFactory();
        connection = connectionFactory.createConnection();
        session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Destination updateTopic = JndiUtil.getDestination(UPDATE_TOPIC_NAME);
        updateProducer = session.createProducer(updateTopic);

        new Thread(new GetStateReplier()).start();

        connection.start();
    }

    public void notifyNoState() throws JMSException {
        TextMessage message = session.createTextMessage();
        updateProducer.send(message);
    }

    private class GetStateReplier implements Runnable, MessageListener {

        public static final String GET_STATE_QUEUE_NAME = "jms/GetState";
        private Session session;
        private MessageConsumer requestConsumer;

        public void run() {
            try {
                session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
                Destination getStateQueue = JndiUtil.getDestination(GET_STATE_QUEUE_NAME);
                requestConsumer = session.createConsumer(getStateQueue);
                requestConsumer.setMessageListener(this);
            } catch (Exception e) {
                e.printStackTrace();
            }
        }

        public void onMessage(Message message) {
            try {
                Destination replyQueue = message.getJMSReplyTo();
                MessageProducer replyProducer = session.createProducer(replyQueue);

                Message replyMessage = session.createTextMessage(subject.getState());
                replyProducer.send(replyMessage);
            } catch (JMSException e) {
                e.printStackTrace();
            }
        }
    }
}

與推模型版本的差異:

  • 拉模型版本持有主體的參考,好在觀察者索取時向主體查詢狀態
  • notify(String) 變成 notifyNoState()——拉模型只送出通知而不帶任何狀態(也因為 Java 已經用掉了 notify() 這個方法名)
  • 最大的新增是 GetStateReplier 這個內部類別:它實作 Runnable 以便跑在自己的執行緒中,同時是 MessageListener(因此是 Event-Driven Consumer)。它的 onMessage 從 GetState 佇列讀取請求,並把含有主體狀態的回覆送到請求所指定的佇列(見 Request-Reply)。

PullObserverGateway#

public class PullObserverGateway implements MessageListener {

    public static final String UPDATE_TOPIC_NAME = "jms/Update";
    public static final String GET_STATE_QUEUE_NAME = "jms/GetState";
    private PullObserver observer;
    private QueueConnection connection;
    private QueueSession session;
    private MessageConsumer updateConsumer;
    private QueueRequestor getStateRequestor;

    protected void initialize(PullObserver observer) throws JMSException, NamingException {
        this.observer = observer;

        QueueConnectionFactory connectionFactory = JndiUtil.getQueueConnectionFactory();
        connection = connectionFactory.createQueueConnection();
        session = connection.createQueueSession(false, Session.AUTO_ACKNOWLEDGE);
        Destination updateTopic = JndiUtil.getDestination(UPDATE_TOPIC_NAME);
        updateConsumer = session.createConsumer(updateTopic);
        updateConsumer.setMessageListener(this);

        Queue getStateQueue = (Queue) JndiUtil.getDestination(GET_STATE_QUEUE_NAME);
        getStateRequestor = new QueueRequestor(session, getStateQueue);
    }

    public void onMessage(Message message) {
        try {
            // message's contents are empty
            updateNoState();
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }

    private void updateNoState() throws JMSException {
        TextMessage getStateRequestMessage = session.createTextMessage();
        Message getStateReplyMessage = getStateRequestor.request(getStateRequestMessage);
        TextMessage textMsg = (TextMessage) getStateReplyMessage;
        String newState = textMsg.getText();
        observer.update(newState);
    }
}

initialize 不只設定 updateConsumer 監聽更新,還設定 getStateRequestor 來送 GetState() 請求。拉模型版本的 onMessage 忽略訊息內容(訊息是空的)——訊息的存在告訴觀察者主體變了,但沒告訴它新狀態是什麼,所以它只呼叫 updateNoState()

推模型與拉模型對觀察者的差別,在 updateNoState() vs. update(String) 的實作上一目了然:推模型版本直接拿到新狀態當參數,更新觀察者即可;拉模型版本必須先去取得新狀態才能更新觀察者

注意:在這個簡化實作中,gateway 是單執行緒的,因此它在送出 get-state 請求並等待回覆的期間,不會處理任何新的更新。若請求或回覆訊息傳輸得很久,gateway 就會卡住,之後發生的更新只能排隊。

若這些代價在你的應用程式中可以接受,拉模型是可行的作法。但如果拿不定主意,就從推模型開始——它比較簡單。

通道設計:到底需要幾個通道?#

真實的企業應用複雜得多:一個應用程式可能有許多主體要宣告變更;每個主體常有數個能獨立變化的面向(aspects);單一個觀察者可能對數個主體的數個不同面向都有興趣,而那些主體不只是同一類別的多個實例,還可能是不同類別的實例。

案例一:地址變更#

企業可能有數套系統各自保存客戶聯絡資訊。地址在其中一套被更新時,該系統應通知其他可能需要這份新資訊的系統。

這個問題很簡單:只需要一個宣告地址變更的 Publish-Subscribe Channel。每個能改地址的應用程式也負責在該通道上發布訊息宣告變更;每個想收到通知的應用程式訂閱該通道。

<AddressChange customer_id="12345">
    <OldAddress>
        <Street>123 Wall Street</Street>
        <City>New York</City>
        <State>NY</State>
        <Zip>10005</Zip>
    </OldAddress>
    <NewAddress>
        <Street>321 Sunset Blvd</Street>
        <City>Los Angeles</City>
        <State>CA</State>
        <Zip>90012</Zip>
    </NewAddress>
</AddressChange>

案例二:缺貨公告#

同樣的問題、同樣的解法——用一個 Publish-Subscribe Channel 做缺貨公告:

<OutOfProduct>
    <ProductID>12345</ProductID>
    <StoreID>67890</StoreID>
    <QuantityRequested>100</QuantityRequested>
</OutOfProduct>

能不能用同一個通道傳地址變更與缺貨公告?大概不行,理由有二:

  • Datatype Channel 告訴我們同一通道上的所有訊息必須是同一型別——這裡意味著它們得符合同一份 XML schema。<AddressChange><OutOfProduct> 顯然是很不同的元素型別。
  • 就算重新設計資料格式讓兩者共用同一 schema、讓接收端能分辨,對地址變更有興趣的應用程式,多半不是對商品更新有興趣的那些——共用通道會讓應用程式頻繁收到自己不關心的通知。

因此分成兩個通道是合理的

案例三:信用評等變更#

<CreditRatingChange customer_id="12345">
    <OldRating>AAA</OldRating>
    <NewRating>BBB</NewRating>
</CreditRatingChange>

很誘人的作法是再開一個「信用評等變更」通道。

大量通道會對訊息系統造成負擔:通道多而每個流量都小,會浪費資源並讓負載難以分散;通道多而小訊息也多,則增加訊息開銷。相依者會搞不清該訂閱哪一個;多通道又需要多個收發端,可能導致一堆執行緒在檢查一堆通常是空的通道

折衷方案#

比較好的作法是:把地址變更與信用評等變更送在同一個通道上——兩者都關乎「客戶」的變化,對其中一種有興趣的應用程式可能對另一種也有興趣。但缺貨通道仍該獨立,因為關心客戶的應用程式未必關心商品,反之亦然。

由於 Datatype Channel 要求同通道訊息同格式,在 XML 中這意味著所有訊息必須有相同的根元素型別,但可以有不同的可選巢狀元素

<CustomerChange customer_id="12345">
    <AddressChange>
        <OldAddress>
            <Street>123 Wall Street</Street>
            <City>New York</City>
            <State>NY</State>
            <Zip>10005</Zip>
        </OldAddress>
        <NewAddress>
            <Street>321 Sunset Blvd</Street>
            <City>Los Angeles</City>
            <State>CA</State>
            <Zip>90012</Zip>
        </NewAddress>
    </AddressChange>
</CustomerChange>
<CustomerChange customer_id="12345">
    <CreditRatingChange>
        <OldRating>AAA</OldRating>
        <NewRating>BBB</NewRating>
    </CreditRatingChange>
</CustomerChange>

仍可能有這個問題:關心地址變更的出貨應用程式不關心信用評等變更,帳務應用程式則相反。這些應用程式可以用 Selective Consumer 只取自己要的訊息。

若 selective consumer 被證實太複雜,而訊息系統又能輕鬆支撐更多通道,那也許分開通道畢竟還是比較好。

一如企業架構與設計中的許多議題,沒有簡單答案,只有一堆取捨。使用 Publish-Subscribe Channel(以及任何訊息通道)的目標,是確保觀察者只收到它們需要的通知,同時既不造成通道爆炸,也不讓典型的觀察者背上「一堆執行緒跑一堆消費者監控一堆通道」的負擔

結論#

本章顯示 Publish-Subscribe Channel 是 Observer 模式的一種實作,讓這個模式在分散式環境中好用得多

  • 用了通道之後,Subject.Notify()Observer.Update() 都變得簡單許多——它們要做的只是收送訊息
  • 訊息系統負責處理分散與並行,同時讓遠端通知更可靠
  • 推模型比拉模型簡單、通常也更有效率,尤其在分散式通知、尤其搭配訊息傳遞時;但拉模型同樣能用訊息傳遞實作
  • 在資料變動繁多的複雜應用中,很容易想為每一種變化各開一個通道;但把送往同一群觀察者的類似通知放在同一個通道上,往往更務實