這是一個用 JMS 實作的簡單訊息傳遞範例。它示範如何實作 Request-Reply——請求方送出請求,回覆方收到請求並傳回回覆,請求方再收到回覆——同時也示範無效訊息如何被改道送往特殊通道。

本範例以 JMS 1.1 開發,並在 J2EE 1.4 參考實作上執行。

組成#

範例由兩個主要類別構成:

  1. Requestor——一個 Message Endpoint,送出請求訊息並等待收到回覆
  2. Replier——一個 Message Endpoint,等待收到請求訊息;收到後以回覆訊息回應

Requestor 與 Replier 各自跑在不同的 JVM 中,這正是讓這場通訊成為分散式的原因。

範例假設訊息系統已定義三個佇列:

  • jms/RequestQueue——Requestor 用來把請求訊息送給 Replier
  • jms/ReplyQueue——Replier 用來把回覆訊息送回 Requestor
  • jms/InvalidMessages——雙方收到無法解讀的訊息時,把它移到這裡

跑一遍會看到什麼#

在命令列視窗啟動 Requestor,它印出:

Sent request
        Time:        1048261736520 ms
        Message ID:  ID:_XYZ123_1048261766139_6.2.1.1
        Correl. ID:  null
        Reply to:    com.sun.jms.Queue: jms/ReplyQueue
        Contents:    Hello world.

注意:即使 Replier 根本還沒啟動、無法接收請求,這件事仍然成功了——這正是非同步訊息的價值。

在另一個視窗啟動 Replier:

Received request
        Time:        1048261766790 ms
        Message ID:  ID:_XYZ123_1048261766139_6.2.1.1
        Correl. ID:  null
        Reply to:    com.sun.jms.Queue: jms/ReplyQueue
        Contents:    Hello world.
Sent reply
        Time:        1048261766850 ms
        Message ID:  ID:_XYZ123_1048261758148_5.2.1.1
        Correl. ID:  ID:_XYZ123_1048261766139_6.2.1.1
        Reply to:    null
        Contents:    Hello world.

這段輸出有幾個值得注意的地方:

  • 請求的送出與接收時間戳——請求是在送出之後才被收到
  • 兩邊的 message ID 相同,因為那是同一則訊息
  • 內容 “Hello world.” 相同——這正是被傳輸的資料,兩邊當然必須一致
  • 請求訊息中指定了 jms/ReplyQueue 作為回覆的目的地——這是 Return Address 模式的實例

本例的請求相當陽春,它基本上是一則 Document Message;真實情境下的請求通常會是 Command Message

再對照「收到請求」與「送出回覆」:

  • 回覆是在請求被收到之後才送出的
  • 回覆的 message ID 與請求不同——因為請求與回覆是兩則各自獨立的訊息
  • 請求的內容被取出並放進了回覆
  • 回覆的 reply-to 目的地未指定,因為不預期還有回覆(回覆本身不使用 Return Address)
  • 回覆的 correlation ID 等於請求的 message ID——這是 Correlation Identifier 模式

最後回到第一個視窗,請求方收到了回覆:

Received reply
        Time:        1048261737060 ms
        Message ID:  ID:_XYZ123_1048261758148_5.2.1.1
        Correl. ID:  ID:_XYZ123_1048261766139_6.2.1.1
        Reply to:    null
        Contents:    Hello world.

請求方被設計成「送請求、收回覆、然後結束」,所以收到回覆後它就不再執行。回覆方則不知道什麼時候會有請求進來,因此永不停止——要停掉它,就到它的命令視窗按 Enter。

程式碼:Requestor
import javax.jms.Connection;
import javax.jms.Destination;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageConsumer;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.naming.NamingException;

public class Requestor {

    private Session session;
    private Destination replyQueue;
    private MessageProducer requestProducer;
    private MessageConsumer replyConsumer;
    private MessageProducer invalidProducer;

    protected Requestor() {
        super();
    }

    public static Requestor newRequestor(Connection connection, String requestQueueName,
            String replyQueueName, String invalidQueueName)
            throws JMSException, NamingException {

        Requestor requestor = new Requestor();
        requestor.initialize(connection, requestQueueName, replyQueueName, invalidQueueName);
        return requestor;
    }

