Sven Erik Matzen

Software Architect | Cloud & Security Expert | AI-enabled Solutions

Der lange Schwanz der Latenz: Tail Latency und warum Mittelwerte in der Cloud lügen

🎧 Listen to this article

Cloud Computing · 2026-09-27

EU-Kennzeichnung: vollständig KI-generierter Inhalt Vollständig KI-generierter Artikel (ohne Vorabprüfung).

Der Aufhänger: Jeder Teil ist schnell, das Ganze ist langsam

Stell dir vor, du betreibst einen Dienst, der auf 100 Maschinen verteilt ist. Eine Suchanfrage wird an alle 100 geschickt, jede Maschine durchsucht ihren Teil des Datenbestandes, und erst wenn die letzte Antwort eingetroffen ist, kann das Ergebnis an den Nutzer gehen. Du hast sauber gemessen: Jede einzelne Maschine antwortet in 99 von 100 Fällen innerhalb von 10 Millisekunden. Nur in einem von 100 Fällen dauert es eine ganze Sekunde – irgendein Hintergrundprozess, ein Garbage-Collection-Lauf, ein ungünstig getimter Festplattenzugriff. Ein Prozent Ausreißer. Das klingt nach einem System, das man guten Gewissens in Produktion geben kann.

Rechne nach. Die Wahrscheinlichkeit, dass alle 100 Maschinen schnell sind, beträgt 0,99 hoch 100 – etwa 0,366. Das heißt: In 63 Prozent aller Nutzeranfragen wartet mindestens eine Maschine eine Sekunde, und weil du auf die letzte Antwort warten musst, wartet der Nutzer diese Sekunde mit. Ein System, dessen Teile in 99 Prozent der Fälle blitzschnell sind, ist als Ganzes in fast zwei von drei Fällen träge.

Dieses Rechenbeispiel steht am Anfang eines der einflussreichsten Aufsätze der Cloud-Ära: „The Tail at Scale" von Jeffrey Dean und Luiz André Barroso, 2013 in den Communications of the ACM erschienen. Die beiden Google-Ingenieure formulierten darin eine Einsicht, die das Denken über verteilte Systeme verschoben hat: Sobald ein System aus vielen Komponenten besteht, die alle antworten müssen, bestimmt nicht die typische, sondern die seltene schlechte Antwortzeit das Nutzererlebnis. Und je größer das System wird, desto stärker wirkt dieser Effekt. Skalierung verbessert nicht nur den Durchsatz – sie verstärkt systematisch den Einfluss der Ausreißer.

Der Begriff dafür lautet Tail Latency – die Latenz am rechten Rand, im „Schwanz" der Verteilung. Und die vielleicht unangenehmste Konsequenz: Fast alle Dashboards, die Entwicklerinnen und Entwickler täglich betrachten, zeigen genau die Zahl, die hier nichts aussagt – den Mittelwert.


Teil 1: Was Latenz eigentlich ist – und warum der Mittelwert die falsche Frage beantwortet

Latenz ist eine Verteilung, keine Zahl

Der erste Denkfehler ist sprachlicher Natur. Wir sagen „die Latenz des Dienstes beträgt 40 Millisekunden", als wäre Latenz eine Eigenschaft wie Masse oder Länge. Sie ist es nicht. Latenz ist eine Zufallsvariable mit einer Verteilung, und diese Verteilung ist in praktisch jedem realen System rechtsschief und schwerschwänzig: Es gibt eine untere Schranke (schneller als die Summe der physikalisch notwendigen Schritte geht es nicht), aber keine obere. Nach unten ist der Weg kurz, nach oben offen.

Bei einer solchen Verteilung ist der Mittelwert eine der am wenigsten informativen Kennzahlen, die man bilden kann. Er vermischt zwei völlig unterschiedliche Populationen – den dichten Hauptteil der schnellen Antworten und die dünne, weit ausgezogene Fahne der langsamen – zu einer einzigen Zahl, die keine von beiden beschreibt. Ein Dienst mit einem Mittelwert von 40 ms kann bedeuten: „fast alle Anfragen brauchen 38 bis 42 ms" oder „95 Prozent brauchen 10 ms und 5 Prozent brauchen 600 ms". Das sind völlig verschiedene Systeme, die sich in der Produktion völlig verschieden verhalten – und sie sehen auf dem Mittelwert-Dashboard identisch aus.

Noch schlimmer: Der Mittelwert ist gegenüber dem, was die Nutzer erleben, systematisch zu optimistisch. Wenn 1 Prozent der Anfragen 100-mal so lange dauert wie der Rest, dann trägt dieses Prozent nur etwa die Hälfte zum Mittelwert bei – der Mittelwert verdoppelt sich also nur, während für die betroffenen Nutzer die Welt stehenbleibt. Das arithmetische Mittel ist gebaut, um Ausreißer zu glätten. Genau das darf man hier nicht wollen.

Perzentile: die richtige Frage, richtig gestellt

Die brauchbare Alternative sind Perzentile. Das 99. Perzentil (üblich abgekürzt p99) ist der Wert, unter dem 99 Prozent aller gemessenen Antwortzeiten liegen. „p99 = 250 ms" heißt: Eine von hundert Anfragen dauert länger als 250 Millisekunden. Das ist eine Aussage über Nutzer, nicht über Statistik.

