待機キューの実装による高トラフィック制御
大量のユーザーが特定のサービスに一斉にアクセスすると、システムのリソースが不足し、障害発生に繋がります。限られたリソースでサービスを安定して維持するためには、同時にアクセスできるユーザー数を制限する技術が必要です。
例えば、全体のユーザー数が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