👨‍🎤

待機キューの実装による高トラフィック制御

に公開

大量のユーザーが特定のサービスに一斉にアクセスすると、システムのリソースが不足し、障害発生に繋がります。限られたリソースでサービスを安定して維持するためには、同時にアクセスできるユーザー数を制限する技術が必要です。

例えば、全体のユーザー数が1000万人のサービスでも、同時に利用できるのは10万人や100万人といった単位になるよう、システムを設計します。

今回は、この課題を解決するため、RedisのSorted Setを利用して、ユーザー待機キュー機能を実装します。

まず、実装の核となるRedisの「Sorted Set」について簡単に説明します。

Redisの「Sorted Set」のコマンド

1. ZADD: 待機列にユーザーを追加する

ZADDは "Z-ADD" と読み、「Z(Sorted Set)にメンバーをADD(追加)する」という意味です。

役割: Sorted Setに、新しいユーザーを時刻付きで追加します。
ZADD <待機列の名前> <時刻> <ユーザーID>
例: ZADD users:queue:event-sale:wait 1678886400 12345

これは、「users:queue:event-sale:waitという名前の待機列に、ユーザー12345を、1678886400という時刻で追加する」という命令です。

2. ZRANK: ユーザーの順番を確認する

ZRANKは "Z-RANK" と読み、「Z(Sorted Set)のメンバーのRANK(順位)を調べる」という意味です。

役割: Sorted Setの中で、特定のユーザーが前から何番目にいるかを教えてくれます。
ZRANK <待機列の名前> <ユーザーID>
例: ZRANK users:queue:event-sale:wait 12345

これは、「users:queue:event-sale:waitという待機列の中で、ユーザー12345は何番目ですか?」と問い合わせる命令です。

例:コンサートチケットの予約のサービス実装

例えば、有名なロックバンドのチケット発売日を例に考えてみましょう。この日、数万人のユーザーが一度にアクセスする可能性があり、正確なアクセス数を予測することは困難です。このような状況でサーバーダウンを防ぎ、限られたユーザーに安定してサービスを提供するため、前述の待機キューを実装します。

以下はシステムの流れを描いたシーケンス図です。

1. 待機列にユーザー登録。

ユーザーがアクセスするとき、見えるスクリーン。

これは、ユーザーが最初にアクセスした際に表示される待機画面です。システムはまずユーザーを待機キューに登録し、サービスへのアクセスが許可されたかを確認するため、3秒ごとにサーバーへ問い合わせを行います。以下は、この処理を行うクライアント側のJavascriptコードの一部です。

<script>
  async function registerAndPoll(queueName, userId) {
    try {
      // 1. まず待機列に登録を試みる
      const registerResponse = await fetch(`/api/v1/queue`, {
          method: 'POST',
          headers: { 'Content-Type': 'application/json' },
          body: JSON.stringify({ queue: queueName, userId: userId })
      });

      if (registerResponse.ok) {
          console.log("Successfully registered to the queue.");
      } else {
          // 既に登録済みの場合など、エラーが発生しても処理を続ける
          console.log("Already registered or failed to register. Proceeding to poll status.");
      }

    } catch (error) {
        console.error("Registration error:", error);
    } finally {
        // 2. 登録処理の成否にかかわらず、ポーリングを開始する
        startPolling(queueName, userId);
    }
  }
</script>

サーバー側では、クライアントからの登録リクエストを受け取ると、ユーザーIDとキュー名を使ってRedisのSorted Setのキーを生成し、ユーザーを待機リストに追加します。

public class UserQueueService {
    private final ReactiveRedisTemplate<String, String> reactiveRedisTemplate;

    // ZSETのキー: users:queue:{queueName}:wait
    // ユーザーの待機リスト。スコアは待機開始時間(unix timestamp)。
    private final String USER_QUEUE_WAIT_KEY = "users:queue:%s:wait";

    // SCANコマンドで使用するキーのパターン
    private final String USER_QUEUE_WAIT_KEY_FOR_SCAN = "users:queue:*:wait";

    // ZSETのキー: users:queue:{queueName}:proceed
    // 実際にアクセスが許可されたユーザーのリスト。
    private final String USER_QUEUE_PROCEED_KEY = "users:queue:%s:proceed";

    /**
     * 指定されたキューにユーザーを登録します。
     *
     * @param queue  キューの名称 (例: "event-2025-summer-sale")
     * @param userId 登録するユーザーのID
     * @return 登録後のユーザーの待機順位。既に登録済みの場合はエラーを返す。
     */
    public Mono<Long> registerWaitQueue(final String queue, final Long userId) {
        var unixTimestamp = Instant.now().getEpochSecond();
        // ZSETに追加
        return reactiveRedisTemplate.opsForZSet().add(USER_QUEUE_WAIT_KEY.formatted(queue), userId.toString(), unixTimestamp)
                .filter(i -> i)
                .switchIfEmpty(Mono.error(ErrorCode.QUEUE_ALREADY_REGISTERED_USER.build()))
                .flatMap(i -> reactiveRedisTemplate.opsForZSet().rank(USER_QUEUE_WAIT_KEY.formatted(queue), userId.toString()))
                .map(i -> i >= 0 ? i + 1 : i); // Redisのランクは0から始まるため、1を加算して人間が分かりやすい順位にする
    }
}

