這是同一個請求/回覆範例的 .NET 與 C# 版本。它示範如何實作 Request-Reply,並示範無效訊息如何被改道送往特殊通道。
本範例以 Microsoft .NET Framework SDK 開發,在安裝了 MSMQ 的 Windows XP 電腦上執行。
組成#
- Requestor——一個 Message Endpoint,送出請求訊息並等待收到回覆
- Replier——一個 Message Endpoint,等待收到請求訊息;收到後以回覆訊息回應
Requestor 與 Replier 各自作為獨立的 .NET 程式執行,這使得通訊成為分散式的。
範例假設訊息系統已定義三個佇列:
.\private$\RequestQueue——Requestor 用來把請求訊息送給 Replier.\private$\ReplyQueue——Replier 用來把回覆訊息送回 Requestor.\private$\InvalidQueue——雙方收到無法解讀的訊息時,把它移到這裡
跑一遍會看到什麼#
啟動 Requestor:
Sent request
Time: 09:11:09.165342
Message ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\149
Correl. ID:
Reply to: .\private$\ReplyQueue
Contents: Hello world.同樣地,即使 Replier 還沒啟動,這一步仍然成功。
啟動 Replier:
Received request
Time: 09:11:09.375644
Message ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\149
Correl. ID: <n/a>
Reply to: FORMATNAME:DIRECT=OS:XYZ123\private$\ReplyQueue
Contents: Hello world.
Sent reply
Time: 09:11:09.956480
Message ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\150
Correl. ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\149
Reply to: <n/a>
Contents: Hello world.觀察重點與 JMS 版本一致:
- 請求在送出之後才被收到;兩邊的 message ID 相同,因為是同一則訊息;內容也相同
- 請求訊息中指定了回覆佇列——Return Address
- 回覆的 message ID 與請求不同(兩則獨立訊息);回覆的 reply-to 未指定(不預期再有回覆)
- 回覆的 correlation ID 等於請求的 message ID——Correlation Identifier
最後請求方收到回覆:
Received reply
Time: 09:11:10.156467
Message ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\150
Correl. ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\149
Reply to: <n/a>
Contents: Hello world.程式碼:Requestor
using System;
using System.Messaging;
public class Requestor
{
private MessageQueue requestQueue;
private MessageQueue replyQueue;
public Requestor(String requestQueueName, String replyQueueName)
{
requestQueue = new MessageQueue(requestQueueName);
replyQueue = new MessageQueue(replyQueueName);
replyQueue.MessageReadPropertyFilter.SetAll();
((XmlMessageFormatter)replyQueue.Formatter).TargetTypeNames =
new string[]{"System.String,mscorlib"};
}
public void Send()
{
Message requestMessage = new Message();
requestMessage.Body = "Hello world.";
requestMessage.ResponseQueue = replyQueue;
requestQueue.Send(requestMessage);
// ...印出訊息細節
}
public void ReceiveSync()
{
Message replyMessage = replyQueue.Receive();
// ...印出訊息細節
}
}建構子#
應用程式指定兩個佇列的路徑名稱(請求與回覆),requestor 用它們連上訊息系統:
- 用佇列名稱查出
MessageQueue(名稱是 MSMQ 資源的路徑名) - 設定回覆佇列的屬性過濾器,讓訊息被讀出時所有屬性也一併讀出
- 把佇列的 formatter 設為
XmlMessageFormatter,使訊息內容被解讀為字串
Send()#
- 建立
Message,內容設為 “Hello world.” - 把
ResponseQueue屬性設為回覆佇列——這是告訴回覆方怎麼送回回覆的 Return Address - 送出訊息
同樣地,印細節放在送出之後——message ID 由訊息系統設定,訊息真正送出前它並不存在。
ReceiveSync()#
執行佇列的 Receive() 取得訊息,它會同步阻塞直到有訊息被遞送並讀出——因此請求方是一個 Polling Consumer。
程式碼:Replier
using System;
using System.Messaging;
class Replier {
private MessageQueue invalidQueue;
public Replier(String requestQueueName, String invalidQueueName)
{
MessageQueue requestQueue = new MessageQueue(requestQueueName);
invalidQueue = new MessageQueue(invalidQueueName);
requestQueue.MessageReadPropertyFilter.SetAll();
((XmlMessageFormatter)requestQueue.Formatter).TargetTypeNames =
new string[]{"System.String,mscorlib"};
requestQueue.ReceiveCompleted += new ReceiveCompletedEventHandler(OnReceiveCompleted);
requestQueue.BeginReceive();
}
public void OnReceiveCompleted(Object source, ReceiveCompletedEventArgs asyncResult)
{
MessageQueue requestQueue = (MessageQueue)source;
Message requestMessage = requestQueue.EndReceive(asyncResult.AsyncResult);
try
{
// ...印出請求細節
string contents = requestMessage.Body.ToString();
MessageQueue replyQueue = requestMessage.ResponseQueue;
Message replyMessage = new Message();
replyMessage.Body = contents;
replyMessage.CorrelationId = requestMessage.Id;
replyQueue.Send(replyMessage);
// ...印出回覆細節
}
catch ( Exception ) {
// 無效訊息:保存原始 ID 後改送到無效訊息佇列
requestMessage.CorrelationId = requestMessage.Id;
invalidQueue.Send(requestMessage);
}
requestQueue.BeginReceive();
}
}與 Requestor 的兩個差異#
- 回覆方不查回覆佇列——它不假設自己永遠往那個佇列送回覆,而是讓請求訊息告訴它送去哪。
- 回覆方是 Event-Driven Consumer,因此設定了一個
ReceiveCompletedEventHandler。訊息一被遞送到請求佇列,訊息系統就自動呼叫OnReceiveCompleted。
OnReceiveCompleted 的處理#
source就是請求佇列這個MessageQueue- 用佇列的
EndReceive取得訊息本身 - 接著實作 Return Address:取出請求訊息的 response-queue 屬性值,用它參照到正確的
MessageQueue
重點同樣是:回覆方沒有把回覆佇列寫死,它用的是每則請求訊息各自指定的那一個。
- 建立回覆訊息時,把 correlation-id 設成與請求的 message-id 相同——Correlation Identifier
- 若訊息收得到卻處理不成功、拋出 Exception,回覆方就把它重送到無效訊息佇列,過程中把新訊息的 correlation id 設為原訊息的 message id
- 處理完後執行
BeginReceive開始監聽下一則訊息
無效訊息範例#
同樣地,InvalidMessenger 刻意在請求通道上送出格式不正確的訊息:
Sent request
Type: 768
Time: 09:39:44.223729
Message ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\168
Correl. ID: 00000000-0000-0000-0000-000000000000\0
Reply to: .\private$\ReplyQueueReplier 收到後重送到無效訊息佇列:
Invalid message detected
Type: 768
Message ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\168
Correl. ID: <n/a>
Sent to invalid message queue
Type: 768
Message ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\169
Correl. ID: 8b0fc389-f21f-423b-9eaa-c3a881a34808\168同一個洞見:訊息被移到無效訊息佇列時其實是被重送,因此拿到新的 message ID。所以回覆方一判定訊息無效,就把主 ID 複製到 correlation ID,保留原始 ID 的紀錄。
結論#
與 JMS 版本相同:兩個 Message Endpoint 用 Request-Reply 交換訊息;請求用 Return Address 指定回覆通道,回覆用 Correlation Identifier 指明對應的請求;Requestor 是 Polling Consumer,Replier 是 Event-Driven Consumer;請求與回覆佇列都是 Datatype Channel,型別不符的訊息會被改道送往 Invalid Message Channel。