    protected void initialize(Connection connection, String requestQueueName,
            String replyQueueName, String invalidQueueName)
            throws NamingException, JMSException {

        session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);

        Destination requestQueue = JndiUtil.getDestination(requestQueueName);
        replyQueue = JndiUtil.getDestination(replyQueueName);
        Destination invalidQueue = JndiUtil.getDestination(invalidQueueName);

        requestProducer = session.createProducer(requestQueue);
        replyConsumer = session.createConsumer(replyQueue);
        invalidProducer = session.createProducer(invalidQueue);
    }

    public void send() throws JMSException {
        TextMessage requestMessage = session.createTextMessage();
        requestMessage.setText("Hello world.");
        requestMessage.setJMSReplyTo(replyQueue);
        requestProducer.send(requestMessage);
        // ...印出訊息細節
    }

    public void receiveSync() throws JMSException {
        Message msg = replyConsumer.receive();
        if (msg instanceof TextMessage) {
            TextMessage replyMessage = (TextMessage) msg;
            // ...印出訊息細節
        } else {
            // 無效訊息:先保存原始 ID,再改送到無效訊息佇列
            msg.setJMSCorrelationID(msg.getJMSMessageID());
            invalidProducer.send(msg);
        }
    }
}

initialize 做了什麼#

應用程式提供一個到訊息系統的 Connection,外加三個佇列的 JNDI 名稱(請求、回覆、無效訊息):

  • Connection 建立一個 Session
  • 用佇列名稱查出各個 Destination(名稱是 JNDI 識別碼,由 JndiUtil 執行查詢)
  • 建立三個端點物件:送請求的 producer、收回覆的 consumer、以及把訊息移往無效佇列的 producer

send()#

  • 建立一則 TextMessage,內容設為 “Hello world.”
  • 把訊息的 reply-to 屬性設為回覆佇列——這就是告訴回覆方「回覆怎麼送回來」的 Return Address
  • requestProducer 送出訊息(producer 連的是請求佇列,訊息因此送到那裡)

印出訊息細節的動作放在送出之後——因為 message ID 是由訊息系統設定的,訊息真正送出前它並不存在

receiveSync()#

  • replyConsumer 接收回覆。receive()同步阻塞直到有訊息被遞送到佇列並被讀出——因此請求方是一個 Polling Consumer(方法名稱 receiveSync() 即由此而來)
  • 訊息應該是 TextMessage;若是,就取出內容並印出細節
  • 若不是,訊息就無法處理。與其直接丟棄,請求方把它重送到無效訊息佇列

重送會改變訊息的 message ID,因此重送前請求方會先把原本的 message ID 存進 correlation ID(見 Correlation Identifier)。

程式碼:Replier
import javax.jms.Connection;
import javax.jms.Destination;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageConsumer;
import javax.jms.MessageListener;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.TextMessage;
import javax.naming.NamingException;

public class Replier implements MessageListener {

    private Session session;
    private MessageProducer invalidProducer;

    protected Replier() {
        super();
    }

    public static Replier newReplier(Connection connection, String requestQueueName,
            String invalidQueueName) throws JMSException, NamingException {

        Replier replier = new Replier();
        replier.initialize(connection, requestQueueName, invalidQueueName);
        return replier;
    }

    protected void initialize(Connection connection, String requestQueueName,
            String invalidQueueName) throws NamingException, JMSException {

        session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Destination requestQueue = JndiUtil.getDestination(requestQueueName);
        Destination invalidQueue = JndiUtil.getDestination(invalidQueueName);

        MessageConsumer requestConsumer = session.createConsumer(requestQueue);
        MessageListener listener = this;
        requestConsumer.setMessageListener(listener);

        invalidProducer = session.createProducer(invalidQueue);
    }

