脈絡Splitter 能把單一訊息拆成一連串可個別處理的子訊息;Recipient ListPublish-Subscribe Channel 能把請求訊息平行轉送給多個接收者,以取得多個回應供挑選。

在多數這類情境中,後續處理仰賴子訊息被成功處理——例如我們想從多家廠商的回應中挑出最好的出價,或想在所有品項都從倉庫揀出之後才向客戶請款。

非同步讓「收齊」變得很難#

訊息系統的非同步本質,讓「跨多則訊息收集資訊」變得棘手:

  • 到底有幾則訊息? 若我們把訊息廣播到一個廣播通道,可能根本不知道有多少接收者在聽,因而不知道該期待幾個回應。
  • 回應可能亂序抵達。 即使用了 Splitter,回應訊息也未必依原本建立的順序回來——個別訊息可能走不同網路路徑,訊息基礎設施通常能保證每則訊息會送達,卻不保證送達順序;加上個別訊息可能由處理速度不同的各方處理,亂序幾乎必然。
  • 該等多久? 多數訊息基礎設施運作在「保證送達,但不保證何時」的模式。等太久會拖延後續處理;不等就往下走,又得設法在資訊不完整的情況下工作。
  • 遲到的訊息怎麼辦? 有時可以單獨處理,但一般而言那會導致重複處理;反過來說,若忽略遲到者,那些訊息中的資訊就永久遺失了。

所有這些問題都讓「多則相關訊息的合併處理」變複雜。若能有一個獨立元件承擔這些複雜性,只把單一一則訊息交給後續的業務處理,商業邏輯會好寫得多。

解法#

Aggregator 是一種特殊的過濾器:它接收一串訊息、辨識出彼此相關的訊息;一旦收到完整的一組,就從每則相關訊息收集資訊,並把單一一則聚合訊息發布到輸出通道供後續處理。

圖 7-7:Aggregator 解法示意

它是有狀態的#

與先前的路由模式不同,Aggregator 是有狀態元件

像 Content-Based Router 這類簡單路由模式通常是無狀態的:元件逐則處理訊息,訊息之間不必保存任何資訊,處理完後狀態與處理前相同。

Aggregator 不可能無狀態——它必須保存每則進站訊息,直到所有屬於同一組的訊息都收到,再把各訊息的資訊提煉進聚合訊息。

Aggregator 不一定要完整保存每則進站訊息。例如處理進站的拍賣出價時,我們可能只需保留最高出價與對應的出價者 ID,而不必保存所有出價訊息的歷史。但它終究要跨訊息保存資訊,因此是有狀態的。

設計時要定的三件事#

  • 關聯(Correlation)——哪些進站訊息屬於同一組?
  • 完成條件(Completeness Condition)——什麼時候可以發布結果訊息?
  • 聚合演算法(Aggregation Algorithm)——如何把收到的訊息合併成單一結果訊息?

關聯通常靠進站訊息的型別或一個明確的 Correlation Identifier 達成。

實作細節#

由於訊息系統是事件驅動的,Aggregator 可能在任何時間、以任何順序收到相關訊息。為了關聯訊息,它維護一份活躍聚合體的清單(也就是已經收到部分訊息的那些聚合體):

  1. 收到新訊息時,檢查它是否屬於某個既有聚合體
  2. 若沒有相關的聚合體,就假設這是一組訊息中的第一則,建立新的聚合體並把訊息加進去
  3. 若已有聚合體,就直接把訊息加進去
  4. 加入後評估該聚合體的完成條件——為真就組出聚合訊息發布到輸出通道;為假就不發布,讓聚合體保持活躍等待更多訊息

這種策略在收到「無法關聯到既有聚合體」的訊息時就新建一個,因此 Aggregator 不需要事先知道自己會產出哪些聚合體——本書稱這種變體為 Self-starting Aggregator(自啟式聚合器)

處理已結案的聚合體#

視聚合策略而定,Aggregator 可能得處理「進站訊息屬於一個已經結案(聚合訊息已發布)的聚合體」的情況。為了避免又開一個新聚合體,它必須保存一份已結案聚合體的清單

這份清單需要定期清除機制,否則會無限成長。這假設我們能對「相關訊息會在多長的時間窗內抵達」做出基本假設。

