脈絡:訊息系統出於必要,內部使用交易行為。而讓外部客戶端能控制影響自身行為的交易範圍,可能是有價值的。
訊息系統本來就在用交易#
單一個 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 Consumer 與 Message 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);
}
}測試載具發布一串訊息到輸入佇列,再確認能從輸出佇列以正確順序收到全部訊息。