Eventy backendu gry znikają pod obciążeniem? Oto architektura trwałego strumieniowania, która to naprawia
W skrócie
Poznaj architekturę trwałego strumieniowania zdarzeń, która eliminuje cichą utratę danych w backendzie gry pod obciążeniem i chroni Twoje systemy.
Twój potok analityczny właśnie stracił 40% zdarzeń śmierci graczy podczas weekendowego szczytu. Tablice wyników są nieaktualne. Raporty awarii nigdy nie dotarły. Przez trzy dni nikt tego nie zauważył.
To cichy zabójca niezawodności backendu gier: sprzężone architektury producent-konsument, które działają dobrze przy 100 żądaniach na sekundę, a przy 10 000 zaczynają tracić dane. Wywołanie RPC z serwera gry do konsumenta analitycznego przekracza limit czasu, połączenie jest resetowane, a zdarzenie po prostu znika — bez błędu, bez ponowienia, bez zapisu.
Ten runbook opisuje, co się psuje, jak to wykryć, jaki wzorzec architektoniczny to naprawia i jak zapobiec nawrotom.
Co się psuje: tryb awarii sprzężonego producenta-konsumenta
Tradycyjne architektury RPC wymuszają na producentach i konsumentach dopasowanie zarówno pod względem skali, jak i czasu. Gdy serwer gry wysyła zdarzenie „player_killed” bezpośrednio do usługi analitycznej przez HTTP:
- Niedopasowanie skali: Jeśli 5 000 graczy umiera jednocześnie podczas wydarzenia serwerowego, Twój endpoint analityczny otrzymuje nagły wzrost, którego nie jest w stanie przetworzyć. Połączenie HTTP wygasa po 30 sekundach. Zdarzenia są gubione.
- Sprzężenie czasowe: Jeśli usługa wykrywania oszustw wdraża nową wersję i jest offline przez 90 sekund, każde zdarzenie wyprodukowane w tym oknie znika.
- Fan-out do wielu konsumentów: To samo zdarzenie „player_killed” musi trafić do trzech niezależnych systemów — aktualizatora tablicy wyników, potoku analitycznego i pulpitu live-ops. Każdy konsument ma inną przepustowość. Najwolniejszy staje się wąskim gardłem dla producenta.
Oto jak to wygląda w praktyce:
[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
Efekt: cicha, częściowa utrata danych, która psuje analitykę, nieaktualne tablice wyników i niewidoczne wzorce awarii. Odkrywasz to tygodnie później, gdy liczby w lejku się nie zgadzają.
Konkretne liczby awarii
W typowym backendzie multiplayer dla indie obsługującym 5 000 równoczesnych graczy:
- ~2400 zdarzeń rozgrywki/sekundę w szczycie (zabójstwa, aktualizacje wyników, zmiany ekwipunku, przejścia między strefami)
- Średni rozmiar zdarzenia: ~200 bajtów
- Bezpośredni fan-out HTTP do 3 konsumentów: 7200 żądań wychodzących/sekundę
- Próg timeoutu konsumenta: 30 sekund
- Zaobserwowany współczynnik gubienia w szczycie: 15–45% w zależności od zdrowia konsumenta
Matematyka jest bezlitosna. Jeden konsument, który pozostaje niezdrowy przez 60 sekund, gubi 72 000 zdarzeń. Te zdarzenia przepadają, chyba że zbudowałeś bufor.
Jak wykryć cichą utratę zdarzeń
Cicha utrata zdarzeń jest z definicji trudna do wykrycia. Oto runbook wykrywania:
Krok 1: Instrumentacja numerów sekwencyjnych
Każdy producent powinien oznaczać każde zdarzenie monotonicznie rosnącym numerem sekwencyjnym dla danego źródła. Jeśli serwer gry wysyła zdarzenia dla player_abc, sekwencja wygląda tak: 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!
Luka w numerach sekwencyjnych po stronie konsumenta oznacza potwierdzoną utratę zdarzeń.
Krok 2: Monitoruj lag konsumenta
Śledź różnicę między najnowszą wyprodukowaną sekwencją a najnowszą skonsumowaną sekwencją dla każdego konsumenta. Progi alertów:
- Lag < 1000 zdarzeń: Zdrowy
- Lag 1000–10 000 zdarzeń: Ostrzeżenie — konsument nie nadąża
- Lag > 10 000 zdarzeń: Krytyczny — konsument jest praktycznie offline lub przeciążony
Krok 3: Krzyżowa walidacja sum
Porównuj liczby zdarzeń między logami producenta a licznikami przyjęcia u konsumenta co godzinę. Rozbieżność większa niż 1% wymaga zbadania.
// 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);
}
}
Jeśli liczba wyprodukowana na godzinę wynosi 8 400 000, a konsument analityczny przyjął 5 100 000, straciłeś 39% zdarzeń. To jest Twój sygnał.
Rozwiązanie architektoniczne: dekoupling przez trwały log zdarzeń
Rozwiązaniem jest odseparowanie producentów od konsumentów poprzez wstawienie trwałego bufora między nimi. Zamiast:
Game Server --direct HTTP--> Analytics
Game Server --direct HTTP--> Leaderboard Service
Game Server --direct HTTP--> Crash Reporter
Zapisujesz do:
Game Server --single write--> [Durable Event Stream] --independent reads--> Analytics
--independent reads--> Leaderboard Service
--independent reads--> Crash Reporter
Trwały strumień przyjmuje zapisy z prędkością producenta. Każdy konsument czyta we własnym tempie. Jeśli konsument jest offline przez 5 minut, zdarzenia gromadzą się w strumieniu, a konsument po powrocie wznawia od miejsca, w którym skończył. Brak utraty danych.
Kluczowe właściwości trwałego strumienia zdarzeń
Strumień zdarzeń klasy produkcyjnej dla backendów gier potrzebuje następujących właściwości:
| Property | Why It Matters for Games |
|---|---|
| Uporządkowanie w obrębie partycji | Wszystkie zdarzenia dla jednego gracza muszą być przetwarzane sekwencyjnie — zabójstwo przed aktualizacją wyniku, nie po |
| Trwałe przechowywanie | Zdarzenia przetrwają restarty konsumentów, okna wdrożeń i awarie infrastruktury |
| Niezależne offsety konsumentów | Konsument analityczny i konsument tablicy wyników czytają z różnymi prędkościami, nie blokując się nawzajem |
| Porządek na poziomie partycji | Zyskujesz równoległość między graczami (różne partycje) i spójność w obrębie gracza (ta sama partycja) |
Strategia partycjonowania dla backendów gier
Klucz partycji określa, które zdarzenia trafiają do którego uporządkowanego logu. Dla backendów gier naturalnym kluczem partycji jest 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]
To gwarantuje, że wszystkie zdarzenia dla pojedynczego gracza są przetwarzane w dokładnej kolejności przez każdego konsumenta, a zdarzenia różnych graczy mogą być przetwarzane równolegle.
Implementacja trwałego strumieniowania zdarzeń: praktyczny przewodnik
Oto konkretny wzorzec implementacji. Możesz go zbudować na bazie przechowywania obiektów (np. S3/R2), zarządzanej usługi strumieniowania lub logu opartego na bazie danych.
Producent: fan-out przy pojedynczym zapisie
Twój serwer gry zapisuje jedno zdarzenie. Infrastruktura strumieniowania obsługuje dostarczenie do wszystkich konsumentów.
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);
}
}
}
Kluczowe decyzje projektowe w tym kodzie:
- Klucz partycji to
playerId: Zapewnia kolejność dla każdego gracza we wszystkich typach zdarzeń. - Pojedyncze wywołanie zapisu: Producent nie wie i nie obchodzi go, ilu jest konsumentów. Dodanie nowego konsumenta (np. trackera wydarzeń sezonowych) nie wymaga żadnych zmian po stronie producenta.
- Lokalna kolejka ponowień: Jeśli strumień jest chwilowo niedostępny, zdarzenia trafiają do lokalnej kolejki i są wysyłane po przywróceniu łączności. To zapobiega utracie danych podczas chwilowych problemów z siecią.
Konsument: niezależne śledzenie offsetu
Każdy konsument utrzymuje własny offset odczytu dla każdej partycji. To podstawowy mechanizm umożliwiający niezależne prędkości przetwarzania.
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
}
}
}
Krytyczne szczegóły implementacji:
- Offset przesuwa się po przetworzeniu, nie przed. Jeśli konsument ulegnie awarii w trakcie przetwarzania partii, po restarcie przetworzy te same zdarzenia ponownie. Twoje handlery muszą być idempotentne — przetworzenie tego samego zdarzenia zabójstwa dwa razy nie może podwoić wpisu w tablicy wyników.
- Rozmiar partii 500 równoważy przepustowość i pamięć. Przy zdarzeniach o rozmiarze 200 bajtów to ~100 KB na partię — pomijalne.
- Interwał odpytywania 100 ms oznacza, że najgorszy przypadek opóźnienia to ~100 ms od wyprodukowania zdarzenia do aktualizacji tablicy wyników. Dla większości przypadków użycia tablic wyników to w pełni akceptowalne. Jeśli potrzebujesz opóźnienia poniżej 10 ms, wkraczasz na teren transportu czasu rzeczywistego, co jest zupełnie inną architekturą.
Idempotentność: siatka bezpieczeństwa konsumenta
Idempotentność jest w tej architekturze niepodlegająca negocjacjom. Oto konkretny idempotentny handler:
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 }
);
});
}
}
Tabela processed_events działa jak magazyn deduplikacji. Kosztuje jeden dodatkowy zapis na zdarzenie, ale gwarantuje, że restarty konsumentów nigdy nie uszkodzą Twoich danych.
Radzenie sobie z przestojem konsumenta: gwarancja buforowania
Głównym powodem używania trwałego strumieniowania zdarzeń jest przetrwanie przestoju konsumenta bez utraty danych. Oto jak działa matematyka buforowania:
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
Przy 144 MB to bez problemu mieści się w każdym nowoczesnym systemie przechowywania. Kluczowy wniosek: rozmiar bufora jest proporcjonalny do szybkości zdarzeń pomnożonej przez maksymalną tolerancję przestoju, a nie do całkowitej historycznej ilości danych.
W przypadku długoterminowego przechowywania (30 dni zdarzeń do odtworzenia lub ponownego przetworzenia) liczby rosną:
30 days × 86,400 seconds × 2,400 events/sec × 200 bytes = ~1.24 TB
To w pełni mieści się w możliwościach backendów przechowywania obiektów. Struktura partycji utrzymuje przewidywalną wydajność odczytu nawet w tej skali — nigdy nie skanujesz całego logu, czytasz z konkretnych partycji na konkretnych offsetach.
Co horizOn oferuje w architekturze opartej na zdarzeniach
Gdy Twój strumień zdarzeń dostarcza dane niezawodnie, potrzebujesz usług, które te zdarzenia konsumują. horizOn dostarcza prymitywy backendowe, które wpinają się w stronę konsumencką tej architektury:
- Tablice wyników konsumują zdarzenia zabójstw, wyników i ukończeń, aby aktualizować rankingi w czasie rzeczywistym
- Logi użytkowników przechwytują strumienie zdarzeń do debugowania problemów zgłaszanych przez graczy — gdy gracz mówi „mój wynik się zresetował”, możesz przejrzeć historię jego zdarzeń
- Raporty awarii przyjmują zdarzenia awarii z pełnym kontekstem, dając Ci stack trace powiązany z sesją gracza
Sedno: strumień zdarzeń dostarcza dane do tych usług niezawodnie. Same usługi muszą być sprawdzone w boju. Możesz przeczytać o tym, jak zaprojektowaliśmy jedną z naszych największych aktualizacji backendu, w tym opisie aktualizacji backendu indie horizOn, który omawia decyzje infrastrukturalne stojące za niezawodnym przyjmowaniem zdarzeń na dużą skalę.
Runbook: zapobieganie nawrotom utraty zdarzeń
Jeśli raz doświadczyłeś utraty zdarzeń, oto lista kontrolna, która zapobiegnie jej ponownemu wystąpieniu:
1. Przeaudytuj każde sprzężenie producent-konsument
Przejdź przez swój kod i zidentyfikuj każde miejsce, w którym serwer gry wykonuje bezpośrednie synchroniczne wywołanie HTTP do usługi backendowej. Każde z nich to potencjalny punkt gubienia danych pod obciążeniem. Wypisz je:
GameServer → AnalyticsService (HTTP POST, no retry) ← RISK
GameServer → LeaderboardService (HTTP POST, no retry) ← RISK
GameServer → CrashReporter (UDP, fire-and-forget) ← RISK
2. Wprowadź strumień zdarzeń jako pośrednika
Zastąp każde bezpośrednie wywołanie pojedynczym zapisem do trwałego strumienia zdarzeń. Każda usługa niższego poziomu staje się niezależnym konsumentem z własnym offsetem.
3. Zaimplementuj monitorowanie zdrowia konsumentów
Dla każdego konsumenta śledź:
- Lag (zdarzenia w tyle za producentem)
- Tempo przetwarzania (zdarzenia konsumowane na sekundę)
- Współczynnik błędów (zdarzenia, których przetwarzanie nie powiodło się na sekundę)
- Ostatni poprawny offset (wykrywanie nieaktualności)
Alertuj, gdy lag przekroczy wyliczone okno bufora. Jeśli Twój strumień przechowuje dane przez 7 dni, a konsument był niedostępny przez 6 dni, masz 24 godziny, zanim zacznie się utrata danych.
4. Przetestuj odzyskiwanie po restarcie konsumenta
Celowo wyłącz konsumenta na 5 minut, przywróć go i zweryfikuj, że nadrabia zaległości bez duplikatów. To Twój test pewności, że architektura działa. Zautomatyzuj go w 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. Skonfiguruj dead-letter queue
Zdarzenia, których przetwarzanie nie powiedzie się po N ponowieniach (zwykle 3–5), trafiają do dead-letter queue. Monitoruj rozmiar DLQ. Rosnący DLQ oznacza błąd w konsumencie, a nie przejściową awarię.
Najlepsze praktyki dla strumieniowania zdarzeń w backendzie gry
Wybierz playerId jako klucz partycji. To daje Ci porządek per gracz (kluczowy dla zdarzeń ekwipunku, wyniku i stanu) przy jednoczesnym pozwoleniu na równoległość między graczami. Nie partycjonuj według typu zdarzenia — „kill” i „score_update” dla tego samego gracza muszą pozostać w kolejności.
Utrzymuj zdarzenia małe i samoopisujące się. Każde zdarzenie powinno mieć 100–500 bajtów. Uwzględnij typ zdarzenia, ID gracza, znacznik czasu i minimalny wymagany payload. Nie osadzaj pełnego stanu gry — odwołuj się do niego po ID.
Projektuj konsumentów jako idempotentnych od pierwszego dnia. Używaj ID zdarzeń i magazynu deduplikacji. Zakładaj, że każde zdarzenie zostanie dostarczone co najmniej raz, a podczas failover prawdopodobnie więcej niż raz.
Dopasuj okno retencji do maksymalnego akceptowalnego przestoju konsumenta. Jeśli najdłuższe wdrożenie trwa 15 minut, przechowuj co najmniej 30 minut gorących zdarzeń. Trzymaj 7–30 dni zimnego przechowywania do odtwarzania i debugowania.
Monitoruj lag konsumenta jako metrykę pierwszej klasy. Lag to serce Twojej architektury opartej na zdarzeniach. Konsument, który pozostaje w tyle o więcej niż 50% swojego okna retencji, to sytuacja awaryjna związana z utratą danych, a nie zadanie „naprawimy to w następnym sprincie”.
Następne kroki
Jeśli obecnie używasz bezpośrednich wywołań RPC z serwerów gry do usług backendowych, przeaudytuj te połączenia w tym tygodniu. Policz, ile z nich po cichu gubiłoby zdarzenia przy 10-krotnym skoku obciążenia. Następnie zaprojektuj prototyp trwałego strumienia zdarzeń między producentami a konsumentami — nawet prosty log oparty na bazie danych jest lepszy niż bezpośrednie sprzężenie.
Po stronie konsumenckiej Twojej architektury opartej na zdarzeniach — tablice wyników, raportowanie awarii, logi użytkowników i konfiguracja zdalna — horizOn dostarcza te elementy jako usługi zarządzane, dzięki czemu możesz skupić się na logice gry, zamiast wymyślać każdego konsumenta od zera. Sprawdź dokumentację API, aby zobaczyć, które prymitywy pasują do Twojego backendu.
Źródło: Ogłoszenie Cloudflare K2: strumienie zdarzeń serverless