Dabei ist ein Punkt wichtig, der in der Praxis oft untergeht: Bei welchem Perzentil man hinschaut, ist keine Geschmacksfrage, sondern ergibt sich aus der Anzahl der Interaktionen pro Nutzersitzung. Wer eine Web-Anwendung baut, in der eine einzige Seitenansicht 100 Backend-Aufrufe erzeugt, für den ist p99 keine seltene Ausnahme, sondern der Normalfall: Bei 100 Aufrufen erwischt praktisch jede Seitenansicht mindestens einen p99-Fall. Dean und Barroso formulieren die Konsequenz drastisch: In Systemen mit hohem Fan-out ist das 99,9. Perzentil (p999) die relevante Kennzahl, weil das 99. Perzentil eines Bausteins bereits zum Median des zusammengesetzten Dienstes werden kann.

Ein zweiter, subtilerer Punkt: Perzentile addieren sich nicht. Wenn Dienst A ein p99 von 50 ms hat und Dienst B ein p99 von 50 ms, dann hat die Kette A→B nicht ein p99 von 100 ms. Die Wahrscheinlichkeit, dass beide gleichzeitig ihren schlechten Fall zeigen, ist klein (bei Unabhängigkeit 0,0001), aber die Wahrscheinlichkeit, dass mindestens einer langsam ist, verdoppelt sich fast: 1 − 0,99² ≈ 1,99 Prozent. Das p99 der Kette liegt also dichter bei 50 ms als bei 100 ms, während sich das p98 verschlechtert. Perzentile sind keine Größen, mit denen man Budgets addieren kann – man muss über Wahrscheinlichkeiten rechnen.

Der Messfehler, der fast jede Benchmark verdirbt: Coordinated Omission

Bevor man über das Zähmen des Schwanzes nachdenkt, muss man ihn überhaupt korrekt messen können. Und hier liegt ein systematischer Fehler, den der Java-Performance-Ingenieur Gil Tene in seinem viel zitierten Vortrag How NOT to Measure Latency als Coordinated Omission bezeichnet hat: Die meisten Lastgeneratoren melden zu gute hohe Perzentile – nicht weil sie falsch rechnen, sondern weil sie die schlimmsten Fälle gar nicht erst in ihre Stichprobe aufnehmen.

Der Mechanismus ist einfach und deshalb so tückisch. Ein typischer Lastgenerator arbeitet in einer Schleife: Anfrage senden, auf Antwort warten, Zeit notieren, nächste Anfrage senden. Wenn das System nun für zwei Sekunden hängt, dann bleibt dieser Thread zwei Sekunden blockiert – und sendet in dieser Zeit keine weiteren Anfragen. Er hätte, bei einer Zielrate von 1.000 Anfragen pro Sekunde, in diesen zwei Sekunden 2.000 Anfragen abschicken sollen, die alle langsam gewesen wären. Stattdessen notiert er eine Messung von 2.000 ms. Die 1.999 schlechten Messungen, die es hätte geben müssen, fehlen. Der Lastgenerator hat mit dem getesteten System stillschweigend „kooperiert" und genau während der Aussetzer aufgehört zu messen.

Die Folge ist keine kleine Ungenauigkeit, sondern eine Verzerrung um Größenordnungen im Bereich, der einen interessiert. Ein Messaufbau, der p999 = 20 ms meldet, kann in Wahrheit p999 = 2.000 ms haben. Tene hat darum mit HdrHistogram ein Werkzeug bereitgestellt, das Latenzen über viele Größenordnungen mit konstanter relativer Genauigkeit protokolliert und Korrekturen für Coordinated Omission ermöglicht. Die praktische Regel, die sich daraus ergibt: Messe nach einem Zeitplan, nicht nach einer Schleife. Eine Anfrage, die um 10:00:00,000 hätte abgehen sollen und erst um 10:00:01,500 beantwortet wird, hat 1.500 ms Latenz – auch wenn der Server sie erst um 10:00:01,498 überhaupt gesehen hat. Das ist die Latenz, die der Nutzer erlebt.


Teil 2: Die Mathematik der Fan-out-Verstärkung

Das Grundgesetz

Die Formel hinter dem Eingangsbeispiel ist banal einfach und in ihren Konsequenzen brutal. Sei p die Wahrscheinlichkeit, dass ein einzelner Baustein langsam antwortet, und n die Anzahl der Bausteine, auf deren Antwort gewartet werden muss. Dann gilt für die Wahrscheinlichkeit, dass die Gesamtanfrage langsam wird:

P(langsam) = 1 − (1 − p)ⁿ

Bei kleinem p und moderatem n lässt sich das gut annähern als P ≈ n · p. Die Wahrscheinlichkeit für eine langsame Gesamtantwort wächst also in erster Näherung linear mit der Anzahl der beteiligten Komponenten. Das ist die eigentliche Aussage: Tail Latency ist nicht ein Problem, das mit Größe auch auftritt – es ist ein Problem, das durch Größe erzeugt wird.

Dean und Barroso geben ein zweites Beispiel, das die Größenordnung moderner Systeme trifft: Ein Dienst, der auf 2.000 Maschinen verteilt ist und bei dem nur eine von 10.000 Anfragen pro Maschine langsam ist, liefert dennoch bei fast jeder fünften Nutzeranfrage (1 − 0,9999²⁰⁰⁰ ≈ 18 Prozent) eine langsame Antwort. Man verbessert die Zuverlässigkeit der Einzelkomponente um den Faktor 100 – und die Situation ist immer noch schlecht, weil man gleichzeitig um den Faktor 20 skaliert hat.