好消息是:我們不必保存完整的聚合體,只需保存「它已結案」這個事實,因此清單能存得相當精簡,也能在清除演算法中留出足夠的安全邊際。我們也可以用 Message Expiration 忽略延遲過久的訊息。

加上控制通道#

為了提高整體方案的強健性,可以讓 Aggregator 監聽一個控制通道,允許手動清除全部或特定的活躍聚合體——這在「想從錯誤狀態復原、又不想重啟元件」時很有用。

同理,允許 Aggregator 依請求把活躍聚合體清單發布到特殊通道,是非常好用的除錯功能。這兩項都是典型會被納入 Control Bus 的功能。

完成條件的策略#

可用策略主要取決於我們知不知道該期待幾則訊息(Aggregator 可能因為收到原始複合訊息的副本、或因為每則訊息都含有總數而知道):

  • Wait for All(等全部)——等到所有回應都收到。前面的訂單例子多半屬於這種:不完整的訂單沒有意義,所以若在逾時期限內沒收齊,Aggregator 就該拋出錯誤。

這給出最好的決策基礎,但也最慢、最脆弱(而且我們得知道該期待幾則)。單單一則遺失或延遲的訊息,就會擋住整個聚合體的後續處理。

在鬆散耦合的非同步系統中,解決這類錯誤狀況很複雜——非同步的訊息流讓錯誤狀況難以可靠偵測(到底要等多久才算「遺失」?)。一種對策是重新索取該訊息,但這要求 Aggregator 知道訊息的來源,又引入了額外的相依

  • Time Out(逾時)——等待指定時間,然後根據時限內收到的回應做決定。若完全沒收到回應,系統可以回報例外或重試。

當進站回應會被評分、而我們只取分數最高的一則(或少數幾則)時,這個啟發式很有用——常見於「競標」情境

  • First Best(取最快)——只等第一個(最快的)回應,忽略其他全部。最快,但也忽略了大量資訊;在回應時間至關重要的競標或報價情境中可能很實用。
  • Time Out with Override(逾時但可提前結束)——等待指定時間,或直到收到分數達到預設門檻的訊息為止。也就是說:遇到非常有利的回應就提早收工;否則就等到時間用完,屆時若沒有明顯贏家,就對目前收到的所有訊息做排序。
  • External Event(外部事件)——有時聚合由外部業務事件的到來而結束。例如在金融業,收盤可能就代表進站報價聚合的終點。

用固定計時器來代表這類事件會降低彈性,因為它無法容納變動;而以 Event Message 形式出現的指定業務事件,則讓系統能被集中控制。Aggregator 可以在特殊控制通道上監聽該事件訊息,或接收一則格式特殊、代表聚合結束的訊息。

聚合演算法的策略#

  • 挑「最好」的答案——假設存在單一個最佳解,例如同一項商品的最低出價。

但現實中的挑選準則很少這麼單純:某項商品的「最佳」出價,可能還取決於交期、可供數量、廠商是否在優先供應商名單上等等。

  • 濃縮資料——Aggregator 可用來降低高流量來源的訊息量,例如計算個別訊息的平均值、或把每則訊息的數值欄位加總成單一訊息。這在每則訊息代表一個數值(例如收到的訂單數)時效果最好。
  • 收集資料留待後續評估——並不總是能由 Aggregator 決定如何挑出最佳答案。這種情況下仍值得用它把個別訊息收集並合併成單一訊息(可能只是把各訊息資料彙編在一起),聚合決策則稍後由另一個元件或人來做

自啟式 vs. 初始化式#

聚合策略常由參數驅動:等待策略要設定最長等待時間;門檻策略要事先讓 Aggregator 知道門檻值。

若這些參數要能在執行期設定,Aggregator 可以多開一個輸入來接收控制訊息。控制訊息也可以帶上「預期會有幾則相關訊息」這類資訊,幫助 Aggregator 實作更有效的完成條件

在這種情境下,Aggregator 不是等第一則訊息到來才新建聚合體,而是事先收到與一系列預期訊息相關的資訊——這份資訊可以是原始請求訊息(例如一則 Scatter-Gather 訊息)的副本,再加上必要的參數。Aggregator 於是配置一個新聚合體並把參數資訊存在其中;個別訊息進來時再與對應的聚合體關聯。

