脈絡:應用程式正在使用訊息傳遞,但它處理訊息的速度趕不上訊息被加進通道的速度

循序消費太慢#

訊息在通道上循序抵達,消費者自然傾向循序處理。

原因可能是:

  • 通道上有多個寄件端
  • 網路中斷造成訊息積壓,恢復後一次全部湧入
  • 接收端故障造成積壓
  • 消費並執行一則訊息,所需的工夫遠大於建立與送出它

多通道不是好解法#

應用程式可以用多個通道,但一個通道可能塞爆、另一個卻閒著,而且寄件端不知道該用哪一個等價通道

多通道的優點是能有多個消費者(一通道一個)並行處理;但即使如此,通道數量仍會限制吞吐量。

我們需要的是讓一個通道能有多個消費者的方式。

解法#

Competing Consumers 是為了從單一個 Point-to-Point Channel 接收訊息而建立的多個消費者。通道遞送訊息時,任何一個消費者都有可能收到它——由訊息系統的實作決定實際是誰收到,效果上消費者們在競爭成為接收者

圖 10-6:Competing Consumers 解法示意

運作方式#

  • 每個競爭消費者跑在自己的執行緒中,才能並行消費
  • 通道遞送訊息時,訊息系統的交易控制確保只有其中一個消費者成功收到它
  • 那個消費者處理訊息的同時,通道可以遞送其他訊息給其他消費者並行處理

於是瓶頸從「消費者處理一則訊息要多久」變成「通道能多快把訊息餵給消費者」。消費者數量有限仍可能是瓶頸,但只要還有可用的運算資源,增加消費者就能緩解

圖 10-7:Competing Consumers 的循序圖

執行緒需求#

  • Polling Consumers每個消費者都得有自己的執行緒才能並行輪詢
  • Event-Driven Consumers訊息系統必須為每個並行消費者使用一條執行緒——那條執行緒用來把訊息交給消費者,消費者也用它來處理訊息

訊息系統的成熟度差異#

這會讓交易式客戶端很沒效率#

允許多個消費者嘗試消費同一則訊息的訊息系統,會讓 Transactional Client 極度沒效率:客戶端以為自己拿到訊息、消費它、花力氣處理它,然後嘗試提交卻失敗(因為訊息已被競爭者消費掉)

頻繁地做完工作又回滾,反而傷害吞吐量——而這個解法的重點本來就是提升吞吐量。

因此交易式競爭消費者的效能必須仔細測量:它在不同訊息系統實作與設定下的表現可能差異極大。

跨應用程式的擴散#

Competing Consumers 不只能把負載分散到單一應用程式的多條消費者執行緒上,也能把消費負載分散到多個應用程式

若一個應用程式消費得不夠快,多個消費者應用程式(每個又各自使用多條消費執行緒)可以一起上。「多個應用程式跑在多台電腦上、用多條執行緒消費訊息」的能力,提供了幾乎無上限的訊息處理容量——唯一的限制是訊息系統把訊息從通道遞送到消費者的能力。

競爭消費者之間的協調取決於各訊息系統的實作。若客戶端想自己實作這份協調,就該使用 Message Dispatcher。

範例:簡單的 JMS 競爭消費者

一個外部的 driver/manager 物件(未列出)會跑好幾個這種消費者,各自在自己的執行緒中,並呼叫 stopRunning() 讓它停止。

JMS 規格並未規定並行 QueueReceiver(也就是 Competing Consumers)該如何運作,甚至沒有要求這種作法必須可行。因此使用這個技巧的應用程式不能假設具可攜性,在不同 JMS 供應者上可能表現不同。

消費者類別實作 Runnable 以便在自己的執行緒中執行。所有消費者共用同一個 Connection,但各自建立自己的 Session

import javax.jms.Connection;
import javax.jms.Destination;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.MessageConsumer;
import javax.jms.Session;
import javax.naming.NamingException;

public class CompetingConsumer implements Runnable {

    private int performerID;
    private MessageConsumer consumer;
    private boolean isRunning;

    protected CompetingConsumer() {
        super();
    }

    public static CompetingConsumer newConsumer(int id, Connection connection, String queueName)
            throws JMSException, NamingException {
        CompetingConsumer consumer = new CompetingConsumer();
        consumer.initialize(id, connection, queueName);
        return consumer;
    }

    protected void initialize(int id, Connection connection, String queueName)
            throws JMSException, NamingException {
        performerID = id;
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Destination dispatcherQueue = JndiUtil.getDestination(queueName);
        consumer = session.createConsumer(dispatcherQueue);
        isRunning = true;
    }

    public void run() {
        try {
            while (isRunning())
                receiveSync();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }

    private synchronized boolean isRunning() {
        return isRunning;
    }

    public synchronized void stopRunning() {
        isRunning = false;
    }

    private void receiveSync() throws JMSException, InterruptedException {
        Message message = consumer.receive();
        if (message != null)
            processMessage(message);
    }

    private void processMessage(Message message) throws JMSException, InterruptedException {
        int id = message.getIntProperty("cust_id");
        System.out.println(System.currentTimeMillis() + ": Performer #" + performerID
            + " starting; message ID " + id);
        Thread.sleep(500);
        System.out.println(System.currentTimeMillis() + ": Performer #" + performerID
            + " processing.");
        Thread.sleep(500);
        System.out.println(System.currentTimeMillis() + ": Performer #" + performerID
            + " finished.");
    }
}

實作一個簡單的競爭消費者很容易——主要訣竅就是把消費者做成 Runnable,並讓它跑在自己的執行緒中。 >