脈絡:在許多企業整合情境中,單一事件會觸發一連串處理步驟,每一步各司其職。

以「一筆新訂單以訊息形式抵達企業」為例,可能有這些要求:

  • 訊息必須加密,防止竊聽者窺看客戶的訂單
  • 訊息必須帶有數位憑證形式的驗證資訊,確保只有受信任的客戶能下單
  • 外部可能送來重複訊息(還記得購物網站上「請只按一次『立即訂購』」的警語嗎?)——為避免重複出貨與客戶不滿,必須在後續訂單處理啟動前把重複訊息剔除

也就是說:我們要把一串「可能重複、經過加密、帶著額外驗證資料」的訊息,轉換成一串「唯一、單純、不含多餘欄位的明文訂單訊息」。

為什麼不寫成一個大模組#

一種可能的解法是寫一個包山包海的「進站訊息處理模組」,把所有必要功能都做進去。但這種作法缺乏彈性且難以測試

  • 如果要新增或移除一個步驟怎麼辦?例如某些大客戶走私有網路下單,根本不需要加密。
  • 把所有功能塞進單一元件也減少了重用的機會。訂單狀態訊息可能同樣需要解密,卻不需要去重(重複的狀態查詢一般無害);把解密拆成獨立模組,就能在其他訊息上重用它。

異質環境的現實限制#

整合方案通常是一組異質系統,因此不同處理步驟可能得在不同的實體機器上執行:

解密進站訊息所需的私鑰,可能因安全理由只存在於某台指定機器上、無法從其他機器存取。這表示解密元件必須跑在那台機器,其他步驟則可能跑在別處。

同樣地,不同處理步驟可能用不同的程式語言或技術實作,導致它們無法跑在同一個行程、甚至同一台電腦上。

拆成元件還不夠——還要拆掉相依性#

把每個功能實作成獨立元件,仍可能在元件之間引入相依性。例如若解密元件直接呼叫驗證元件並傳入解密結果,我們就無法在不帶驗證功能的情況下使用解密功能

要化解這件事,我們得能把既有元件「組合」成一串處理步驟,而且每個元件都獨立於系統中的其他元件——這意味著元件必須暴露通用的對外介面,好讓彼此可替換

此外,既然用的是非同步訊息,就該善用它的非同步特性:元件可以把訊息送給下一個元件繼續處理,而不必等待結果——藉此在每個元件內部並行處理多則訊息

解法#

每個過濾器暴露的介面極其簡單:從入站管道收訊息、處理它、把結果發布到出站管道。管道把一個過濾器接到下一個,把前者的輸出訊息送給後者。

因為所有元件使用相同的對外介面,只要接上不同的管道,就能組合出不同的方案——新增過濾器、省略既有過濾器、或重新排列順序,全都不必修改過濾器本身

過濾器與管道之間的連接處有時稱作埠(port)。在最基本的形式中,每個過濾器有一個輸入埠與一個輸出埠。

套用到本章的範例問題上,Pipes and Filters 架構產出三個過濾器、兩條中間管道;再加上一條把訊息送進解密元件的管道,以及一條把明文訂單訊息從去重器送到訂單管理系統的管道,總共四條管道

圖 3-3:以 Pipes and Filters 處理進站訂單(解密→驗證→去重)

這是訊息系統的根本架構風格#

Pipes and Filters 描述的是訊息系統的一種根本架構風格:個別處理步驟(「過濾器」)透過訊息通道(「管道」)串接起來。

本書後續許多模式——例如路由與轉換模式——都建立在這個風格之上,這正是我們能把個別模式輕鬆組合成更大方案的原因

管道的實作可以抽換#

Pipes and Filters 用抽象的管道讓元件彼此解耦:一個元件把訊息送進管道,稍後由它並不認識的另一個行程消費。最直覺的實作就是本章開頭談的 Message Channel——多數訊息通道能在過濾器之間提供語言、平台與位置的獨立性,讓我們能因相依性、維護或效能理由把某個處理步驟搬到別台機器。

