JavaでのRocketMQ順次消費の性能向上方法

1. 順次消費の仕組みと課題

RocketMQでは、メッセージキューごとに順次消費を保証します。例えば、あるトピックに4つのキューがある場合、同一キュー内でのみメッセージの順序が維持されますが、異なるキュー間では順序保証がありません。

順次消費の実装には次の3つのフェーズがあります:

  1. キューの再バランス処理: RebalanceServiceスレッドが定期的にキューの割り当てを計算します。新規に割り当てられたキューに対してはBrokerにロックを取得し、そのキューのメッセージを取得できるようにします。
  2. メッセージの取得: PullMessageServiceがpullRequestQueueからタスクを取り出し、Brokerからメッセージを取得します。取得したメッセージはProcessQueueに格納され、消費スレッドプールに送られます。
  3. 順次消費プロセス: 同一キュー内のメッセージはロックによりシングルスレッドで処理され、メッセージの順序が保証されます。

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. 改良モデルの利点

  • 従来のモデルでは並列度がキュー数に依存
  • 改良モデルでは、メッセージのキーに基づくスレッド選択により、キュー数に依存せずスレッド数に比例して並列度が向上

タグ: RocketMQ Java メッセージキュー 並列処理 ロック最適化

7月21日 19:53 投稿