1. 順次消費の仕組みと課題
RocketMQでは、メッセージキューごとに順次消費を保証します。例えば、あるトピックに4つのキューがある場合、同一キュー内でのみメッセージの順序が維持されますが、異なるキュー間では順序保証がありません。
順次消費の実装には次の3つのフェーズがあります:
- キューの再バランス処理: RebalanceServiceスレッドが定期的にキューの割り当てを計算します。新規に割り当てられたキューに対してはBrokerにロックを取得し、そのキューのメッセージを取得できるようにします。
- メッセージの取得: PullMessageServiceがpullRequestQueueからタスクを取り出し、Brokerからメッセージを取得します。取得したメッセージはProcessQueueに格納され、消費スレッドプールに送られます。
- 順次消費プロセス: 同一キュー内のメッセージはロックによりシングルスレッドで処理され、メッセージの順序が保証されます。
2. パフォーマンスのボトルネック
RocketMQの順次消費では、以下の3つのロックが発生します:
- 再バランス後のキュー割り当て時にBrokerでロックを取得
- メッセージ消費前にMessageQueueのロックを取得
- ProcessQueueのロックを取得
これらのロックにより、並列度がキュー数に制限され、高パフォーマンスな処理が難しい状況になります。
3. 順次消費の最適化アプローチ
実際の業務では、関連順序性(同じキーに対する順序保証)が必要なケースがほとんどです。たとえば、銀行口座の残高変更通知では、同一ユーザーの複数メッセージが順序通りに届けば十分であり、他ユーザーとの順序は不要です。
この観点から、以下の最適化モデルを設計できます:
- 同じキーを持つメッセージは同一スレッドで処理(シングルスレッド)
- 異なるキーのメッセージは別スレッドで並列処理(マルチスレッド)
4. 新しい順次消費モデルの実装
以下に、RocketMQ 4.6をベースにした改良モデルの実装概要を示します。
4.1 クラス構成
DefaultMQLitePushConsumer:PullコンシューマーをベースにPushモデルを実装ConsumeMessageQueueService:キュー単位の消費サービスAbstractConsumeMessageService:メッセージルーティング戦略を定義AbstractConsumerTask:メッセージ消費の基本処理を定義
4.2 コンシューマースレッドプールの作成
private void startConsumerThreads() {
String threadPrefix = isOrderly ? "OrderlyConsumerThread_" : "ConcurrentConsumerThread_";
AtomicInteger threadNumIndex = new AtomicInteger(0);
consumerThreadGroup = new ThreadPoolExecutor(threadCount, threadCount, 0, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>(), r -> {
Thread t = new Thread(r);
t.setName(threadPrefix + threadNumIndex.incrementAndGet());
return t;
});
msgByKeyBlockQueue = new ArrayList<>(threadCount);
consumerRunningTasks = new ArrayList<>(threadCount);
for (int i = 0; i < threadCount; i++) {
msgByKeyBlockQueue.add(new LinkedBlockingQueue<>());
AbstractConsumerTask task = isOrderly
? new OrderlyConsumerTask(this, msgByKeyBlockQueue.get(i), listener)
: new ConcurrentlyConsumerTask(this, msgByKeyBlockQueue.get(i), listener);
consumerRunningTasks.add(task);
consumerThreadGroup.submit(task);
}
}
4.3 コンシューマースレッドの実行処理
public void run() {
try {
while (isRunning) {
List<MessageExt> msgs = new ArrayList<>(consumer.getBatchSize());
while (msgQueue.drainTo(msgs, consumer.getBatchSize()) <= 0) {
Thread.sleep(20);
}
doTask(msgs);
}
} catch (Exception e) {
logger.error("Error in consumer task", e);
}
}
4.4 メッセージのルーティング戦略
public class ConsumeMessageQueueOrderlyService extends AbstractConsumeMessageService {
private final String NO_KEY = "__nokey";
public ConsumeMessageQueueOrderlyService(DefaultMQLitePushConsumer consumer, MessageQueue messageQueue) {
super(consumer, messageQueue);
}
@Override
protected int selectTaskQueue(MessageExt msg, int totalQueues) {
String key = msg.getKeys();
if (key == null || key.isEmpty()) {
key = NO_KEY;
}
return Math.abs(key.hashCode()) % totalQueues;
}
}
5. 改良モデルの利点
- 従来のモデルでは並列度がキュー数に依存
- 改良モデルでは、メッセージのキーに基づくスレッド選択により、キュー数に依存せずスレッド数に比例して並列度が向上