但若所有元件其實都能待在同一台機器上,訊息基礎設施提供的 Message Channel 會顯得太笨重——用一個簡單的記憶體內佇列來實作管道會有效率得多。

因此,把元件設計成與抽象的管道介面溝通很有價值:介面的實作可以在 Message Channel 與記憶體內佇列之間抽換。Messaging Gateway 說明如何為這種彈性設計元件。

代價:通道變多#

Pipes and Filters 架構的潛在缺點之一,是所需的通道數量變多:

  • 通道並非無限資源——它們提供緩衝等功能,會消耗記憶體與 CPU 週期。
  • 把訊息發布到通道有一定開銷,因為資料必須從應用程式內部格式轉換成訊息基礎設施自己的格式,接收端還得反向轉一次。

若使用很長的過濾器鏈,我們是在用反覆的訊息資料轉換所帶來的效能損失,換取彈性上的收穫。

在訊息情境下放寬純粹形式#

純粹的 Pipes and Filters 只允許每個過濾器有單一輸入埠與單一輸出埠。在訊息傳遞的情境中,這個性質可以放寬一些:

  • 一個元件可以從多個通道消費訊息,也可以輸出到多個通道(例如 Message Router)。
  • 多個過濾器元件也可以從同一個 Message Channel 消費訊息;此時 Point-to-Point Channel 確保每則訊息只被其中一個過濾器消費。

常被忽略的好處:可測試性#

使用 Pipes and Filters 也改善了可測試性。我們可以對個別處理步驟送入 Test Message 並比對結果與預期,隔離地測試與除錯每個核心功能——而且測試機制能針對該功能量身打造

  • 測試加解密:送入大量含隨機資料的訊息,加密再解密後與原文比對。
  • 測試驗證:則需要提供帶有特定驗證碼、且對應到系統中已知使用者的訊息。

管線化處理#

用非同步 Message Channel 連接元件,讓鏈上每個單元都能在自己的執行緒或行程中運作。一個單元處理完一則訊息後,就把它送到輸出通道並立刻開始處理下一則,不必等後續元件讀取與處理。

於是多則訊息能在通過各階段時並行處理:第一則訊息解密完交給驗證元件的同時,下一則已經可以開始解密。

這種配置稱為處理管線(processing pipeline)——訊息像液體流過管子一樣流過各過濾器。相較於嚴格循序的處理,處理管線能顯著提高系統吞吐量

圖 3-4:Pipes and Filters 的管線化處理

平行處理#

然而整體吞吐量受限於鏈上最慢的那個行程。要改善它,可以部署該行程的多個平行實例。

此時需要一個搭配 Competing ConsumersPoint-to-Point Channel,保證通道上的每則訊息恰好被 N 個可用處理器中的一個消費。

這種配置可能導致訊息亂序處理。若訊息順序至關重要,就只能讓每個元件跑單一實例,或改用 Resequencer

例如若解密比驗證慢得多,就可以跑三個平行的解密元件實例。

歷史:Pipes and Filters 的來歷

Pipes and Filters 架構絕非新概念。它簡潔優雅,又兼具彈性與高吞吐量,難怪一直很受歡迎;簡單的語意也讓形式化方法得以用來描述這個架構。

  • Kahn 於 1974 年提出 Kahn Process Networks:一組由無界 FIFO 通道連接的平行行程 [Kahn]。
  • [Garlan] 有一章談各種架構風格,包含 Pipes and Filters。
  • [Monroe] 詳細處理了架構風格與設計模式之間的關係。
  • [PLOPD1] 收錄 Regine Meunier 的〈The Pipes and Filters Architecture〉,構成了 [POSA] 中該模式的基礎。

幾乎所有與整合相關的 Pipes and Filters 實作,都遵循 [POSA] 提出的「Scenario IV」——使用主動過濾器,自行從排隊管道拉取、處理、再推送。

