脈絡:訊息系統出於必要,內部使用交易行為。而讓外部客戶端能控制影響自身行為的交易範圍,可能是有價值的

訊息系統本來就在用交易#

單一個 Message Channel 可以有多個寄件端與多個收件端,因此訊息系統必須協調訊息,確保:

  • 寄件端不會覆寫彼此的訊息
  • 多個點對點收件端不會收到同一則訊息
  • 多個發布訂閱收件端各收到一份副本

訊息系統也必須運用交易(最好是兩階段的分散式交易)把訊息從寄件端電腦複製到收件端電腦,使得任一時刻訊息都「真正」只存在其中一台上

收送訊息的 Message Endpoint 其實一直都在交易中,即使它們沒意識到send 方法在交易中把訊息加進通道,好把它與同時被加入或移出的其他訊息隔離;receive 方法也用交易,防止其他點對點收件端拿到同一則訊息。

邊欄:訊息交易與 ACID

交易常被描述為具備 ACID(原子性、一致性、隔離性、持久性)。

  • 只有 Guaranteed Delivery 的交易是持久的
  • 訊息依定義就是原子的
  • 但所有訊息交易都必須一致且隔離——訊息不可能「有點在通道裡」,它要嘛在、要嘛不在;而應用程式的收送也必須與其他執行緒與應用程式的收送隔離

什麼時候內部交易不夠用#

訊息系統的內部交易,對「只想送或收單一則訊息」的客戶端已經足夠且方便。但應用程式可能需要更寬的交易,來協調數則訊息、或把訊息與其他資源協調起來

  • 送-收訊息對——收一則訊息、送另一則。例如 Request-Reply 情境,或實作 Message Router、Message Translator 這類過濾器時。
  • 訊息群組——送或收一組相關訊息,例如 Message Sequence
  • 訊息/資料庫協調——把收送訊息與更新資料庫結合起來(常由 Channel Adapter 實作)。例如應用程式接收並處理訂購商品的訊息時,也得更新商品庫存資料庫;同樣地,Document Message 的寄件端可能想在訊息成功送出時才刪掉已持久化的文件,而收件端可能想在訊息真正被視為消費之前先把文件持久化。
  • 訊息/工作流程協調——用一對 Request-Reply 訊息執行一個工作項目,並用交易確保「請求沒送出就不取得工作項目」「回覆沒收到就不完成或中止工作項目」

這些情境需要更大的原子交易,牽涉的不只是單一則訊息,還可能包含訊息系統以外的其他交易性儲存。有交易,才能在「一部分成功(例如收到訊息)、另一部分失敗(例如更新資料庫或送出另一則訊息)」時,把全部回滾成從未發生過,讓應用程式重試。

解法#

寄件端與收件端都可以是交易性的

  • 對寄件端而言,訊息在提交交易之前不會「真正」被加進通道
  • 對收件端而言,訊息在提交交易之前不會「真正」從通道移除

交易式接收的價值#

有了交易式接收,應用程式可以在不真正把訊息從佇列移除的情況下接收它

  • 此時若應用程式崩潰,復原後訊息仍在佇列上,不會遺失
  • 收到訊息後應用程式處理它;確定要消費它時才提交交易,成功提交才把訊息從通道移除
  • 這之後若應用程式崩潰,復原時訊息已不在通道上——所以應用程式最好真的已經處理完了

四種情境的具體作法#

送-收訊息對#

  • 怎麼做:開啟交易 → 接收並處理第一則訊息 → 建立並送出第二則訊息 → 提交
  • 效果第一則訊息在第二則成功加入通道之前不會被移除
  • 交易型別:若兩個通道屬於同一套訊息系統,這是簡單交易;若由兩套不同的訊息系統管理(例如透過 Messaging Bridge),就是協調兩套系統的分散式交易

訊息群組#

  • 怎麼做:開啟交易 → 送出或接收群組中的所有訊息 → 提交
  • 效果:送出時,在全部成功送出之前沒有任何一則被加進通道;接收時,在全部收到之前沒有任何一則被移除
  • 交易型別:因為都在單一通道上,由單一訊息系統管理,屬於簡單交易

在許多訊息系統實作中,在單一交易中送出一組訊息,能確保它們在通道另一端以送出的順序被接收。

訊息/資料庫協調#

  • 怎麼做:開啟交易 → 接收訊息 → 更新資料庫 → 提交;或更新資料庫、送出訊息回報變更 → 提交
  • 效果資料庫沒更新,訊息就不會被移除(或訊息送不出去,資料庫變更就不算數)
  • 交易型別:訊息系統與資料庫各有自己的交易管理員,因此是分散式交易

訊息/工作流程協調#

  • 怎麼做:開啟交易 → 取得工作項目 → 送出請求訊息 → 提交;另開一個交易 → 接收回覆訊息 → 完成或中止工作項目 → 提交
  • 效果請求沒送出,工作項目就不會被提交;工作項目沒更新,回覆就不會被移除
  • 交易型別:同樣是分散式交易

與事件驅動消費者不合#

消費者通常必須先提交「接收訊息」的交易,才把訊息交給應用程式。之後若應用程式檢視訊息後決定不想消費它、或遇到錯誤想回滾接收動作,它做不到——因為它拿不到那個交易