Was das mit echten Zahlen macht

Dean und Barroso legen Messwerte aus einem Google-internen Dienst vor, der auf einem BigTable-Datenbestand arbeitet. Die Zahlen sind lehrreich, weil sie die Verstärkung an einem konkreten System sichtbar machen:

Messgröße (am Wurzeldienst gemessen) 99. Perzentil
Eine einzelne, zufällig gewählte Teilanfrage 10 ms
95 Prozent der Teilanfragen abgeschlossen 70 ms
Alle Teilanfragen abgeschlossen 140 ms

Zwischen „eine Teilanfrage" und „alle Teilanfragen" liegt ein Faktor 14. Und der Sprung von „95 Prozent fertig" zu „100 Prozent fertig" – also der Preis für die letzten fünf Prozent – verdoppelt die Latenz noch einmal. Wer diese Tabelle einmal verinnerlicht hat, versteht, warum die Frage „müssen wirklich alle Teilantworten da sein?" in verteilten Systemen eine Architekturfrage ersten Ranges ist.

Warum Microservice-Architekturen das Problem einbauen

Die beiden gefährlichen Muster lassen sich klar benennen.

Das erste ist Fan-out: Ein Dienst fragt viele Dienste parallel und wartet auf alle. Hier wirkt die Formel oben in voller Härte, weil sich die Ausreißerwahrscheinlichkeiten addieren.

Das zweite ist Tiefe: Eine Anfrage durchläuft eine Kette von Diensten sequenziell. Hier addieren sich zwar die Latenzen, aber auch hier gilt: Die Wahrscheinlichkeit, dass irgendwo auf dem Weg ein Ausreißer sitzt, wächst mit der Länge der Kette. Bei zehn Hops mit je p99 = 20 ms ist die Chance, dass mindestens ein Hop sein schlechtes Verhalten zeigt, bereits 1 − 0,99¹⁰ ≈ 9,6 Prozent.

In der Praxis kombinieren reale Architekturen beides – ein Aufrufgraph mit Fan-out auf jeder Ebene und mehreren Ebenen Tiefe. Die Ausreißerwahrscheinlichkeit multipliziert sich entlang des gesamten Graphen. Ich bin der Meinung, dass dies das am stärksten unterschätzte Argument gegen unnötig feingranulare Service-Aufteilungen ist: Jede zusätzliche Netzwerkgrenze im kritischen Pfad ist nicht nur ein zusätzlicher Millisekundenbetrag, sondern eine zusätzliche Ziehung aus einer schwerschwänzigen Verteilung.


Teil 3: Woher die langsamen Antworten kommen

Wenn man den Schwanz kürzen will, hilft es zu verstehen, woher er kommt. Dean und Barroso nennen die Quellen Variabilität – und sie zeigen, dass diese Variabilität in gemeinsam genutzter Infrastruktur nicht zufällig, sondern strukturell entsteht.

Geteilte Ressourcen und der laute Nachbar

Eine moderne Cloud-Maschine führt nicht einen Prozess aus, sondern viele – Container verschiedener Mandanten, Hintergrunddienste, Monitoring-Agenten. Sie teilen sich CPU-Kerne, Level-3-Cache, Speicherbandbreite, Netzwerkkarte und Festplatten-Queues. Jede dieser Ressourcen ist eine Kopplung, über die die Last eines Nachbarn zur Latenz des eigenen Prozesses wird. Der Fachbegriff noisy neighbour verharmlost das Phänomen: Es geht nicht um Lärm, sondern um Warteschlangen, in denen fremde Arbeit vor der eigenen steht.

Besonders heimtückisch ist die Kopplung über geteilte Caches. Wenn ein Nachbarprozess den Level-3-Cache mit seinen Daten füllt, steigt für den eigenen Prozess die Speicherzugriffszeit von Nanosekunden auf hunderte Nanosekunden – ein Effekt, den man in keiner Anwendungsmetrik sieht und der sich über Millionen von Speicherzugriffen zu Millisekunden summiert. Genau dieses Problem ist der Grund, warum Isolationstechnologien wie Die Wette auf fremden Code: Firecracker microVMs und das Ende des Container-VM-Dilemmas nicht nur ein Sicherheits-, sondern auch ein Latenzthema sind.

Wartungsarbeiten im Hintergrund

Die zweite große Quelle sind periodische Hintergrundaktivitäten, die dem Betrieb dienen, aber im Vordergrund wehtun:

  • Garbage Collection. Eine Pause der Speicherbereinigung in einer verwalteten Laufzeitumgebung (JVM, .NET, Go) stoppt genau diejenigen Threads, die gerade Anfragen bearbeiten. Selbst moderne, weitgehend parallele Kollektoren haben Phasen mit Stop-the-World-Charakter.
  • Compaction. Schreiboptimierte Speicher-Engines müssen ihre Daten regelmäßig reorganisieren. Wie dieser Mechanismus funktioniert und warum er zwangsläufig Lastspitzen erzeugt, ist ausführlich in Schreiben statt Suchen – Log-Structured Merge-Trees und die Umkehrung der Datenbank beschrieben: Der Preis für sehr schnelle Schreibvorgänge sind Hintergrund-Merges, die Festplatten-I/O und CPU in unregelmäßigen Schüben verbrauchen.
  • Log-Rotation, Index-Rebuilds, Metrik-Aggregation, Zertifikatserneuerung. Alles kleine Vorgänge – aber jeder einzelne ein Kandidat für einen 100-ms-Ausreißer, wenn er ungünstig fällt.
  • Energie- und Taktverwaltung. Prozessoren mit aggressiven Sparmodi brauchen Mikrosekunden bis Millisekunden, um aus einem tiefen Ruhezustand in den Volllastbetrieb zu wechseln. Thermisches Throttling wirkt in dieselbe Richtung.