Buschmann 描述的模式假設每個元素在過濾器之間傳遞時都經歷相同的處理步驟;在整合情境下通常並非如此。許多時候訊息必須依內容或外部控制動態路由——路由在企業整合中如此常見,值得擁有自己的模式:Message Router

Pipes and Filters 與 CSP(Communicating Sequential Processes) 有些相似之處。CSP 由 Hoare 於 1978 年提出 [CSP],提供描述平行處理系統中同步問題的簡單模型:其基礎機制是兩個行程透過 I/O 同步——當行程 A 表示準備好輸出給 B、且 B 表示準備好從 A 輸入時,I/O 才發生;若只有一方成立,該行程就被放進等待佇列直到另一方就緒。

CSP 與整合方案的差別在於:它們沒那麼鬆散耦合,「管道」也不提供任何排隊機制。不過學術界對 CSP 的深入處理仍值得我們借鑑。

邊欄:關於「過濾器」這個詞

討論 Pipes and Filters 架構時,得對「filter」這個詞格外小心。

本書稍後會定義另外兩個模式:Message FilterContent Filter。這兩者都是通用 filter 的特例,但這個模式語言中的許多其他模式也是。換句話說:一個模式不必真的執行過濾功能(例如剔除欄位或訊息),也能是 Pipes and Filters 意義下的 filter。

作者本可以替 Pipes and Filters 架構風格改名來避免混淆,但考量它是如此重要且被廣泛討論的概念,改名只會更亂。因此書中會謹慎使用「filter」一詞,並盡量講清楚指的是哪一種;在仍可能混淆的地方,通用 filter 一律改稱「元件(component)」——這個詞夠通用(也夠常被濫用),應該不會惹出麻煩。

範例:用 C# 與 MSMQ 寫一個簡單過濾器

以下是一個「單一輸入埠、單一輸出埠」過濾器的通用基底類別。基礎實作只是印出收到訊息的內文並送到輸出埠;更有意思的過濾器會繼承 Processor 並覆寫 ProcessMessage,對訊息做額外處理(轉換內容、或路由到不同輸出通道)。

注意 Processor 在實例化時需要輸入與輸出通道的參考,但類別本身既不綁定特定通道、也不綁定任何其他過濾器——這讓我們能實例化多個過濾器,並以任意組態串接它們。

using System;
using System.Messaging;

namespace PipesAndFilters
{
    public class Processor
    {
        protected MessageQueue inputQueue;
        protected MessageQueue outputQueue;

        public Processor (MessageQueue inputQueue, MessageQueue outputQueue)
        {
            this.inputQueue = inputQueue;
            this.outputQueue = outputQueue;
        }

        public void Process()
        {
            inputQueue.ReceiveCompleted += new ReceiveCompletedEventHandler(OnReceiveCompleted);
            inputQueue.BeginReceive();
        }

        private void OnReceiveCompleted(Object source, ReceiveCompletedEventArgs asyncResult)
        {
            MessageQueue mq = (MessageQueue)source;

            Message inputMessage = mq.EndReceive(asyncResult.AsyncResult);
            inputMessage.Formatter = new System.Messaging.XmlMessageFormatter(
                new String[] {"System.String,mscorlib"});

            Message outputMessage = ProcessMessage(inputMessage);

            outputQueue.Send(outputMessage);

            mq.BeginReceive();
        }

        protected virtual Message ProcessMessage(Message m)
        {
            Console.WriteLine("Received Message: " + m.Body);
            return (m);
        }
    }
}

這個實作是一個 Event-Driven ConsumerProcess 方法註冊對進站訊息的關注,指示訊息系統每當訊息抵達就呼叫 OnReceiveCompleted;該方法從事件物件取出訊息資料,再呼叫虛擬方法 ProcessMessage

這個簡單範例不具交易性。若處理訊息時(在送到輸出通道之前)發生錯誤,訊息就遺失了——這在正式環境中通常不可接受。解法見 Transactional Client