本書稱這種變體為 Initialized Aggregator(初始化式聚合器),相對於 Self-starting Aggregator。這種配置顯然只在我們能取得原始訊息時才可行,而那並非總是如此。

Aggregator 常與 SplitterRecipient List 耦合成複合模式——見 Composed Message ProcessorScatter-Gather

圖 7-9:初始化式 Aggregator(Initialized Aggregator)

範例:Loan Broker 與「遺失訊息偵測器」

圖 7-8:帶逾時的 Aggregator 可偵測遺失訊息

Loan Broker#

本章末的組合式訊息間奏用 Aggregator 從各銀行回傳的貸款報價中挑出最好的一筆。這個例子用的是初始化式 Aggregator——由 Recipient List 告知 Aggregator 該期待幾則報價訊息。間奏中提供 Java、C# 與 TIBCO 三種實作。

把 Aggregator 當成遺失訊息偵測器#

Joe Walnes 示範了 Aggregator 的一種創意用法。他的系統把訊息送過一連串相當不可靠的元件。

即使用 Guaranteed Delivery 也解決不了這個問題——因為典型的失敗發生在系統消費了訊息之後才掛掉;而由於這些應用程式不是 Transactional Client,處理中的訊息就這樣沒了

他的解法是把進站訊息沿兩條平行路徑送出:一條穿過那些必要但不可靠的元件,另一條用 Guaranteed Delivery 繞過它們;再用 Aggregator 把兩條路徑的訊息重新組合。

Aggregator 採用 Time Out with Override 完成條件,因此逾時或兩則相關訊息都收到,都會使它完成;聚合演算法則取決於哪個條件先被滿足:

  • 收到兩則訊息 → 把處理過的訊息原封不動往下傳
  • 發生逾時事件 → 我們就知道有某個元件失敗、「吃掉」了訊息;於是指示 Aggregator 發布一則錯誤訊息告警維運人員

這個配置中元件仍得手動重啟,但更講究的作法很可能可以自動重啟元件並重送遺失的訊息。

範例:用 JMS 實作 Aggregator

這個範例在一個通道上接收出價訊息,聚合所有相關出價,並把最低出價發布到另一個通道。出價透過 AuctionID 屬性關聯(它扮演 Correlation Identifier)。聚合策略是至少收到 3 筆出價;這個 Aggregator 是自啟式的,不需要外部初始化。

方案由四個主要類別組成:

  • Aggregator——收訊息、聚合、送出結果訊息的邏輯,透過 Aggregate 介面與聚合體互動
  • AuctionAggregate——實作 Aggregate 介面,扮演 Aggregate 介面與 Auction 類別之間的 Adapter(見 [GoF]),Auction 類別完全不必參照 JMS API
  • Auction——已收到的相關出價集合,實作聚合策略(找最低價、判斷是否完成)
  • Bid——一個便利類別,持有出價相關的資料項
public class Aggregator implements MessageListener
{
    static final String PROP_CORRID = "AuctionID";

    Map activeAggregates = new HashMap();

    Destination inputDest = null;
    Destination outputDest = null;
    Session session = null;

    MessageConsumer in = null;
    MessageProducer out = null;

    public Aggregator (Destination inputDest, Destination outputDest, Session session)
    {
        this.inputDest = inputDest;
        this.outputDest = outputDest;
        this.session = session;
    }

    public void run()
    {
        try {
            in = session.createConsumer(inputDest);
            out = session.createProducer(outputDest);
            in.setMessageListener(this);
        } catch (Exception e) {
            System.out.println("Exception occurred: " + e.toString());
        }
    }

    public void onMessage(Message msg)
    {
        try {
            String correlationID = msg.getStringProperty(PROP_CORRID);
            Aggregate aggregate = (Aggregate)activeAggregates.get(correlationID);
            if (aggregate == null) {
                aggregate = new AuctionAggregate(session);
                activeAggregates.put(correlationID, aggregate);
            }
            //--- ignore message if aggregate is already closed
            if (!aggregate.isComplete()) {
                aggregate.addMessage(msg);
                if (aggregate.isComplete()) {
                    MapMessage result = (MapMessage)aggregate.getResultMessage();
                    out.send(result);
                }
            }
        } catch (JMSException e) {
            System.out.println("Exception occurred: " + e.toString());
        }
    }
}

