脈絡:應用程式正在使用訊息傳遞,但它處理訊息的速度趕不上訊息被加進通道的速度。
循序消費太慢#
訊息在通道上循序抵達,消費者自然傾向循序處理。
原因可能是:
- 通道上有多個寄件端
- 網路中斷造成訊息積壓,恢復後一次全部湧入
- 接收端故障造成積壓
- 消費並執行一則訊息,所需的工夫遠大於建立與送出它
多通道不是好解法#
應用程式可以用多個通道,但一個通道可能塞爆、另一個卻閒著,而且寄件端不知道該用哪一個等價通道。
多通道的優點是能有多個消費者(一通道一個)並行處理;但即使如此,通道數量仍會限制吞吐量。
我們需要的是讓一個通道能有多個消費者的方式。
解法#
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,並讓它跑在自己的執行緒中。 >