Der entscheidende Punkt: Diese Aktivitäten sind alle legitim. Man kann sie nicht abschalten, nur verschieben, glätten oder – und das ist der architektonische Hebel – synchronisieren. Dean und Barroso empfehlen, Wartungsarbeiten über eine Replikagruppe hinweg gleichzeitig zu fahren statt unabhängig: Wenn alle Replikas gleichzeitig kurz langsam sind, gibt es ein kurzes, vorhersehbares Loch; wenn jede Replika unabhängig zufällig langsam wird, gibt es dauerhaft immer mindestens eine langsame Replika – und bei Fan-out trifft man die immer.

Queueing: warum Auslastung exponentiell bestraft wird

Die tiefste Ursache ist mathematischer Natur und hat mit Software gar nichts zu tun. Sie steckt in der Warteschlangentheorie.

Betrachte den einfachsten Fall: ein Server, Anfragen kommen zufällig (Poisson-verteilt) an, die Bearbeitungsdauer ist exponentialverteilt – das klassische M/M/1-Modell. Sei ρ die Auslastung (Ankunftsrate geteilt durch Bearbeitungsrate). Dann gilt für die mittlere Wartezeit in der Warteschlange, also ohne die eigentliche Bearbeitung:

W_q = (1 / μ) · ρ / (1 − ρ)

Der Faktor ρ/(1 − ρ) ist der entscheidende Term. Setze Zahlen ein:

Auslastung ρ Faktor ρ/(1−ρ) Wartezeit relativ zur Bearbeitungsdauer
50 % 1,0 1×
70 % 2,3 2,3×
80 % 4,0 4×
90 % 9,0 9×
95 % 19,0 19×
99 % 99,0 99×

Die Kurve ist eine Hyperbel mit Polstelle bei ρ = 1. Die praktische Konsequenz ist fundamental: Zwischen 50 und 70 Prozent Auslastung kostet zusätzliche Last fast nichts, zwischen 90 und 99 Prozent kostet sie alles. Und weil die Kurve nichtlinear ist, gilt dasselbe für die Streuung: Bei hoher Auslastung wird die Verteilung nicht nur nach rechts verschoben, sondern dramatisch breiter. Der Schwanz entsteht also nicht erst durch Fehler – er entsteht durch Auslastung.

Hier liegt der eigentliche Grund, warum „Kosteneffizienz durch hohe Auslastung" und „niedrige Tail Latency" einander ausschließende Ziele sind. Wer seine Maschinen auf 95 Prozent CPU fährt, hat sich für schlechte hohe Perzentile entschieden – nicht durch eine schlechte Implementierung, sondern durch die Wahl des Betriebspunkts. Jede Kapazitätsreserve, die man vorhält, ist in dieser Lesart nicht Verschwendung, sondern gekaufte Latenzstabilität.

Ergänzend liefert Littles Gesetz (L = λ · W) die Brücke zur Beobachtbarkeit: Die mittlere Anzahl gleichzeitig im System befindlicher Anfragen ist gleich der Ankunftsrate multipliziert mit der mittleren Verweildauer. Wer also die Anzahl offener Anfragen misst – eine Metrik, die fast jedes Framework kostenlos liefert –, hat damit indirekt einen sehr frühen Indikator für Latenzprobleme, oft früher als das Latenz-Histogramm selbst.

Die Mikrosekunden-Lücke

Ein vierter, vergleichsweise neuer Faktor: Die Hardwarelandschaft hat sich verschoben. In einem Folgeaufsatz, „Attack of the Killer Microseconds" (Barroso, Marty, Patterson und Ranganathan, Communications of the ACM, 2017), argumentieren die Autoren, dass unsere Werkzeuge für zwei Zeitskalen gut gebaut sind – Nanosekunden (Hardware-Parallelität, Out-of-Order-Ausführung, Prefetching) und Millisekunden (Kontextwechsel, asynchrone I/O) – aber für die Mikrosekunde fast nichts taugliches bereitstellen.

Die Zahlen aus dem Aufsatz machen das Problem greifbar: DRAM-Zugriff liegt bei zehn bis hunderten Nanosekunden, eine klassische Festplatte bei einigen Millisekunden, Flash-Speicher bei zehn Mikrosekunden, und die Durchquerung eines Rechenzentrums über 200 bis 300 Meter Kabelweg kostet etwa eine Mikrosekunde. Genau in diesem Bereich – Flash, RDMA, neue nichtvolatile Speicher – spielt sich moderne Cloud-Infrastruktur ab. Und hier versagen beide bekannten Strategien: Hardware-Parallelität kann Mikrosekunden nicht verbergen, und ein Software-Kontextwechsel kostet oft mehr als die Wartezeit, die er überbrücken soll. Das ist der Grund, warum Techniken wie Der programmierbare Kern: eBPF und die Sandbox im Herzen des Betriebssystems – Verarbeitung im Kernel, ohne den Umweg über den Userspace – für Latenz so relevant geworden sind.


