脈絡:Message Router 可以依訊息內容或其他準則把訊息路由到不同通道。因為個別訊息可能走不同路徑,有些訊息會比其他訊息更早通過處理步驟,導致訊息亂序。但某些後續處理步驟確實需要訊息依序處理(例如維護參照完整性)。

為什麼會亂序#

顯而易見的解法是一開始就別讓訊息亂序。保持有序確實比事後回復順序容易——這也是為什麼許多大學圖書館寧可不讓讀者自己把書放回架上:只要控制插入過程,任何時刻的順序都(幾乎)有保證。但在非同步訊息方案中保持順序,難度大概不亞於說服青少年「把房間保持整齊其實比較有效率」。

亂序的一個常見成因是不同訊息走不同的處理路徑

假設有一串編號訊息,所有偶數號訊息得經過一次特殊轉換,奇數號則直接放行。那麼奇數號訊息會很快出現在結果通道上,而偶數號訊息還在轉換那裡排隊。若轉換相當慢,所有奇數訊息可能在任何一則偶數訊息出現之前就全部通過,順序完全亂掉。

圖 7-10:不同處理路徑造成訊息亂序

一個行不通的解法#

要避免亂序,我們可以引入一個回送(確認)機制,確保系統中一次只有一則訊息通過——上一則處理完才送下一則。

這個保守作法確實能解決問題,但有兩個重大缺點:

  • 會嚴重拖慢系統。 若有大量平行處理單元,運算能力會被嚴重閒置。很多時候採用平行處理正是為了提升效能,把流量節流到一次一則,等於完全抵銷了方案的目的。
  • 它要求我們能控制送進處理單元的訊息。 但我們常常只是站在「亂序訊息流」的接收端,對訊息來源根本沒有控制權

解法#

Resequencer 內含一個內部緩衝區,把亂序訊息存起來直到取得完整序列;然後把依序的訊息發布到輸出通道。

和多數路由器一樣,Resequencer 通常不修改訊息內容

序號#

Resequencer 要能運作,每則訊息都必須有唯一的序號(Sequence Number)(見 Message Sequence)。

  • Message Identifier 是唯一識別每則訊息的特殊屬性,但多數情況下它不可比較——基本上是隨機值,甚至常常不是數字。即使碰巧是數值,把序號語意疊加在既有的 Message Identifier 上通常是壞主意
  • Correlation Identifier 的設計目的是把進站訊息與原始出站請求配對,唯一的要求是唯一性,不必是數字、也不必連續。

所以若要保存一串訊息的順序,就該另外定義一個欄位來追蹤序號——通常這個欄位屬於訊息表頭。

產生序號比產生唯一識別碼更麻煩#

唯一識別碼常能以分散方式產生——結合唯一的位置資訊(例如網卡 MAC 位址)與當前時間,多數 GUID 演算法就是這麼做的。

但要產生依序的號碼,一般需要一個跨系統指派號碼的單一計數器;而且多數情況下光是遞增還不夠,號碼還必須連續,否則很難辨識出遺失的訊息。一不小心,這個序號產生器就會成為訊息流的瓶頸。

若個別訊息是使用 Splitter 的產物,最好把編號直接做進 Splitter 裡

內部運作#

序號確保 Resequencer 能偵測到亂序訊息。亂序意味著序號較大的訊息比序號較小的先抵達。Resequencer 必須把序號較大的訊息存起來,直到收到所有「缺漏」的較小序號訊息;期間可能還會收到其他亂序訊息,同樣得存起來。一旦緩衝區內含一段連續的訊息序列,就把它送到輸出通道並從緩衝區移除。

走一個簡單例子——Resequencer 依序收到序號 1、3、5、2 的訊息(假設序列從 1 開始):

  1. 收到 1:可以立刻送出並從緩衝區移除
  2. 收到 3:缺 2,因此把 3 存起來
  3. 收到 5:同樣存起來
  4. 收到 2:緩衝區現在有 2、3 這段完整序列,於是發布它們並從緩衝區移除5 則繼續留在緩衝區,直到序列中剩下的「缺口」被補上