Aggregator 是 Event-Driven Consumer。對每則進站訊息,它取出 correlation ID、檢查是否已有對應的活躍聚合體;沒有就新建 AuctionAggregate。接著檢查聚合體是否仍活躍——已結案就丟棄該訊息;仍活躍就加入並測試終止條件,滿足就取出最佳出價並發布。

這段程式碼相當通用,只有兩行與本例綁定:假設 correlation ID 存在 AuctionID 屬性中,以及建立 AuctionAggregate 實例。後者可以用一個回傳 Aggregate 型別的工廠來消除——但因為這是一本談企業整合而非物件導向設計的書,就保持簡單了。

Aggregate 介面很簡單,只有三個方法:

public interface Aggregate {
    public void addMessage(Message message);
    public boolean isComplete();
    public Message getResultMessage();
}

聚合策略並未實作在 AuctionAggregate 中,而是獨立成不相依於 JMS APIAuction 類別:

public class Auction
{
    ArrayList bids = new ArrayList();

    public void addBid(Bid bid)
    {
        bids.add(bid);
        System.out.println(bids.size() + " Bids in auction.");
    }

    public boolean isComplete()
    {
        return (bids.size() >= 3);
    }

    public Bid getBestBid()
    {
        Bid bestBid = null;

        Iterator iter = bids.iterator();
        if (iter.hasNext())
            bestBid = (Bid) iter.next();
        while (iter.hasNext()) {
            Bid b = (Bid) iter.next();
            if (b.getPrice() < bestBid.getPrice()) {
                bestBid = b;
            }
        }
        return bestBid;
    }
}

Auction 提供與 Aggregate 介面類似的三個方法,但簽章用的是強型別的 Bid 而非 Message。本例的聚合策略非常簡單——等到收滿三筆出價;但因為策略與 JMS API 分離,要把 Auction 強化成更複雜的邏輯很容易。

AuctionAggregate 則扮演 Aggregate 介面與 Auction 之間的 Adapter:

public class AuctionAggregate implements Aggregate {
    static String PROP_AUCTIONID = "AuctionID";
    static String ITEMID = "ItemID";
    static String VENDOR = "Vendor";
    static String PRICE = "Price";

    private Session session;
    private Auction auction;

    public AuctionAggregate(Session session)
    {
        this.session = session;
        auction = new Auction();
    }

    public void addMessage(Message message) {
        Bid bid = null;
        if (message instanceof MapMessage) {
            try {
                MapMessage mapmsg = (MapMessage)message;
                String auctionID = mapmsg.getStringProperty(PROP_AUCTIONID);
                String itemID = mapmsg.getString(ITEMID);
                String vendor = mapmsg.getString(VENDOR);
                double price = mapmsg.getDouble(PRICE);
                bid = new Bid(auctionID, itemID, vendor, price);
                auction.addBid(bid);
            } catch (JMSException e) {
                System.out.println(e.getMessage());
            }
        }
    }

    public boolean isComplete()
    {
        return auction.isComplete();
    }

    public Message getResultMessage() {
        Bid bid = auction.getBestBid();
        try {
            MapMessage msg = session.createMapMessage();
            msg.setStringProperty(PROP_AUCTIONID, bid.getCorrelationID());
            msg.setString(ITEMID, bid.getItemID());
            msg.setString(VENDOR, bid.getVendorName());
            msg.setDouble(PRICE, bid.getPrice());
            return msg;
        } catch (JMSException e) {
            System.out.println("Could not create message: " + e.getMessage());
            return null;
        }
    }
}

這個簡化範例假設 Auction ID 是全域唯一的,因此不必操心清理未結案的拍賣清單——就讓它一直長。真實應用中必須決定何時清除舊拍賣紀錄,以避免記憶體洩漏

因為程式只參照 JMS Destination,主題與佇列都能跑。正式環境中大概更可能用 Point-to-Point Channel(等同 JMS queue),因為一筆出價只該有一個接收者,就是 Aggregator

但如 Publish-Subscribe Channel 所述,主題能簡化測試與除錯——在不影響訊息流的前提下加一個 listener 非常容易。許多 JMS 實作允許主題名稱使用萬用字元,因此 listener 可以用 * 訂閱所有主題;有個能顯示主題上所有訊息、並把訊息記錄到檔案供事後分析的簡單監聽工具,非常好用。