這是同一個請求/回覆範例的 .NET 與 C# 版本。它示範如何實作 Request-Reply,並示範無效訊息如何被改道送往特殊通道。

本範例以 Microsoft .NET Framework SDK 開發,在安裝了 MSMQ 的 Windows XP 電腦上執行。

組成#

  1. Requestor——一個 Message Endpoint,送出請求訊息並等待收到回覆
  2. 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$\ReplyQueue

Replier 收到後重送到無效訊息佇列:

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