脈絡:作者在做那個簡單的 JMS 請求/回覆範例時,遇到一個簡單卻有意思的問題。

範例由一個送訊息給 Replier 並等待回應的 Requestor 組成,使用 RequestQueueReplyQueue 兩個點對點通道。

「我們在程式碼裡貼滿了 println 好看清楚發生了什麼。先啟動 Replier、再啟動 Requestor——然後怪事發生了:Requestor 的主控台宣稱在 Replier 承認收到請求之前就已經收到回應了。 主控台輸出延遲嗎?

想不出好主意,我們決定關掉 Replier 再跑一次 Requestor。奇怪的是,我們仍然收到了請求的回應! 魔法嗎?不,只是持久化訊息的副作用——ReplyQueue 上有一則多餘的訊息。Requestor 啟動時把新訊息放上 RequestQueue,然後立刻取回了那則躺在 ReplyQueue 上的多餘回覆訊息,而我們從沒察覺那並不是剛才那個請求的回覆。等 Replier 收到請求訊息後,它又在 ReplyQueue 上放了一則新的回覆訊息,於是下一次測試又重演同樣的『魔法』

持久化、非同步的訊息傳遞,即使在最簡單的情境下也能耍得你團團轉,實在令人驚奇(或者說令人抓狂)。

問題的根源#

Message Channel 被設計成即使接收元件不可用也要可靠地遞送訊息,為此通道必須沿途把訊息持久化。

這些訊息讓我們在待處理訊息被消費完之前無法把測試資料送進系統。若那些待處理訊息是價值數百萬美元的訂單,這是好事;但若我們正在測試或除錯,而通道裡塞滿查詢或回覆訊息,那就夠讓人頭痛的了。

幾種緩解方式(但都不夠)#

用 Correlation Identifier#

在上面那個簡單例子中,若我們用了 Correlation Identifier,除錯的痛苦本可減輕:Requestor 就會認出「進來的訊息其實不是剛才那個請求的回應」,然後把這則舊回覆丟棄、或路由到 Invalid Message Channel——那實際上就把卡住的訊息清掉了。

但在其他情境下,要偵測重複或不想要的訊息就沒那麼容易。例如若某則訊息格式錯誤而導致接收端失敗,那麼在那則「壞訊息」被移除之前,接收端根本無法重新啟動——因為它一啟動就會立刻再次失敗。

當然,這個例子中接收端的缺陷本身該被修正(任何格式錯誤的訊息都不該讓元件掛掉),但把那則訊息移除,能讓系統在缺陷修好之前先恢復運作。

用臨時通道#

另一種避免殘留訊息的方式是使用臨時通道(例如 JMS 的 createTemporaryQueue)。這類通道專為請求-回覆應用設計,應用程式一關閉與訊息系統的連線,所有訊息就消失。

靠交易管理?#

很容易以為交易管理能消除「多餘訊息」的情境——因為訊息消費、處理與發布都被涵蓋在一個交易中:元件若在處理途中中止,訊息就不算被消費;回覆訊息也要等到元件送出最終 commit 才會被發布。

在那個簡單的請求-回覆例子中,程式設計師的錯誤可能導致 Requestor 根本沒去讀 ReplyQueue 上的回應。結果是:即使有交易性,訊息仍然卡在那個通道上,造成前述的症狀。

解法#

  • 簡單的 Channel Purger 就是把通道上所有訊息移除。對「想把系統重置到一致狀態」的測試情境而言,這可能就夠了。
  • 若是在除錯正式系統,我們可能得依特定準則(message ID、特定訊息欄位的值)移除個別訊息或一組訊息

圖 11-14:Channel Purger 解法示意

刪掉還是留著#

許多情況下讓 Channel Purger 直接刪除訊息就好。

這可能還包括在重新注入前編輯訊息內容的需求——這類功能結合了 Message Store 與 Channel Purger 的部分特性。

範例:用 JMS 實作 Channel Purger

這個簡單範例單純移除通道上的所有訊息,並參照兩個外部類別:

  • JmsEndpoint——任何 JMS 參與者的基底類別,提供預先初始化的 ConnectionSession 實例變數
  • JndiUtil——實作透過 JNDI 查找 JMS 物件的輔助函式
import javax.jms.JMSException;
import javax.jms.MessageConsumer;
import javax.jms.Queue;

public class ChannelPurger extends JmsEndpoint
{
    public static void main(String[] args)
    {
        if (args.length != 1) {
            System.out.println("Usage: java ChannelPurger <queue_name>");
            System.exit(1);
        }
        String queueName = new String(args[0]);
        System.out.println("Purging queue " + queueName);

        ChannelPurger purger = new ChannelPurger();

        purger.purgeQueue(queueName);
    }

    private void purgeQueue(String queueName)
    {
        try {
            initialize();
            connection.start();
            Queue queue = (Queue) JndiUtil.getDestination(queueName);

            MessageConsumer consumer = session.createConsumer(queue);

            while (consumer.receiveNoWait() != null)
                System.out.print(".");
            connection.stop();
        } catch (Exception e) {
            System.out.println("Exception occurred: " + e.toString());
        } finally {
            if (connection != null) {
                try {
                    connection.close();
                } catch (JMSException e) {
                    // ignore
                }
            }
        }
    }
}

關鍵就是那個 while (consumer.receiveNoWait() != null) 迴圈——它不斷以非阻塞方式接收訊息,直到佇列被清空為止。 >