    public void onMessage(Message message) {
        try {
            if ((message instanceof TextMessage) && (message.getJMSReplyTo() != null)) {
                TextMessage requestMessage = (TextMessage) message;
                // ...印出請求細節

                String contents = requestMessage.getText();
                Destination replyDestination = message.getJMSReplyTo();
                MessageProducer replyProducer = session.createProducer(replyDestination);
                TextMessage replyMessage = session.createTextMessage();
                replyMessage.setText(contents);
                replyMessage.setJMSCorrelationID(requestMessage.getJMSMessageID());
                replyProducer.send(replyMessage);
                // ...印出回覆細節
            } else {
                // 無效訊息:保存原始 ID 後改送到無效訊息佇列
                message.setJMSCorrelationID(message.getJMSMessageID());
                invalidProducer.send(message);
            }
        } catch (JMSException e) {
            e.printStackTrace();
        }
    }
}

與 Requestor 的兩個差異#

  • 回覆方不查回覆佇列、也不為它建立 producer——因為它不假設自己永遠往那個佇列送回覆,而是讓請求訊息告訴它該送去哪
  • 回覆方是 Event-Driven Consumer,因此實作 MessageListener。訊息一被遞送到請求佇列,訊息系統就自動呼叫它的 onMessage

onMessage 的處理#

  • 請求訊息應該是 TextMessage,而且應該指定回覆佇列。不符合這兩項要求,就移到無效訊息佇列(與請求方相同)。
  • 符合的話,這裡就是回覆方實作 Return Address 的地方:它取出請求訊息的 reply-to 屬性值,用它在正確的佇列上建立 MessageProducer

重點在於:回覆方沒有把回覆佇列寫死,它用的是每則請求訊息各自指定的那一個。

  • 接著建立回覆訊息,並把回覆的 correlation-id 設成與請求的 message-id 相同——這就是 Correlation Identifier 模式。

無效訊息範例#

順帶看一下 Invalid Message Channel 模式的實例。我們需要的其中一個佇列名為 jms/InvalidMessages,它的存在是為了讓 JMS 客戶端(一個 Message Endpoint)在收到無法處理的訊息時,能把這則怪訊息移到特殊通道。

為了示範,我們設計了一個 InvalidMessenger 類別,它刻意在請求通道上送出一則格式不正確的訊息

請求通道和任何通道一樣是 Datatype Channel——請求的接收者預期請求是某種特定格式。InvalidMessenger 送的是不同格式的訊息;回覆方收到時認不出格式,於是把它移到無效訊息佇列。

在一個視窗跑 Replier、另一個跑 InvalidMessenger。後者送出時顯示:

Sent invalid message
       Type:        com.sun.jms.ObjectMessageImpl
       Time:        1048288516959 ms
       Message ID:  ID:_XYZ123_1048288516639_7.2.1.1
       Correl. ID:  null
       Reply to:    com.sun.jms.Queue: jms/ReplyQueue

這顯示訊息是 ObjectMessage 的實例(而回覆方期待的是 TextMessage)。Replier 收到並重送到無效訊息佇列:

Invalid message detected
       Type:        com.sun.jms.ObjectMessageImpl
       Message ID:  ID:_XYZ123_1048288516639_7.2.1.1
       Correl. ID:  null
Sent to invalid message queue
       Type:        com.sun.jms.ObjectMessageImpl
       Message ID:  ID:_XYZ123_1048287020267_6.2.1.2
       Correl. ID:  ID:_XYZ123_1048288516639_7.2.1.1

一個值得記住的洞見:訊息被移到無效訊息佇列時,其實是「被重送」,所以它拿到一個新的 message ID。

正因如此我們才套用 Correlation Identifier 模式——回覆方一判定訊息無效,就把它的主 ID 複製到 correlation ID,藉此保留原始 ID 的紀錄

結論#

我們看到了如何實作兩個 Message Endpoint(Requestor 與 Replier),用 Request-Reply 交換請求與回覆訊息:

  • 請求訊息用 Return Address 指定回覆該送到哪個佇列
  • 回覆訊息用 Correlation Identifier 指明它回應的是哪個請求
  • Requestor 以 Polling Consumer 接收回覆,Replier 以 Event-Driven Consumer 接收請求
  • 請求與回覆佇列都是 Datatype Channel——消費者收到型別不對的訊息時,會把它改道送往 Invalid Message Channel