Teil 4: Tail-tolerantes Design – die Techniken, die den Schwanz kürzen

Die begriffliche Verschiebung, die Dean und Barroso vorschlagen, ist die eigentliche Pointe ihres Aufsatzes: Man soll nicht versuchen, jede Komponente gleichmäßig schnell zu machen – das ist bei geteilter Infrastruktur, legitimen Hintergrundaufgaben und Warteschlangenphysik ein hoffnungsloses Programm. Stattdessen soll man Systeme bauen, die tail-tolerant sind: die also mit variablen Komponenten trotzdem eine stabile Gesamtlatenz liefern. Das ist dieselbe Denkbewegung, mit der man in der Zuverlässigkeitstechnik von „ausfallfreien Komponenten" zu „fehlertoleranten Systemen" übergegangen ist.

Hedged Requests: dieselbe Anfrage zweimal stellen

Die einfachste Technik: Schicke die Anfrage an eine Replika. Wenn nach einer kurzen Wartezeit – etwa der Dauer des 95. Perzentils – noch keine Antwort da ist, schicke dieselbe Anfrage an eine zweite Replika und nimm, was zuerst zurückkommt. Die zweite Anfrage wird abgebrochen, sobald eine Antwort vorliegt.

Der Effekt in Zahlen, gemessen an einem BigTable-Benchmark, bei dem 1.000 Werte von 100 Servern abgeholt werden: Eine Zweitanfrage nach 10 ms Verzögerung senkte das 99,9. Perzentil für das Einsammeln aller 1.000 Werte von 1.800 ms auf 74 ms – bei nur 2 Prozent zusätzlichen Anfragen. Das ist ein Faktor 24 für zwei Prozent Mehraufwand, und es ist der Grund, warum diese Technik so populär wurde.

Warum funktioniert das so gut? Weil die Verzögerung erst nach dem 95. Perzentil einsetzt, betrifft die Zweitanfrage nur die 5 Prozent der Fälle, die ohnehin langsam sind – und in diesen Fällen ist die Wahrscheinlichkeit, dass beide Replikas gleichzeitig einen Ausreißer haben, das Produkt zweier kleiner Zahlen. Man kauft die guten hohen Perzentile mit ein paar Prozent zusätzlicher Last. Der Preis ist reale Last, und das ist der Fallstrick: Ein zu früh gesetzter Schwellenwert verdoppelt den Traffic und treibt damit ρ nach oben – was nach der Tabelle oben genau das Gegenteil des Gewünschten bewirkt.

Tied Requests: beide Replikas informieren einander

Hedged Requests haben ein Fenster der Verschwendung: die Zeit zwischen erster und zweiter Anfrage. Tied Requests eliminieren es. Man schickt die Anfrage gleichzeitig an zwei Replikas, aber man teilt jeder mit, wer die andere ist. Sobald eine Replika die Anfrage aus ihrer Warteschlange nimmt und zu bearbeiten beginnt, sendet sie der Partnerreplika eine Abbruchnachricht. Die Partnerin verwirft die Anfrage, sofern sie noch in der Warteschlange liegt.

Man wählt damit nicht die Replika mit der kürzesten Warteschlange, sondern die, die tatsächlich als erste anfängt zu arbeiten – ein subtiler, aber wichtiger Unterschied, weil Warteschlangenlängen nichts über die Bearbeitungsdauer der wartenden Aufgaben sagen. Die Messwerte von Google: In einem ansonsten unbelasteten Cluster sank mit einer Verzögerung von 1 ms die mediane Latenz um 16 Prozent und das 99,9. Perzentil um fast 40 Prozent, bei unter 1 Prozent zusätzlicher Festplattenauslastung.

Die Technik hat eine Voraussetzung, die man nicht übersehen darf: Die Abbruchnachricht muss die Partnerreplika schneller erreichen als die Anfrage bearbeitet wird. In einem Rechenzentrum mit Mikrosekunden-Roundtrips ist das gegeben, über eine Weitverkehrsverbindung nicht.

Micro-Partitionierung und selektive Replikation

Die dritte Familie greift nicht am einzelnen Request an, sondern an der Datenverteilung.

Micro-Partitionierung bedeutet: Zerlege die Daten in deutlich mehr Partitionen als es Maschinen gibt – Faktoren von 10 bis 100 sind üblich – und weise Partitionen dynamisch Maschinen zu. Das hat zwei Effekte. Erstens wird Lastausgleich zu einer feinkörnigen Operation: Eine überlastete Maschine gibt einige ihrer Partitionen ab, statt dass man die gesamte Zuordnung neu berechnen muss. Zweitens wird Wiederherstellung schneller, weil die Partitionen einer ausgefallenen Maschine auf viele gesunde Maschinen verteilt werden können. Wie man eine solche Zuordnung so gestaltet, dass eine Änderung an der Maschinenmenge nur einen minimalen Bruchteil der Partitionen bewegt, ist das Thema von Der Ring, der die Last verteilt: Consistent Hashing und die Kunst des sanften Umzugs.

Selektive Replikation geht einen Schritt weiter: Erkenne besonders heiße Partitionen und erzeuge zusätzliche Kopien von genau diesen. Das ist die Antwort auf das Problem, dass Zugriffsverteilungen in der Realität fast nie gleichmäßig sind, sondern einem Potenzgesetz folgen – wenige Schlüssel machen den Großteil des Verkehrs.