因此事件驅動消費者不論客戶端是否具交易性,行為都差不多。

分散式交易的支援#

訊息系統有能力參與分散式交易,但有些實作可能不支援

  • 在 JMS 中,供應者可以扮演 XA 資源、參與 JTA 交易(由 javax.jms 中的 XA 類別,特別是 javax.jms.XASession,以及 javax.transaction.xa 套件定義)。JMS 規格建議客戶端不要自己處理分散式交易,應用程式應使用 J2EE 應用伺服器提供的分散式交易支援。
  • MSMQ 同樣能參與 XA 交易,在 .NET 中透過 MessageQueue.Transactional 屬性與 MessageQueueTransaction 類別暴露。

與其他模式的搭配#

  • 搭配良好:Request-Reply、Pipes and Filters 中的訊息過濾器、Message Sequence、Channel Adapter;Event Message 的接收端也可能想在完全移除訊息之前先處理完事件。
  • 搭配不佳Event-Driven ConsumerMessage Dispatcher
  • 可能出問題Competing Consumers
  • 搭配良好:單一個 Polling Consumer
範例:JMS 交易式 session、.NET 交易性佇列與 MSMQ 交易式過濾器

JMS 交易式 Session#

在 JMS 中,客戶端在建立 session 時決定自己是否具交易性:

Connection connection = // Get the connection
Session session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);

第一個參數設為 true 就讓 session 具交易性。之後收送都必須顯式提交才算數:

Queue queue = // Get the queue
MessageConsumer consumer = session.createConsumer(queue);
Message message = consumer.receive();

此時訊息只在這個消費者的交易視角中被消費了;對其他擁有自己交易視角的消費者而言,訊息仍然可用

session.commit();

提交若沒拋例外,消費者的交易視角就成為訊息系統的視角,訊息才真正被視為已消費。

.NET 交易性佇列#

在 .NET 中佇列預設不具交易性,因此要用交易式客戶端就得在建立時指定:

MessageQueue.Create("MyQueue", true);

佇列具交易性之後,每次客戶端動作(send 或 receive)都可以是交易性或非交易性的:

MessageQueue queue = new MessageQueue("MyQueue");
MessageQueueTransaction transaction = new MessageQueueTransaction();
transaction.Begin();
Message message = queue.Receive(transaction);
transaction.Commit();

雖然客戶端已經收到訊息,訊息系統直到客戶端成功提交交易,才讓該訊息在佇列上不再可用。

MSMQ 交易式過濾器#

以下範例把 Pipes and Filters 中的基本過濾器強化成交易式的,實作送-收訊息對情境:在同一個交易中接收與送出訊息

其實只要加幾行程式碼:用 MessageQueueTransaction 管理交易——消費輸入訊息前開啟交易、發布輸出訊息後提交;一旦發生例外就中止交易,回滾所有訊息消費與發布動作,並把輸入訊息退回佇列供其他消費者取用

public class TransactionalFilter
{
    protected MessageQueue inputQueue;
    protected MessageQueue outputQueue;
    protected Thread receiveThread;
    protected bool stopFlag = false;

    public TransactionalFilter (MessageQueue inputQueue, MessageQueue outputQueue)
    {
        this.inputQueue = inputQueue;
        this.inputQueue.Formatter = new System.Messaging.XmlMessageFormatter(
            new String[] {"System.String,mscorlib"});
        this.outputQueue = outputQueue;
    }

    public void Process()
    {
        ThreadStart receiveDelegate = new ThreadStart(this.ReceiveMessages);
        receiveThread = new Thread(receiveDelegate);
        receiveThread.Start();
    }

    private void ReceiveMessages()
    {
        MessageQueueTransaction myTransaction = new MessageQueueTransaction();

        while (!stopFlag)
        {
            try
            {
                myTransaction.Begin();
                Message inputMessage = inputQueue.Receive(myTransaction);
                Message outputMessage = ProcessMessage(inputMessage);
                outputQueue.Send(outputMessage, myTransaction);
                myTransaction.Commit();
            }
            catch (Exception e)
            {
                Console.WriteLine(e.Message + " - Transaction aborted ");
                myTransaction.Abort();
            }
        }
    }

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

怎麼驗證它真的有效?#

作者子類別化出一個名副其實的 RandomlyFailingFilter:對每則消費的訊息抽一個 0 到 10 的隨機數,小於 3 就拋出例外

若把這個過濾器實作在非交易式的基本過濾器之上,我們大約會遺失三分之一的訊息。

public class RandomlyFailingFilter : TransactionalFilter
{
    Random rand = new Random();

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

    protected override Message ProcessMessage(Message m)
    {
        string text = (string)m.Body;
        Console.WriteLine("Received Message: " + text);

        if (rand.Next(10) < 3)
        {
            Console.WriteLine("EXCEPTION");
            throw (new ArgumentNullException());
        }
        if (text == "end")
            stopFlag = true;
        return(m);
    }
}

測試載具發布一串訊息到輸入佇列,再確認能從輸出佇列以正確順序收到全部訊息。