2. 待機列から順番を確認すること。

<script>
  function startPolling(queueName, userId) {
      const intervalId = setInterval(async () => {
          try {
              const queryParam = new URLSearchParams({queue: queueName, userId: userId});
              
              // 3. アクセス許可を確認
              const allowedResponse = await fetch(`/api/v1/queue/allowed?` + queryParam);
              const allowedData = await allowedResponse.json();

              if (allowedData.allowed === true) {
                  // 許可された場合、リダイレクト
                  console.log("Access granted! Redirecting...");
                  clearInterval(intervalId);
                  window.location.href = "http://127.0.0.1:9090";
              } else {
                  // 待機中の場合、順位を更新
                  const rankResponse = await fetch('/api/v1/queue/rank?' + queryParam);
                  const rankData = await rankResponse.json();
                  document.querySelector('#number').innerHTML = rankData.rank >= 0 ? rankData.rank : 'N/A';
                  document.querySelector('#updated').innerHTML = new Date().toLocaleString();
              }
          } catch (error) {
              console.error("Polling error:", error);
              clearInterval(intervalId);
          }
      }, 3000); // 3秒ごとにチェック
  }
</script>

待機キューに追加されたユーザーは、アクセスした時間順(スコア)でソートされます。クライアントは定期的にRedisのZRANKコマンドを使い、自身の待機順位を取得して画面に表示します。

        /**
     * ユーザーの現在の待機順位を取得します。
     *
     * @param queue  キュー名
     * @param userId ユーザーID
     * @return 待機順位。待機列にいない場合は-1を返す。
     */
    public Mono<Long> getRank(final String queue, final Long userId) {
        return reactiveRedisTemplate.opsForZSet().rank(USER_QUEUE_WAIT_KEY.formatted(queue), userId.toString())
                .defaultIfEmpty(-1L)
                .map(rank -> rank >= 0 ? rank + 1 : rank);
    }

この実装のバックグラウンドでは、スケジューラが定期的に動作します。スケジューラは、待機キューの先頭から一定数のユーザーを「アクセス許可キュー(proceedキュー)」に移動させます。つまり、時間が経過するにつれて、待機していたユーザーが順次サービスを利用できるようになる仕組みです。

以下のコードは、待機中のユーザーをアクセス許可キューに移動させるallowUserメソッドのロジックです。

    /**
     * 定期的に実行されるスケジューラ。
     * 全ての待機キューをスキャンし、一定数のユーザーをアクセス許可状態に移行させます。
     * Redisへの負荷を考慮し、KEYSではなくSCANコマンドを使用しています。
     */
    @Scheduled(initialDelay = 5000, fixedDelay = 10000)
    public void scheduleAllowUser() {
        if (!scheduling) {
            log.info("passed scheduling...");
            return;
        }
        log.info("called scheduling...");

        var maxAllowUserCount = 100L;
        reactiveRedisTemplate.scan(ScanOptions.scanOptions()
                        .match(USER_QUEUE_WAIT_KEY_FOR_SCAN)
                        .count(100)
                        .build())
                .map(key -> key.split(":")[2])
                .flatMap(queue -> allowUser(queue, maxAllowUserCount).map(allowed -> Tuples.of(queue, allowed)))
                .doOnNext(tuple -> log.info("Tried %d and allowed %d members of %s queue".formatted(maxAllowUserCount, tuple.getT2(), tuple.getT1())))
                .subscribe();
    }

待機列にいる指定されたユーザーを処理完了されたキューに移動させます。

        /**
     * 指定されたキューの待機列から、指定された数のユーザーをアクセス許可状態にします。
     *
     * @param queue 処理対象のキュー名
     * @param count 許可するユーザーの数
     * @return 実際に許可されたユーザーの数
     */
    public Mono<Long> allowUser(final String queue, final Long count) {
        return reactiveRedisTemplate.opsForZSet().popMin(USER_QUEUE_WAIT_KEY.formatted(queue), count)
                .flatMap(member -> reactiveRedisTemplate.opsForZSet().add(USER_QUEUE_PROCEED_KEY.formatted(queue), member.getValue(), Instant.now().getEpochSecond()))
                .count();
    }

まとめ

このシステムでは、まずクライアントがサーバーに対して定期的にサービスが利用可能かどうかを問い合わせ(ポーリングし)ます。もし待機キューに入った場合は、許可が下りるまで現在のページで待機し続けます。

一方サーバーは、自身の可用性をチェックし、処理可能な分だけ待機キューからユーザーを順次取り出し、サービス利用を許可します。

この一連の流れを繰り返すことで、システムは安定的に運営されます。

Discussion