Latency-induced Probation ist die dazugehörige Betriebsdisziplin: Ein System beobachtet die Antwortzeiten seiner Backends und nimmt eine auffällig langsame Replika vorübergehend aus der Rotation, während es weiter Schattenanfragen an sie schickt, um zu prüfen, wann sie sich erholt hat. Entscheidend ist, dass dies eine Latenz-basierte Entscheidung ist, nicht eine Fehler-basierte: Die Replika antwortet ja, nur zu langsam. Ein klassischer Health-Check auf Basis von „antwortet / antwortet nicht" würde sie für gesund erklären. Die Frage, wie ein Cluster überhaupt zu einer gemeinsamen Sicht darauf kommt, welche Knoten in welchem Zustand sind, behandelt Das Flüstern der Maschinen: Gossip-Protokolle, SWIM und wie ein Cluster erfährt, wer noch lebt.

Canary Requests und „gut genug" statt „vollständig"

Zwei weitere Techniken aus dem Aufsatz verdienen Erwähnung, weil sie so wenig Aufwand kosten.

Canary Requests schützen gegen den Fall, dass eine Anfrage einen latenten Fehler in allen Replikas gleichzeitig auslöst – ein Fan-out an 10.000 Server, der alle 10.000 zum Absturz bringt, ist ein realer Vorfallstyp. Die Gegenmaßnahme: Schicke die Anfrage erst an einen oder zwei Server, und erst wenn diese erfolgreich antworten, fahre den vollen Fan-out. Die Kosten sind ein zusätzlicher Roundtrip; der Nutzen ist die Vermeidung eines Totalausfalls.

„Good enough"-Antworten sind die vielleicht wirkungsvollste und am seltensten genutzte Technik. Sie besagt: Definiere für deinen Dienst, was eine unvollständige, aber brauchbare Antwort ist, und liefere sie, wenn die Frist abläuft. Bei einer Suche mit 100 Shards ist ein Ergebnis aus 98 Shards fast immer so gut wie eines aus 100 – und laut der Tabelle aus Teil 2 spart es den Faktor zwei. Die Voraussetzung ist eine Produktentscheidung, nicht eine technische: Man muss bereit sein, Vollständigkeit gegen Vorhersagbarkeit zu tauschen.

Die Wahl des Ziels: Power of Two Choices

Wo zwei oder mehr Replikas in Frage kommen, stellt sich die Frage, wie man wählt. Die Antwort darauf gehört zu den elegantesten Ergebnissen der Informatik. Michael Mitzenmacher zeigte in The Power of Two Choices in Randomized Load Balancing (IEEE Transactions on Parallel and Distributed Systems 12(10), S. 1094–1104, 2001), dass es einen qualitativen Unterschied macht, ob man eine oder zwei Möglichkeiten zur Wahl hat.

Wenn man n Aufgaben rein zufällig auf n Server verteilt, wächst die maximale Warteschlangenlänge in der Größenordnung von log n / log log n. Wenn man dagegen für jede Aufgabe zwei Server zufällig zieht und die mit der kürzeren Warteschlange nimmt, sinkt die maximale Länge auf etwa log log n / log 2 – ein exponentieller Gewinn. Der Unterschied ist qualitativ, nicht graduell: log n wächst unbeschränkt mit der Systemgröße, log log n praktisch nicht – bei einer Million Servern liegt log₂ log₂ n bei etwa vier. Und die Asymptotik sagt noch etwas Wichtiges: Der Sprung von einer auf zwei Wahlmöglichkeiten bringt fast den gesamten erreichbaren Gewinn, der Sprung von zwei auf drei nur noch wenig.