圖 7-11:Resequencer 的內部運作

避免緩衝區溢位#

緩衝區該多大?

若處理的是很長的訊息串,緩衝區可能變得相當大。更糟的是:假設有多個處理單元、各自處理特定訊息型別,只要其中一個單元故障,我們就會得到一長串亂序訊息,緩衝區溢位幾乎必然發生。

有些情況下可以用訊息佇列本身來吸收待處理的訊息——但這只在「訊息基礎設施允許依選擇準則從佇列讀取訊息、而不是永遠先讀最舊的那則」時可行。那樣我們就能輪詢佇列、看看第一則「缺漏」的訊息到了沒,而不必消費掉中間所有訊息。

但到了某個時點,連配置給訊息佇列的儲存空間也會被填滿

用主動確認節流生產端#

一個強健的作法是節流訊息生產者。如前所述,一次只送一則訊息效率太差,所以我們得聰明一點:

但它要求我們能存取原始的有序訊息串,才能插入送出緩衝與節流閥。

這與 TCP/IP 極為相似#

這種作法非常像 TCP/IP 協定的運作方式。TCP 的關鍵特性之一就是確保封包在網路上依序遞送——實際上每個封包可能走不同的網路路徑,亂序相當常見。

接收端維護一個當作滑動視窗的環形緩衝區;收發雙方協商「每次確認前要送幾個封包」。因為送出方要等接收方的確認,快速的送出方無法超前接收方,也無法讓緩衝區溢位。特定規則還防止所謂的 Silly Window Syndrome——收發雙方陷入極沒效率的「一次一個封包」模式。

用替身訊息填補#

另一種解法是為缺漏的訊息算出替身訊息。這在「接收端能容忍『夠好』的訊息資料、不要求每則訊息都精確」或「速度比準確度更重要」時可行。

例如在 VoIP 傳輸中,填入一個空白封包所帶來的使用者體驗,比為遺失封包發出重送請求要好

我們這些應用開發者多半把可靠的網路通訊視為理所當然。設計訊息方案時,深入看看 TCP 的內部其實很有幫助——因為 IP 流量的本質就是非同步且不可靠的,它要面對的許多議題與企業整合方案完全相同。

範例:用 .NET 與 MSMQ 實作 Resequencer

測試配置#

測試由四個 C# 類別組成的元件構成,透過 MSMQ 佇列通訊:

  • MQSend——測試訊息產生器。訊息內文是簡單文字字串,並把序號放進每則訊息的 AppSpecific 屬性(序號從 1 開始,訊息數量由命令列傳入),發布到私有佇列 inQueue
  • DelayProcessor——從 inQueue 讀訊息,唯一的「處理」就是延遲一段時間再把同樣的訊息重新發布到 outQueue。我們平行跑三個 DelayProcessor 模擬負載平衡的處理單元;它們扮演 Competing Consumers,因此每則訊息恰好被一個處理器消費。由於處理速度不同,outQueue 上的訊息就亂序了。
  • Resequencer——把亂序的進站訊息緩衝起來,依序重新發布到 sequenceQueue
  • MQSequenceReceive——從 sequenceQueue 讀訊息,驗證 AppSpecific 中的序號是否遞增。

實際跑起來會看到:抵達 Resequencer 的訊息確實亂序(例如 3、4、1、5、7、2……),而從輸出可以看到 Resequencer 如何在缺漏訊息時緩衝進站訊息,一旦缺漏的訊息抵達,就立刻發布已補齊的序列

圖 7-12:Resequencer 的測試配置

圖 7-13:測試元件的輸出

共用的 Processor 基底類別#

觀察測試配置會發現 DelayProcessor 與 Resequencer 有共通點:兩者都從輸入佇列讀訊息、發布到輸出佇列,差別只在中間對訊息做了什麼。因此作者抽出一個共用基底類別,封裝這種泛用 Filter 的基本功能(見 Pipes and Filters)——它含有建立佇列、以及非同步接收/處理/送出訊息的便利方法與樣板方法

