負荷時にゲームバックエンドのイベントが消失? それを解決するDurable Streamingアーキテクチャ
要点まとめ
ゲームバックエンドで負荷時に発生するイベント消失を防ぐDurable Streamingアーキテクチャを解説。プロデューサーとコンシューマーの結合を断ち切り、オフセット管理・冪等性・監視・復旧手順まで実践的に学べる運用ガイドです。直接RPCの危険性と具体的な実装パターンも網羅しています。
アナリティクスパイプラインは、週末のピーク時にプレイヤーの死亡イベントの40%を失いました。リーダーボードは古いまま。クラッシュレポートは届かず、3日間誰も気づきませんでした。
これはゲームバックエンドの信頼性を静かに蝕む問題です: 結合されたプロデューサー・コンシューマーアーキテクチャは、毎秒100リクエストでは正常に動作しても、毎秒10,000リクエストではデータを大量に失います。ゲームサーバーからアナリティクスコンシューマーへのRPC呼び出しがタイムアウトし、接続がリセットされ、イベントは跡形もなく消えます — エラーも、リトライも、記録もありません。
このランクブックでは、何が壊れるのか、その検出方法、それを修正するアーキテクチャパターン、そして再発防止策を解説します。
何が壊れるのか: 結合されたプロデューサー・コンシューマーの障害モード
従来のRPCアーキテクチャでは、プロデューサーとコンシューマーがスケールと時間の両方で一致することを強制されます。ゲームサーバーが「player_killed」イベントをHTTPでアナリティクスサービスに直接送信する場合を考えます:
- スケールの不一致: サーバーイベント中に5,000人のプレイヤーが同時に死亡すると、アナリティクスエンドポイントは処理できないバーストを受け取ります。HTTP接続は30秒でタイムアウトし、イベントは失われます。
- 時間的結合: 不正検出サービスが新バージョンをデプロイして90秒間オフラインになると、その間に生成されたすべてのイベントが消えます。
- マルチコンシューマーファンアウト: 同じ「player_killed」イベントは、リーダーボード更新、アナリティクスパイプライン、ライブオペレーションダッシュボードの3つの独立したシステムに届く必要があります。各コンシューマーのスループットは異なり、最も遅いコンシューマーがプロデューサーのボトルネックになります。
実際の様子は次のとおりです:
[Game Server] --HTTP POST--> [Analytics Service] ✓ works at 200 req/s
[Game Server] --HTTP POST--> [Analytics Service] ✗ 40% drops at 8,000 req/s
[Game Server] --HTTP POST--> [Fraud Detection] ✗ offline during deploy
結果は: 静かな部分的なデータ損失です。アナリティクスを破壊し、リーダーボードを古くし、クラッシュパターンを不可視にします。数週間後、ファネル数値が合わないときに気づくのです。
具体的な障害数値
典型的なインディーマルチプレイヤーバックエンドで5,000人の同時接続プレイヤーを処理する場合:
- ピーク時のゲームプレイイベントは毎秒約2,400件(キル、スコア更新、インベントリ変更、ゾーン遷移)
- 平均イベントペイロード: 約200バイト
- 3つのコンシューマーへの直接HTTPファンアウト: 毎秒7,200件の送信リクエスト
- コンシューマーのタイムアウトしきい値: 30秒
- ピーク時に観測されたドロップ率: コンシューマーの健全性に応じて15〜45%
その計算は厳しいものです。1つのコンシューマーが60秒間不健全になると、72,000イベントが失われます。バッファを構築していない限り、それらのイベントは消えます。
サイレントなイベント消失を検出する方法
サイレントなイベント消失は、定義上、検出が困難です。以下が検出のランクブックです:
ステップ1: シーケンス番号を仕込む
すべてのプロデューサーは、ソースエンティティごとに単調増加するシーケンス番号を各イベントに刻印する必要があります。ゲームサーバーが player_abc のイベントを送信する場合、シーケンスは1、2、3、4...となります。
Event 1: { seq: 1, player: "abc", type: "kill", ts: 1719432000 }
Event 2: { seq: 2, player: "abc", type: "death", ts: 1719432001 }
Event 3: { seq: 4, player: "abc", type: "score", ts: 1719432005 } // seq 3 missing!
コンシューマー側のシーケンス番号のギャップは、イベント消失が確定したことを意味します。
ステップ2: コンシューマーラグを監視する
各コンシューマーについて、最新の生成済みシーケンスと最新の消費済みシーケンスの差を追跡します。アラートしきい値:
- Lag < 1,000 events: 健全
- Lag 1,000–10,000 events: 警告 — コンシューマーが遅れています
- Lag > 10,000 events: 重大 — コンシューマーは事実上オフラインか過負荷です
ステップ3: 合計値を相互検証する
プロデューサーログとコンシューマーの取り込み件数を1時間ごとに比較します。1%を超える差異は調査が必要です。
// Producer-side counter (emit to your monitoring system every 60s)
public class EventProducerMetrics
{
private long _producedCount = 0;
public void RecordProduced()
{
Interlocked.Increment(ref _producedCount);
}
public long GetAndResetCount()
{
return Interlocked.Exchange(ref _producedCount, 0);
}
}
生成された1時間あたりの件数が8,400,000で、アナリティクスコンシューマーが5,100,000を取り込んだ場合、39%のイベントを失ったことになります。それがシグナルです。
アーキテクチャ修正: Durable Event Logによるデカップリング
解決策は、プロデューサーとコンシューマーの間に耐久性のあるバッファを挿入してデカップリングすることです。次の代わりに:
Game Server --direct HTTP--> Analytics
Game Server --direct HTTP--> Leaderboard Service
Game Server --direct HTTP--> Crash Reporter
次のように書き込みます:
Game Server --single write--> [Durable Event Stream] --independent reads--> Analytics
--independent reads--> Leaderboard Service
--independent reads--> Crash Reporter
耐久性のあるストリームは、プロデューサーの速度で書き込みを吸収します。各コンシューマーは自分のペースで読み取ります。コンシューマーが5分間オフラインになっても、イベントはストリームに蓄積され、復帰時に中断した場所から再開します。データ損失はありません。
Durable Event Streamのコア特性
本番環境対応のゲームバックエンド向けイベントストリームには、次の特性が必要です:
| プロパティ | ゲームにとって重要な理由 |
|---|---|
| パーティション内で順序付け | 1人のプレイヤーのすべてのイベントは順番に処理される必要があります — スコア更新の前にキル、その逆ではない |
| 耐久性のあるストレージ | イベントはコンシューマーの再起動、デプロイウィンドウ、インフラ障害を生き延びます |
| 独立したコンシューマーオフセット | アナリティクスコンシューマーとリーダーボードコンシューマーは、互いをブロックせずに異なる速度で読み取ります |
| パーティションレベルの順序付け | プレイヤー間の並列処理(異なるパーティション)とプレイヤー内の一貫性(同じパーティション)の両方を実現します |
ゲームバックエンドのパーティショニング戦略
パーティションキーによって、どのイベントがどの順序付きログに配置されるかが決まります。ゲームバックエンドでは、自然なパーティションキーは playerId です:
Partition 0: [player_abc kill#1] [player_abc death#2] [player_abc score#3]
Partition 1: [player_def zone#1] [player_def kill#2] [player_def loot#3]
Partition 2: [player_ghi death#1] [player_ghi spawn#2]
これにより、単一プレイヤーのすべてのイベントがすべてのコンシューマーによって正確な順序で処理されることが保証され、異なるプレイヤーのイベントは並列に処理できます。
Durable Event Streamingの実装: 実践チュートリアル
以下は具体的な実装パターンです。オブジェクトストレージ(S3/R2など)、マネージドストリーミングサービス、データベースバックアップのログの上に構築できます。
プロデューサー: シングルライト・ファンアウト
ゲームサーバーは1つのイベントを書き込みます。ストリーミングインフラがすべてのコンシューマーへの配信を処理します。
public class GameEventProducer
{
private readonly IEventStreamClient _stream;
private readonly ILogger _logger;
public GameEventProducer(IEventStreamClient stream, ILogger logger)
{
_stream = stream;
_logger = logger;
}
public async Task EmitAsync(string playerId, string eventType,
Dictionary<string, object> payload)
{
var gameEvent = new GameEvent
{
Id = Guid.NewGuid().ToString(),
PlayerId = playerId, // partition key
EventType = eventType,
Payload = payload,
Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()
};
try
{
// Single write — the stream handles fan-out to all consumers
await _stream.AppendAsync(
partitionKey: playerId,
eventData: JsonSerializer.Serialize(gameEvent)
);
}
catch (EventStreamException ex)
{
// Events that fail to write should be queued locally
// and retried, not silently dropped
_logger.LogWarning(ex,
"Failed to emit event {EventType} for {Player}, queueing retry",
eventType, playerId);
await _retryQueue.EnqueueAsync(gameEvent);
}
}
}
このコードの主要な設計判断:
- パーティションキーは
playerId: すべてのイベントタイプにわたってプレイヤーごとの順序を保証します。 - 単一の書き込み呼び出し: プロデューサーはコンシューマーが何個あるかを知る必要も、気にする必要もありません。新しいコンシューマー(シーズンイベントトラッカーなど)を追加しても、プロデューサーの変更はゼロです。
- ローカルリトライキュー: ストリームが一時的に利用できない場合、イベントはローカルでキューに入り、接続が回復するとフラッシュされます。ネットワークの瞬断時のデータ損失を防ぎます。
コンシューマー: 独立オフセット追跡
各コンシューマーは、パーティションごとに独自の読み取りオフセットを維持します。これが独立した処理速度を可能にする中核メカニズムです。
public class LeaderboardConsumer
{
private readonly IEventStreamClient _stream;
private readonly ILeaderboardService _leaderboard;
private readonly IOffsetStore _offsets;
public async Task ProcessEventsAsync(CancellationToken ct)
{
while (!ct.IsCancellationRequested)
{
// Read the next batch from where we left off
var lastOffset = await _offsets.GetOffsetAsync(
consumerName: "leaderboard-updater",
partitionId: 0
);
var batch = await _stream.ReadBatchAsync(
partitionId: 0,
fromOffset: lastOffset,
maxBatchSize: 500
);
foreach (var evt in batch.Events)
{
var gameEvent = JsonSerializer.Deserialize<GameEvent>(evt.Data);
if (gameEvent.EventType == "player_killed")
{
await _leaderboard.IncrementKillsAsync(
gameEvent.PlayerId,
amount: 1
);
}
else if (gameEvent.EventType == "score_updated")
{
await _leaderboard.UpdateScoreAsync(
gameEvent.PlayerId,
gameEvent.Payload["score"].GetInt32()
);
}
// Advance offset AFTER successful processing
await _offsets.SetOffsetAsync(
consumerName: "leaderboard-updater",
partitionId: 0,
offset: evt.Offset + 1
);
}
await Task.Delay(100, ct); // Poll interval
}
}
}
重要な実装詳細:
- オフセットは処理後ではなく、処理が成功した後に進みます。 コンシューマーがバッチの途中でクラッシュした場合、再起動時に同じイベントを再処理します。ハンドラーは冪等でなければなりません — 同じキルイベントを2回処理しても、リーダーボードエントリが二重にカウントされてはいけません。
- バッチサイズ500は、スループットとメモリのバランスを取ります。200バイトのイベントなら、バッチあたり約100KB — 無視できる量です。
- 100msのポーリング間隔は、イベント生成からリーダーボード更新までの最悪ケースのレイテンシが約100msであることを意味します。ほとんどのリーダーボードのユースケースでは、これは完全に許容範囲です。サブ10msのレイテンシが必要なら、それはリアルタイムトランスポートの領域であり、まったく別のアーキテクチャです。
冪等性: コンシューマーのセーフティネット
このアーキテクチャでは冪等性は必須です。以下は具体的な冪等ハンドラーです:
public class IdempotentKillCounter
{
private readonly IDatabase _db;
public async Task ProcessKillAsync(string playerId, string eventId)
{
// Check if we already processed this event
var alreadyProcessed = await _db.ExecuteScalarAsync<bool>(
"SELECT COUNT(*) > 0 FROM processed_events WHERE event_id = @id",
new { id = eventId }
);
if (alreadyProcessed)
{
return; // Skip duplicate — this is the idempotency guard
}
// Process and record in a transaction
await _db.ExecuteInTransactionAsync(async tx =>
{
await tx.ExecuteAsync(
"UPDATE leaderboard SET kills = kills + 1 WHERE player_id = @pid",
new { pid = playerId }
);
await tx.ExecuteAsync(
"INSERT INTO processed_events (event_id, processed_at) VALUES (@id, @now)",
new { id = eventId, now = DateTime.UtcNow }
);
});
}
}
processed_events テーブルは重複排除ストアとして機能します。イベントごとに1回の追加書き込みがコストですが、コンシューマーの再起動がデータを壊さないことを保証します。
コンシューマーダウンタイムへの対応: バッファリングの保証
耐久性のあるイベントストリーミングを使う主な理由は、データ損失なしにコンシューマーのダウンタイムを乗り切ることです。バッファリングの計算は次のようになります:
Producer rate: 2,400 events/second
Consumer offline window: 5 minutes (300 seconds)
Events buffered: 720,000 events
Storage required: 720,000 × 200 bytes = ~144 MB
144MBなら、最新のストレージシステムに簡単に収まります。重要な洞察: バッファサイズは、イベントレート×最大許容ダウンタイムに比例し、履歴データ全体には比例しません。
長期保存(リプレイや再処理のための30日間)の場合、数値は大きくなります:
30 days × 86,400 seconds × 2,400 events/sec × 200 bytes = ~1.24 TB
これはオブジェクトストレージバックエンドの容量に十分収まります。パーティション構造により、この規模でも読み取りパフォーマンスは予測可能です — ログ全体をスキャンすることはなく、特定のパーティションの特定のオフセットから読み取ります。
horizOnがイベント駆動アーキテクチャで提供するもの
イベントストリームが確実にデータを配信するようになったら、それらのイベントを消費するサービスが必要です。horizOnは、このアーキテクチャのコンシューマー側に接続するバックエンドプリミティブを提供します:
- リーダーボードはキル、スコア、完了イベントを消費してランキングをリアルタイムに更新します
- ユーザーログはイベントストリームをキャプチャして、プレイヤーが報告した問題のデバッグに役立てます — プレイヤーが「スコアがリセットされた」と言ったら、イベント履歴を照会できます
- クラッシュレポートは完全なコンテキストを持つクラッシュイベントを取り込み、プレイヤーセッションに結び付けられたスタックトレースを提供します
ポイント: イベントストリームはデータをこれらのサービスに確実に届けます。サービス自体は実戦テスト済みである必要があります。私たちが最大のバックエンドアップデートの1つをどうアーキテクチャしたかは、このhorizOnのインディーゲームバックエンドアップデートの解説で読めます。スケールでの信頼性の高いイベント取り込みの背後にあるインフラ決定をカバーしています。
ランクブック: イベント消失の再発防止
一度イベント消失を経験したら、再発を防ぐためのチェックリストは次のとおりです:
1. すべてのプロデューサー・コンシューマー結合を監査する
コードベースを調べて、ゲームサーバーがバックエンドサービスに対して直接同期HTTP呼び出しを行っているすべての箇所を特定します。それぞれが負荷時のドロップポイントになる可能性があります。リストアップします:
GameServer → AnalyticsService (HTTP POST, no retry) ← RISK
GameServer → LeaderboardService (HTTP POST, no retry) ← RISK
GameServer → CrashReporter (UDP, fire-and-forget) ← RISK
2. イベントストリームを中間層として導入する
各直接呼び出しを、耐久性のあるイベントストリームへの単一書き込みに置き換えます。各ダウンストリームサービスは、独自のオフセットを持つ独立したコンシューマーになります。
3. コンシューマーのヘルスモニタリングを実装する
各コンシューマーについて、以下を追跡します:
- Lag(プロデューサーから遅れているイベント数)
- 処理レート(毎秒消費されるイベント数)
- エラーレート(毎秒処理に失敗したイベント数)
- 最後に成功したオフセット(陳腐化の検出)
計算済みのバッファウィンドウを超えるラグでアラートを出します。ストリームが7日間保持し、コンシューマーが6日間ダウンしていた場合、データ損失が始まるまで24時間あります。
4. コンシューマーの再起動リカバリをテストする
意図的にコンシューマーを5分間オフラインにし、復帰させて、重複なしに追いつくことを検証します。これがアーキテクチャが機能するという信頼テストです。CIで自動化します:
[Test]
public async Task ConsumerResumesAfterDowntime()
{
// Produce 10,000 events
await ProduceEvents(count: 10_000);
// Simulate consumer offline — skip reads for 30 seconds
await Task.Delay(TimeSpan.FromSeconds(30));
// Resume consumer
var processed = await Consumer.ProcessUntilCaughtUp();
// Verify: all events processed, no duplicates
Assert.AreEqual(10_000, processed.UniqueEventCount);
Assert.AreEqual(0, processed.DuplicateCount);
}
5. Dead-Letter Queueを設定する
N回(通常3〜5回)のリトライ後に処理に失敗したイベントは、Dead-Letter Queueに移動します。DLQサイズを監視します。DLQが増え続ける場合は、一時的な障害ではなくコンシューマーのバグを意味します。
ゲームバックエンドのイベントストリーミングのベストプラクティス
パーティションキーには
playerIdを選びます。 これにより、プレイヤーごとの順序(インベントリ、スコア、状態イベントに重要)を保証しつつ、プレイヤー間の並列処理を可能にします。イベントタイプでパーティション分割しないでください — 同じプレイヤーの「kill」と「score_update」は順序を保つ必要があります。イベントは小さく、自己記述的に保ちます。 各イベントは100〜500バイトにします。イベントタイプ、プレイヤーID、タイムスタンプ、必要な最小ペイロードを含めます。完全なゲーム状態を埋め込まず、IDで参照します。
コンシューマーは初日から冪等に設計します。 イベントIDと重複排除ストアを使用します。すべてのイベントが少なくとも1回、フェイルオーバー時には複数回配信されると想定します。
保持ウィンドウは、許容できる最大コンシューマーダウンタイムに合わせます。 最長のデプロイが15分なら、ホットイベントを少なくとも30分保持します。リプレイとデバッグ用に7〜30日のコールドストレージを保持します。
コンシューマーラグを第一級のメトリクスとして監視します。 ラグはイベント駆動アーキテクチャの鼓動です。保持ウィンドウの50%以上遅れるコンシューマーは、「次のスプリントで直す」項目ではなく、データ損失の緊急事態です。
次のステップ
現在、ゲームサーバーからバックエンドサービスへの直接RPC呼び出しを行っているなら、今週中にそれらの接続を監査してください。10倍の負荷スパイクでサイレントにイベントを落とす接続がいくつあるか数えてください。次に、プロデューサーとコンシューマーの間に耐久性のあるイベントストリームのプロトタイプを作成してください — 単純なデータベースバックアップのログでも、直接結合よりはましです。
イベント駆動アーキテクチャのコンシューマー側 — リーダーボード、クラッシュレポート、ユーザーログ、リモート設定 — には、horizOnがマネージドサービスとして提供しているので、各コンシューマーをゼロから再発明する代わりにゲームロジックに集中できます。APIドキュメントで、どのプリミティブがバックエンドに適合するか確認してください。