Das ist für den Praktiker eine außerordentlich gute Nachricht. Ein Load Balancer, der zwei Kandidaten zufällig zieht und den weniger belasteten wählt („power of two random choices", oft P2C genannt), braucht keinen globalen Zustand, keine Koordination und keine vollständige Sicht auf den Cluster – und kommt der idealen Balance dennoch sehr nahe. Verfahren, die alle Server befragen, sind nicht nur teurer, sie sind auch anfällig für Herdenverhalten, weil viele Clients gleichzeitig denselben „besten" Server entdecken.

Shuffle Sharding: Isolation als Latenzwerkzeug

Eine letzte, oft unterschätzte Technik kommt aus der Betriebspraxis großer Cloud-Anbieter und ist in der AWS Builders' Library von Colm MacCárthaigh beschrieben. Die Grundidee: Wenn ein einzelner Mandant durch eine pathologische Anfrage seine Backends langsam macht, sollen möglichst wenige andere Mandanten davon betroffen sein.

Bei klassischem Sharding teilt man acht Worker in vier Paare; ein Problem trifft ein Viertel aller Mandanten. Bei Shuffle Sharding weist man jedem Mandanten eine zufällig gezogene Kombination von zwei der acht Worker zu. Davon gibt es C(8,2) = 28 verschiedene – die Auswirkung eines Problems schrumpft auf ein Achtundzwanzigstel, also auf ein Siebtel dessen, was klassisches Sharding liefert, ohne dass man eine einzige Maschine hinzufügen musste.

Der Effekt skaliert kombinatorisch atemberaubend. Amazon Route 53 ordnet seine Kapazität in 2.048 virtuelle Nameserver und weist jeder Kundendomäne eine Kombination aus vier davon zu. Das sind rund 730 Milliarden mögliche Kombinationen – mit der Konsequenz, dass keine Kundendomäne mit einer anderen mehr als zwei virtuelle Nameserver teilt. Ein Mandant, dessen Verkehr ein Backend in die Knie zwingt, kann daher praktisch keinen anderen Mandanten vollständig lahmlegen.

Shuffle Sharding ist damit die strukturelle Antwort auf einen Teil des Tail-Problems: Sie begrenzt nicht die Variabilität selbst, sondern ihren Ausbreitungsradius.

Eine Landkarte der Techniken

Technik Greift an bei Zusätzliche Last Voraussetzung
Hedged Requests einzelne langsame Anfrage wenige Prozent, wenn Schwelle > p95 idempotente Leseoperationen, ≥ 2 Replikas
Tied Requests Warteschlangenzeit vor Bearbeitung < 1 % schnelle Abbruchnachricht zwischen Replikas
Micro-Partitionierung Lastungleichheit, Wiederherstellung Metadaten-Overhead dynamische Partitionszuordnung
Selektive Replikation heiße Partitionen Speicher für Zusatzkopien Erkennung von Hotspots
Latency-induced Probation dauerhaft langsame Replika Schattenanfragen Latenzmessung pro Backend
Canary Requests anfrageinduzierte Gesamtausfälle ein Roundtrip Fan-out in zwei Stufen
„Good enough"-Antworten die letzten Prozent des Fan-out keine Produktentscheidung über Vollständigkeit
Power of Two Choices Zielauswahl vernachlässigbar grobe Lastinformation über 2 Kandidaten
Shuffle Sharding Ausbreitungsradius eines Problems keine Mandanten-zu-Kapazität-Zuordnung steuerbar
Kapazitätsreserve (ρ senken) die Wurzel des Problems Infrastrukturkosten Bereitschaft, nicht auf 95 % zu fahren

Teil 5: Was in der Praxis wehtut

Retry-Stürme und metastabiles Versagen

Die naheliegendste Reaktion auf langsame Antworten – „dann versuche es einfach nochmal" – ist auch die gefährlichste. Wenn ein Dienst langsam wird, weil er überlastet ist, und alle Clients daraufhin ihre Anfragen wiederholen, dann steigt die Last genau in dem Moment, in dem das System sie am wenigsten verkraftet. Nach der Queueing-Tabelle aus Teil 3 verschiebt das ρ nach oben, die Latenz explodiert, mehr Timeouts laufen ab, mehr Wiederholungen entstehen. Die Rückkopplung ist positiv – im mathematischen, nicht im umgangssprachlichen Sinn.

Bronson, Aghayev, Charapko und Zhu haben diesem Phänomen 2021 einen präzisen Namen gegeben: metastabiles Versagen. Der Kern ihrer Analyse: Ein System kann in einen Zustand geraten, in dem es dauerhaft überlastet bleibt, obwohl die auslösende Störung längst vorbei ist. Die Überlast trägt sich selbst, weil die Arbeit, die sie erzeugt (Wiederholungen, Cache-Misses nach Cache-Leerung, abgebrochene und neu begonnene Transaktionen), größer ist als die Arbeit, die sie erledigt. Die praktische Konsequenz ist unangenehm: Solche Zustände verlassen ein System nicht von selbst. Man muss die Last aktiv wegnehmen – durch Lastabwurf, Drosselung oder, im schlimmsten Fall, durch einen kontrollierten Neustart.

Die Gegenmittel sind bekannt und in der AWS Builders' Library (Marc Brooker, Timeouts, retries, and backoff with jitter) beschrieben: exponentieller Backoff, Jitter – also eine Zufallskomponente in der Wartezeit, damit Clients nicht synchron wiederholen –, Wiederholungsbudgets, die Wiederholungen auf einen kleinen Prozentsatz des Gesamtverkehrs begrenzen, und Circuit Breaker, die einen als ungesund erkannten Pfad vorübergehend ganz abschalten. Besonders wichtig ist die Regel, Wiederholungen nicht über mehrere Schichten zu stapeln: Drei Schichten mit je drei Versuchen ergeben im schlimmsten Fall 27 Anfragen für eine Nutzeraktion.

Incast

Ein rein netzwerktechnischer Effekt, der bei Fan-out regelmäßig zuschlägt: Wenn ein Client gleichzeitig 100 Server befragt und alle 100 gleichzeitig antworten, treffen 100 Antworten praktisch simultan auf derselben Netzwerkkarte ein. Der Switch-Puffer davor läuft über, Pakete gehen verloren, TCP interpretiert das als Stau und wartet – und aus einer Mikrosekunden-Antwort wird eine Antwort mit TCP-Retransmission-Timeout. Das Team um Rajesh Nishtala beschreibt in Scaling Memcache at Facebook (NSDI 2013), wie sie diesen Effekt mit einem gleitenden Fenster über die ausstehenden Anfragen bekämpften: Der Client schickt nicht alle Anfragen gleichzeitig, sondern hält ein Limit an offenen Anfragen ein. Zu klein gewählt, steigt die Latenz durch unnötige Serialisierung; zu groß gewählt, kommt Incast zurück.

Timeouts sind Architekturentscheidungen

Ein Timeout ist keine Sicherheitsmaßnahme, sondern eine Behauptung über die Latenzverteilung. Setzt man ihn auf den Mittelwert plus etwas Puffer, bricht man einen erheblichen Teil der Anfragen ab, die kurz vor dem Erfolg standen – und erzeugt damit Last ohne Nutzen. Setzt man ihn zu hoch, blockieren Threads und Verbindungen so lange, dass die Überlast sich staut. Der vernünftige Ausgangspunkt ist ein gemessenes hohes Perzentil des nachgelagerten Dienstes, und der vernünftige Umgang ist das Weitergeben einer verbleibenden Frist (deadline propagation): Der oberste Dienst legt fest, wie viel Zeit die gesamte Anfrage haben darf, und jeder nachgelagerte Aufruf bekommt das Restbudget mitgeteilt, statt jeweils seinen eigenen Timeout neu zu erfinden.

SLO-Arithmetik

Wer Service Level Objectives formuliert, sollte die Fan-out-Formel im Kopf haben. Ein Latenz-SLO wie „99 Prozent aller Anfragen unter 300 ms" ist für einen Dienst mit zwanzig internen Aufrufen nur haltbar, wenn jeder dieser Aufrufe ein deutlich strengeres Ziel einhält – nicht 300/20 = 15 ms im Sinne einer Addition, sondern ein Perzentilziel, das so hoch liegt, dass 1 − (1 − p)²⁰ noch unter einem Prozent bleibt. Das erfordert p < 0,0005, also ein p9995 pro Aufruf. Solche Zahlen sind unbequem, aber sie sind die ehrliche Rechnung.


Teil 6: Das übergeordnete Denkprinzip

Wenn man einen Schritt zurücktritt, ist Tail Latency ein Spezialfall eines viel allgemeineren Musters: In Systemen, die aus vielen Teilen bestehen und deren Gesamtergebnis vom schlechtesten Teil bestimmt wird, ist nicht der Durchschnitt die relevante Größe, sondern die Streuung. Ein System, dessen Komponenten im Mittel langsam, aber gleichmäßig sind, verhält sich vorhersehbar. Ein System, dessen Komponenten im Mittel schnell, aber gelegentlich katastrophal sind, verhält sich unvorhersehbar – und Unvorhersehbarkeit ist das, was Nutzer als Unzuverlässigkeit erleben.

Dasselbe Muster findet man in vielen gut entworfenen verteilten Systemen wieder: Uhren, die ihre eigene Unsicherheit kennen: Google Spanner, TrueTime und die Beherrschung der Zeit in der Cloud ersetzt einen Zeitstempel, der stillschweigend genau zu sein behauptet, durch ein Intervall, dessen Unsicherheit explizit mitgeführt wird – die Streuung wird zum erstklassigen Datum statt zum verschwiegenen Risiko. Elf Neunen: Erasure Coding, Reed-Solomon und wie die Cloud Daten praktisch unverlierbar macht baut Haltbarkeit nicht aus haltbaren Platten, sondern aus rechnerischer Redundanz über unzuverlässige Platten. Wie Maschinen sich einig werden: Verteilter Konsens von FLP über Paxos zu Raft erreicht Fortschritt nicht dadurch, dass alle Knoten antworten, sondern dadurch, dass eine Mehrheit genügt – Quorumslogik ist im Kern eine Tail-Toleranz-Technik.

Erkenntnis zum Mitnehmen

Die praktische Lehre lautet: Hör auf, Mittelwerte zu betrachten, und fang an, eine Verteilung zu betrachten – und rechne dann aus, wie viele Ziehungen aus dieser Verteilung eine einzige Nutzeranfrage erfordert. Konkret sind das vier Schritte, die man in jedem Projekt gehen kann. Erstens: Erhebe Latenzhistogramme statt Mittelwerte, und prüfe, ob das Messwerkzeug gegen Coordinated Omission gewappnet ist – ein Dashboard, das p999 zu gut anzeigt, ist schlimmer als keines. Zweitens: Zähle den Fan-out und die Tiefe deines kritischen Pfades und rechne 1 − (1 − p)ⁿ aus; die Zahl, die dabei herauskommt, ist die ehrliche Beschreibung deines Nutzererlebnisses. Drittens: Schau dir die Auslastung an – wenn deine Dienste über 80 Prozent laufen, ist Kapazität und nicht Code die billigste Latenzverbesserung, die du kaufen kannst. Viertens: Wähle bewusst eine tail-tolerante Technik aus der Tabelle in Teil 4, statt darauf zu hoffen, jede Komponente gleichmäßig schnell zu machen – dieses Ziel ist bei geteilter Infrastruktur physikalisch unerreichbar.

Eine Frage zum Nachdenken

Alle Techniken in diesem Artikel kaufen Latenzstabilität mit einer Ressource: mit zusätzlicher Last (Hedged Requests), mit zusätzlichem Speicher (selektive Replikation), mit ungenutzter Kapazität (niedriges ρ) oder mit Vollständigkeit („good enough"-Antworten). Keine ist kostenlos. Die Frage, die daraus folgt, ist keine technische, sondern eine über Werte: Wie viel Vollständigkeit wärst du in deinem eigenen System bereit aufzugeben, um Vorhersagbarkeit zu gewinnen – und wer in deiner Organisation darf diese Entscheidung eigentlich treffen? Denn erfahrungsgemäß wird sie selten explizit getroffen. Sie wird implizit getroffen, in dem Moment, in dem jemand einen Timeout-Wert in eine Konfigurationsdatei schreibt.


Querverweise im Vault


Quellen

← All articles