using System;
using System.Messaging;
using System.Threading;

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

        public Processor (MessageQueue inputQueue, MessageQueue outputQueue)
        {
            this.inputQueue = inputQueue;
            this.outputQueue = outputQueue;
            inputQueue.Formatter = new System.Messaging.XmlMessageFormatter(
                new String[] {"System.String,mscorlib"});
            inputQueue.MessageReadPropertyFilter.ClearAll();
            inputQueue.MessageReadPropertyFilter.AppSpecific = true;
            inputQueue.MessageReadPropertyFilter.Body = true;
            inputQueue.MessageReadPropertyFilter.CorrelationId = true;
            inputQueue.MessageReadPropertyFilter.Id = true;
            Console.WriteLine("Processing messages from " + inputQueue.Path
                + " to " + outputQueue.Path);
        }

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

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

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

            ProcessMessage(m);

            mq.BeginReceive();
        }

        protected virtual void ProcessMessage(Message m)
        {
            string body = (string)m.Body;
            Console.WriteLine("Received Message: " + body);
            outputQueue.Send(m);
        }
    }
}

因為很容易忘記在訊息處理結尾呼叫 BeginReceive,作者用樣板方法把這一步包了進去——子類別只要覆寫 ProcessMessage,完全不必操心非同步處理。

圖 7-14:DelayProcessor 與 Resequencer 共用 Processor 基底類別

Resequencer 本身#

using System;
using System.Messaging;
using System.Collections;
using MsgProcessor;

namespace Resequencer
{
    class Resequencer : Processor
    {
        private int startIndex = 1;
        private IDictionary buffer = (IDictionary)(new Hashtable());
        private int endIndex = -1;

        public Resequencer(MessageQueue inputQueue, MessageQueue outputQueue)
            : base (inputQueue, outputQueue) {}

        protected override void ProcessMessage(Message m)
        {
            AddToBuffer(m);
            SendConsecutiveMessages();
        }

        private void AddToBuffer(Message m)
        {
            Int32 msgIndex = m.AppSpecific;
            Console.WriteLine("Received message index {0}", msgIndex);
            if (msgIndex < startIndex)
            {
                Console.WriteLine("Out of range message index! Current start is: {0}", startIndex);
            }
            else
            {
                buffer.Add(msgIndex, m);
                if (msgIndex > endIndex)
                   endIndex = msgIndex;
            }
            Console.WriteLine("   Buffer range: {0} - {1}", startIndex, endIndex);
        }

        private void SendConsecutiveMessages()
        {
            while (buffer.Contains(startIndex))
            {
                Message m = (Message)(buffer[startIndex]);
                Console.WriteLine("Sending message with index {0}", startIndex);
                outputQueue.Send(m);
                buffer.Remove(startIndex);
                startIndex++;
            }
        }
    }
}

緩衝區以 Hashtable 實作,鍵是存在 AppSpecific 中的訊息序號。新訊息加入後,SendConsecutiveMessages 檢查是否已有從下一個待送序號開始的連續序列;有的話就全部送出並從緩衝區移除。

這個實作假設序列從 1 開始。只有在生產者也從 1 開始、且兩個元件在整個生命週期維持同一序列時才行得通。

要更有彈性的話,生產者應該先與 Resequencer 協商序列起始號碼,才送出序列中的第一則訊息——這個過程類比於 TCP 連線建立時交換的 SYN 訊息。

這個實作也沒有處理緩衝區溢位。若某個 DelayProcessor 中止或故障、「吃掉」了一則訊息,Resequencer 會無限期等待那則遺失的訊息,直到緩衝區溢位

在高流量情境下,生產者與 Resequencer 必須協商一個視窗大小,描述 Resequencer 最多能緩衝多少訊息;緩衝區一滿,就得由錯誤處理器決定如何處置遺失的訊息——例如由生產者重送,或注入一則「假」訊息。