Startseite / Artikel / Von Modulo zu Hash Ring: Skalierung einer Node.js-Cache-Flotte ohne Ausfälle

Von Modulo zu Hash Ring: Skalierung einer Node.js-Cache-Flotte ohne Ausfälle

Erfahren Sie, warum das Hash-Mod-N-Sharding Datenbanken zusammenbricht, wenn sich Knoten ändern, wie sich Rendezvous-, Jump- und Ring-Hashing vergleichen lassen, sowie wie man in Node.js einen ausgewogenen gewichteten Ring erstellt.

4407 Wörter

Ein Cache-Cluster, der monatelang reibungslos funktioniert hat, kann bereits wenige Minuten nach einer Routineänderung – dem Hinzufügen eines Nodes – die Datenbank lahmlegen. Die Ursache liegt in der Regel in einer einzigen Zeile Client-Code, die einen Server mit hash(key) % N auswählt. Diese Anleitung erklärt genau, warum diese Zeile versagt, vergleicht die vier ernsthaften Alternativen, erstellt einen für die Produktion geeigneten konsistenten Hash-Ring in Node.js und listet die Betriebsfallrisiken auf, vor denen der Algorithmus selbst keinen Schutz bietet.

Das Ausfallmuster

Stellen Sie sich eine funktionierende Flotte von vier Cache-Nodes mit einer Trefferquote von 94 % vor. Aufgrund eines saisonalen Anstiegs des Verkehrs fügt ein Ingenieur einen fünften Node hinzu. Es handelt sich dabei um eine Änderung in zwei Zeilen, die sorgfältig mitten am Arbeitstag implementiert wird. Etwa eineinhalb Minuten später ist die CPU-Auslastung der Datenbank bei 100 % und die Website ist nicht mehr verfügbar.

Niemand hat unvorsichtig gehandelt. Der Kunde hat die Schlüssel einfach auf die übliche Weise weitergeleitet:

const node = nodes[hash(key) % nodes.length];

Der Ausdruck verteilt die Schlüssel sehr gleichmäßig, wodurch er korrekt erscheint. Das Problem entsteht, sobald sich nodes.length ändert – genau hier beginnt der Rest dieses Leitfadens. In der Diskussion werden vier separate Probleme behandelt, die tatsächlichen Alternativen abgewogen und anschließend ein Ring in Node.js erstellt sowie getestet.

Vier Probleme hinter einer Idee

Konsistentes Hashing wird oft als ein einziger Trick vorgestellt. In der Praxis löst es vier unterschiedliche Probleme, und eine Implementierung, die nur das erste Problem berücksichtigt, wird aufgrund der anderen drei Probleme in der Produktion weiterhin versagen.

Problem 1: Eine Änderung von N verschiebt fast alle Schlüssel

Mit hash(key) % N führt eine Änderung von N nicht dazu, dass nur einige Schlüssel umgesetzt werden. Stattdessen werden fast alle Schlüssel umgesetzt.

Das Durchgehen der Zahlen hilft dabei. Ein Schlüssel mit Hash 1.000.003 wird dem Knoten 3 unter % 4 zugeordnet und zufällig auch dem Knoten 3 unter % 5. Ein Schlüssel mit Hash 1.000.004 wird unter % 4 dem Knoten 0 und unter % 5 dem Knoten 4 zugeordnet. Diese beiden Zuordnungen stehen in keiner Beziehung zueinander, sodass ein Schlüssel nur zufällig an seinem ursprünglichen Platz bleibt – mit einer Wahrscheinlichkeit von etwa 1 zu N.

Betrachtet man eine Million Schlüssel, ist das Muster auffällig: Der Wechsel von 8 auf 9 Knoten bewirkt, dass 88,93 % der Schlüssel verschoben werden, und in einem Cluster mit hundert Knoten macht das Hinzufügen eines Knotens etwa 99 % des Caches ungültig.

Betrachten Sie, in welche Richtung diese Entwicklung geht. Je größer Ihr System wird, desto zerstörerischer wirkt jeder Skalierungsschritt. Es handelt sich um einen Fehler, der erst dann zum Problem wird, wenn das Unternehmen erfolgreich ist.

Jeder umgeleitete Schlüssel ist ein Fehlschlag, jeder Fehlschlag ist eine Datenbankabfrage – und all diese kommen innerhalb von Sekunden an. Eine Datenbank, die ursprünglich für die 6 % der Lesevorgänge ausgelegt ist, die normalerweise fehlschlagen, erhält plötzlich fast alle dieser Abfragen.

Problem 2: dieselbe Umordnung, unplanmäßig

Zumindest tritt das erste Problem auf, wenn man sich für eine Skalierung entscheidet. Das zweite Problem wird durch einen Ausfall ausgelöst – und zwar zum unpassendsten Zeitpunkt.

Ein Knoten erschöpft seinen Speicher, ein Host wird beendet oder eine Netzwerkpartition verdeckt einen Knoten vor der Hälfte des Systems. Die Anzahl der Knoten sinkt von 8 auf 7, und jeder Client leitet etwa 87 % seiner Schlüssel eigenständig auf die verbleibenden Knoten um.

Die Situation ist nun ernst: 12,5 Prozent der Cache-Kapazität sind bereits verloren gegangen, die Datenbank muss mit 87 Prozent Fehlversuchen umgehen, und die sieben verbleibenden Knoten übernehmen den Verkehr des fehlenden Knotens gleichzeitig, während sie fast mit jeder Schlüsselwert-Kombination wieder aufgefüllt werden. Häufig führt dies dazu, dass ein zweiter Knoten zusammenbricht, was wiederum eine erneute vollständige Umzuordnung erzwingt und schließlich einen dritten Knoten zum Kollabieren bringt.

Durch Modulo-Routing verwandelt sich ein Ausfall eines Knotens in einen miteinander verbundenen, im gesamten Cluster auftretenden Ausfall, der sich selbst verstärkt. Diese Kettenreaktion – und nicht ein leerer Cache – stellt die eigentliche Gefahr dar.

Problem 3: Ein naiver Ring ist stark unausgeglichen

Die offensichtliche Lösung für Problem 1 besteht darin, Knoten und Schlüssel in einen numerischen Raum zu bringen und jedem Schlüssel den Knoten zuzuweisen, der im Uhrzeigersinn folgt. Das ist das Prinzip des konsistenten Hashings, und es beseitigt tatsächlich die massiven Umordnungen.

Wenn es jedoch naiv implementiert wird, balanciert es die Last schlecht. Die Knoten befinden sich dort, wo ihre Hash-Werte landen, wodurch die Abstände zwischen ihnen zufällig sind – und zufällige Abstände sind selten gleich. Bei vier Knoten, von denen jeder durch einen Hash positioniert wird, hatte eine Messung einen einzelnen Knoten, der 45% des Schlüsselraums einnahm, während ein anderer nur 13% hatte – das ergibt einen Unterschied von 3,4-fach, ohne dass es irgendwelche Fehler gibt.

Auch das Ungleichgewicht bleibt bestehen. Es resultiert aus den Knotennamen selbst, sodass derselbe Knoten, zum Beispiel cache-04, weiterhin stark belastet bleibt, bis er umbenannt wird – und jeder, der nachforscht, findet Code, der genau so funktioniert, wie er geschrieben ist.

In größerem Maßstab wird die Situation noch schlimmer. Bei einem Ringpunkt pro Knoten über insgesamt 8 Knoten trug der am stärksten belastete Knotel 434% seines fairen Anteils, während der am wenigsten belastete nur 3,7% trug. Das entspricht im Grunde einem überlasteten Server und sieben inaktiven Servern.

Problem 4: Jeder Client muss zustimmen

Das am wenigsten sichtbare Problem betrifft die Autorität. Etwas muss einen Schlüssel auf einen Knoten abbilden und muss auf jedem abfragenden Gerät dieselbe Antwort liefern.

Eine Möglichkeit ist ein Koordinationsdienst, der die autoritative Shard-Karte speichert. Dann kostet jede Abfrage entweder eine Netzwerkübertragung hin und zurück, oder die Clients speichern die Karte im Cache – in diesem Fall braucht man eine Methode, um sie für ungültig zu erklären. Sollten zwei Clients auch nur vorübergehend unterschiedliche Versionen der Karte besitzen, könnte einer user:42 in Knoten A schreiben, während ein anderer es aus Knoten B liest. Dabei geht nichts verloren – was jedoch noch schlimmer sein kann: Es gibt nun zwei mögliche Werte.

Was Sie benötigen, ist eine Zuordnung, die eine reine Funktion der Schlüssel und der aktuellen Mitgliederliste ist. Kein Koordinator, kein Suchdienst, kein gemeinsamer Zustand: Jeder Client führt dieselben Berechnungen durch und kommt zum selben Ergebnis. Konsistentes Hashing bietet genau das, weshalb es trotz der größeren Flexibilität von Suchtabellen vorzuziehen ist.

Ein konkretes Szenario und die Optionen

Um den Vergleich konkret zu machen, betrachten Sie dieses System.

Das System. Eine E-Commerce-API speichert Sitzungs- und Profildaten in einem Pool von Redis-Noden. Die Spitzenlast beträgt etwa 40.000 Abfragen pro Sekunde bei 25 Millionen Schlüsseln, wobei die Trefferquote bei 94 % liegt. Die Datenbank überlebt nur, weil sie die 6 % der Fälle berücksichtigt, bei denen keine Treffer erzielt werden. (Für einen Überblick über die Caching-Muster selbst siehe Grundlagen des Redis-Cachings.)

Die Anforderungen:

  1. Von 8 auf 12 Noden anwachsen, bevor ein Verkauf stattfindet, ohne eine „Miss-Storm“ auszulösen
  2. Den Verlust eines Nodes überstehen, wobei der betroffene Bereich begrenzt und erträglich sein muss
  3. Die Last gleichmäßig verteilen, sodass kein Node mehr als etwa 120 % seines fairen Anteils trägt
  4. Jeden Koordinator vom Leseweg fernhalten
  5. Auf gemischte Hardware Rücksicht nehmen: Einige Node haben 64 GB, andere 16 GB – sie sollten nicht dieselbe Last tragen

Mehrere Algorithmen können einige oder alle dieser Anforderungen erfüllen. Sie unterscheiden sich in bedeutenden Aspekten, und eine schlechte Wahl hat erhebliche Folgen.

Option A: Modulo-Hashing

hash(key) % N gewährleistet ein perfektes Gleichgewicht, erfordert nur eine Anweisung und verbraucht kein Speicherplatz.

Er erfüllt die Anforderungen 1 und 2 völlig nicht. Er wird dennoch erwähnt, weil jeder ihn zuerst verwendet und er sich bis zum Zeitpunkt des Versagens perfekt verhält.

Wählen Sie ihn, wenn N tatsächlich niemals ändert, beispielsweise bei der Aufteilung einer Batch-Aufgabe auf eine feste Anzahl von Prozessen oder beim Sharding innerhalb eines einzigen Prozesses.

Option B: Koordinator und Abfrageschema

Bewahren Sie eine explizite Zuordnung von Schlüsselbereichen zu Knoten in einem Speicher wie etcd oder ZooKeeper auf. Systeme wie Vitess und HBase funktionieren in etwa auf diese Weise.

Der Vorteil ist echt: vollständige Kontrolle. Sie können einen einzelnen Hot Shard verschieben, schrittweise einen bestimmten Bereich neu ausbalancieren, während Sie die Metriken überwachen, oder einen bestimmten Tenant auf spezifische Hardware binden. Kein hash-basiertes Verfahren bietet etwas davon, und ab einer bestimmten Skalierung werden Sie dies benötigen.

Der Preis ist genauso real: ein zu betreibendes Konsenssystem, die Herausforderung, jede gespeicherte Kopie der Karte aktuell zu halten, sowie eine notwendige Abhängigkeit vom Leseweg.

Wählen Sie dieses Verfahren, wenn Sie dauerhafte Daten statt einmalig verwendbarer Cache-Einträge verschieben und die Migration kontrollieren müssen, anstatt alles auf einmal umzustellen.

Option C: Rendezvous-Hashing (HRW)

Das Hashing mit dem höchsten zufälligen Gewicht vergleicht jeden Schlüssel mit jedem Knoten und wählt den Gewinner aus:

function rendezvous(key, nodes) {
  let best = null, bestScore = -1;
  for (const node of nodes) {
    const score = mix(hash(key), hash(node));
    if (score > bestScore) { bestScore = score; best = node; }
  }
  return best;
}

Das ist der gesamte Algorithmus. Es gibt keinen Ring, keine virtuellen Knoten, keine sortierte Struktur und nichts, was bei Änderungen der Mitgliedschaft neu aufgebaut werden müsste.

Was die beiden wichtigsten Metriken angeht, übertrifft er auch den Ring-Algorithmus. Der Wechsel von 8 auf 9 Knoten führte zu einer Verschiebung von 11,09 % der Schlüssel – gegenüber einem theoretischen Minimum von 11,11 % – und das Gleichgewicht war ohne jegliche Anpassungen nahezu perfekt.

Der Nachteil ist die Laufzeit von O(N) pro Abfrage, da jeder Schlüssel gegen jeden Knoten gehäshed wird, und dieser Aufwand nimmt mit zunehmender Anzahl der Knoten schnell zu.

Wählen Sie ihn, wenn Sie weniger als etwa 30 Knoten haben. Viele Teams verwenden etwa acht Cache-Knoten, für die rendezvous Hashing am besten geeignet ist: Es ist einfacher zu implementieren und zu verstehen, und es gewährleistet ein besseres Gleichgewicht. Der Ring-Algorithmus ist zwar besser bekannt, ist aber nicht automatisch der bessere.

Option D: konsistenter Sprung-Hash

Dieser Algorithmus, der 2014 von Google veröffentlicht wurde, passt in etwa zehn Zeilen, benötigt kein Speicher, balanciert fast perfekt und bewegt die minimale Anzahl an Schlüsseln.

Ihre Grenze ist strukturell bedingt. Er ordnet einen Schlüssel einer Bucket-Nummer im Bereich [0, N) zu und kennt kein Konzept einer Identität des Knotens. Buckets können nur am Ende des Bereichs hinzugefügt oder entfernt werden; es gibt keine Möglichkeit, Knoten 3 aus der Mitte zu nehmen, ohne alles andere stabil zu halten.

Wählen Sie ihn, wenn die Buckets austauschbar sind und nur ihre Anzahl sich ändert, wie beispielsweise beim Aufteilen eines Datensatzes auf eine anpassbare Gruppe von Arbeitsknoten. Er eignet sich nicht, wenn spezifische, benannte Server hinzukommen oder wieder gehen – genau so verhalten sich Cache-Knoten.

Option E: ein Hash-Ring mit virtuellen Knoten

Das ist das klassische Design. Knoten und Schlüssel teilen sich einen kreisförmigen Adressraum, wobei ein Schlüssel zum ersten Knoten gehört, der im Uhrzeigersinn davon gefunden wird.

Wählen Sie dieses Design, wenn die Flotte groß genug ist, sodass die linearen Kosten für Suchvorgänge unzumutbar werden, und Sie außerdem gewichtete Knoten sowie die Möglichkeit benötigen, beliebige Knoten zu entfernen.

Auswahl für das Szenario

Betrachtet man einen Wechsel von 8 auf 9 Knoten bei einer Million Schlüsseln, liegen alle Alternativen zum Modulo-Verfahren nahe am theoretischen Minimum. Das Gleichgewicht des Rings hängt stark davon ab, wie viele virtuelle Knoten jeder Server erhält. Für dieses E-Commerce-System mit gemischter Hardware, benannten Knoten, die ausfallen können, und einer Flotte, die voraussichtlich 30 Knoten umfasst, ist der Ring die richtige Wahl. Der Rest dieses Leitfadens zeigt, wie man ihn richtig implementiert.

Wie der Ring funktioniert

Lassen Sie Arrays und Reste beiseite. Stellen Sie sich einen Kreis vor, der von 0 bis 2³² − 1 nummeriert ist und an der Oberseite wieder beginnt.

Zwei Regeln bilden den gesamten Algorithmus:

  1. Hashen Sie jeden Knotennamen auf den Kreis, sodass cache-01 dort landet, wohin sein Hash es führt.
  2. Hashen Sie jeden Schlüssel auf denselben Kreis und bewegen Sie sich anschließend im Uhrzeigersinn. Der erste erreichte Knoten besitzt den Schlüssel.

Die entscheidende Erkenntnis ist, dass Knoten und Schlüssel einen Adressraum teilen. Alles Weitere folgt daraus, einschließlich der Tatsache, warum das Hinzufügen eines Knotens kostengünstig ist.

Warum Änderungen der Mitgliedschaft lokal bleiben

Fügen Sie einen neuen Knoten auf den Kreis hinzu – er liegt dann zwischen zwei bereits vorhandenen Knoten. Er übernimmt nur den Bogen zwischen sich und seinem gegen den Uhrzeigersinn liegenden Nachbarn.

Schlüssel außerhalb dieses Bogens bleiben unberührt und gelangen genau wie zuvor zu ihrem ursprünglichen Besitzer. Der Neuzugang erhält im Durchschnitt etwa 1/(N+1) des Kreises, wodurch dieser Anteil an Schlüsseln verschoben wird. Bei einer Veränderung von 8 auf 9 betrug dieser Wert 11,06%, wobei ein Mindestwert von 11,11% gilt – im Vergleich zu 88,93% bei der Modulo-Methode unter denselben Bedingungen.

Das Entfernen funktioniert umgekehrt: Der Bogen des entfernten Knotens geht an seinen nach rechts gerichteten Nachfolger über. Beim Rückgang von 8 auf 7 Knoten wurden 12,60% der Schlüssel verschoben, was nahe am theoretischen Wert von 12,50% liegt. Die Auswirkungen sind begrenzt und erträglich, und die anderen sechs Knoten bleiben unberührt.

Virtuelle Knoten beheben das Ungleichgewicht

Zurück zu Problem 3: Vier Knoten an vier zufälligen Positionen erzeugen sehr ungleichmäßige Bögen.

Die Lösung ist überraschend einfach: platzieren Sie jeden Knoten nicht nur einmal. Plazieren Sie ihn 160 Mal unter 160 abgeleiteten Namen wie cache-01#0 und cache-01#1. Jeder abgeleitete Name landet an einer anderen Stelle, sodass jeder physische Knoten 160 kleine, verstreute Bögen statt eines großen besitzt – und das Gesetz der großen Zahlen sorgt für Ausgleich.

Betrachtet man eine Million Schlüssel auf 8 Knoten, verbessert sich das Gleichgewicht stetig mit steigender Anzahl an Replikaten. 160 ist der übliche Standard, weil dort die Verbesserungskurve in etwa flacher wird – doch 500 ist immer noch deutlich besser. Ein Ringpunkt benötigt etwa 12 Bytes (eine 4-Byte-Position plus eine 8-Byte-Eigentümerreferenz); daher nehmen acht Knoten mit 500 Replikaten weniger als 50 KB in Anspruch. Wenn Ihnen das Gleichgewicht mehr zählt als dieser Speicherbedarf, wählen Sie eine höhere Zahl – das ist eine Berechnung, die nur wenige Teams sich machen.

Virtuelle Knoten machen das Gewichtungssystem im Grunde kostenlos. Ein Knoten mit doppelter Speicherkapazität erhält doppelte Punkte und somit in etwa auch doppelt so viel Traffic. Bei Gewichten von 4:4:1:1 lag die tatsächliche Verteilung bei 40,6 %, 41,1 %, 9,4 % und 9,0 %, im Vergleich zum Idealwert von 40/40/10/10.

Implementierung des Rings in Node.js

Die Implementierung besteht aus vier Schritten: Die Namen der virtuellen Knoten werden gehasht, platziert, sortiert und anschließend mittels binärer Suche wird der Eigentümer einer Schlüsselwert-Kombination ermittelt. Die untenstehende Klasse speichert die Mitgliedschaften in einem Map, rebuildet die sortierten Arrays bei Änderungen der Mitgliedschaften und stellt get(key) zur Abfrage bereit.

export class ConsistentHashRing {
  #positions = new Uint32Array(0); // sorted ring positions
  #owners = []; // owners[i] owns #positions[i]
  #nodes = new Map(); // id -> { weight, points }

  constructor({ replicas = 160, hash = defaultHash } = {}) {
    if (replicas < 1) throw new RangeError("replicas must be >= 1");
    this.replicas = replicas;
    this.hash = hash;
  }

  addNode(id, weight = 1) {
    if (typeof id !== "string" || id.length === 0)
      throw new TypeError("node id must be a non-empty string");
    if (weight <= 0) throw new RangeError("weight must be > 0");
    if (this.#nodes.has(id)) return this;
    this.#nodes.set(id, {
      weight,
      points: Math.max(1, Math.round(this.replicas * weight)),
    });
    this.#rebuild();
    return this;
  }

  removeNode(id) {
    if (this.#nodes.delete(id)) this.#rebuild();
    return this;
  }

  #rebuild() {
    const pairs = [];
    for (const [id, { points }] of this.#nodes) {
      for (let i = 0; i < points; i++)
        pairs.push([this.hash(`${id}#${i}`), id]);
    }
    pairs.sort((a, b) => a[0] - b[0]);
    this.#positions = Uint32Array.from(pairs, (p) => p[0]);
    this.#owners = pairs.map((p) => p[1]);
  }

  /** Index of the first ring point >= h, wrapping to 0. */
  #successor(h) {
    const pos = this.#positions;
    let lo = 0,
      hi = pos.length;
    while (lo < hi) {
      const mid = (lo + hi) >>> 1;
      if (pos[mid] < h) lo = mid + 1;
      else hi = mid;
    }
    return lo === pos.length ? 0 : lo;
  }

  get(key) {
    if (this.#positions.length === 0) return null;
    return this.#owners[this.#successor(this.hash(key))];
  }
}

Drei Gestaltungsentscheidungen verdienen besondere Aufmerksamkeit.

Parallel arrays anstelle eines Arrays von Objekten. Das Speichern der Positionen in einem Uint32Array ermöglicht es dem binären Suchalgorithmus, auf kompaktem, zusammenhängendem Speicher zu arbeiten, der im CPU-Cache verfügbar bleibt. Bei 8 Knoten und 160 Replikaten gibt es 1.280 Punkte, etwa 5 KB, und eine Suche erfordert rund 11 Vergleiche.

Die Umkehrung in lo === pos.length ? 0 : lo. Ein Schlüssel, dessen Hash-Wert über den letzten Punkt hinausgeht, gehört zum ersten Knoten im Kreis. Das Auslassen dieser Bedingung ist der häufigste Fehler bei selbstgebauten Ringstrukturen: Das Ergebnis ist für fast jeden Schlüssel korrekt, doch unerklärlicherweise falsch für die wenigen Schlüssel am oberen Ende des Bereichs.

Aufrüsten bei Änderung der Zugehörigkeit, nicht bei Suche. Änderungen der Zugehörigkeit kommen selten vor, während Suchvorgänge zehntausende Male pro Sekunde stattfinden; daher kostet das gelegentliche Sortieren von 1.280 Einträgen nichts Wesentliches.

Replikation, das heißt die Suche nach den nächsten Eigentümern für einen Schlüssel, erfolgt durch einen Uhrzeigersinn-Weg, bei dem unterschiedliche physische Knoten gesammelt werden. Das Wort „unterschiedlich“ ist wichtig, denn benachbarte Punkte auf dem Ring gehören oft zum selben Server:

getReplicas(key, count = 1) {
  const n = this.#positions.length;
  if (n === 0) return [];
  const wanted = Math.min(count, this.#nodes.size);
  const out = [];
  const start = this.#successor(this.hash(key));
  for (let step = 0; step < n && out.length < wanted; step++) {
    const owner = this.#owners[(start + step) % n];
    if (!out.includes(owner)) out.push(owner);
  }
  return out;
}

Beachten Sie, dass wanted auf die Anzahl der physischen Knoten begrenzt ist, sodass eine Anfrage nach mehr Replikaten als Servern nicht endlos weiterlaufen kann – der Weg stoppt in jedem Fall nach einer vollen Runde.

Auswahl der Hash-Funktion mit Bedacht

Viele Tutorials überspringen diesen Teil, obwohl er das aufschlussreichste Ergebnis der gesamten Übung liefert.

Die meisten Ringimplementierungen verwenden standardmäßig MD5. Es funktioniert zwar, ist aber langsam: Ein Ring mit dieser Methode erreichte lediglich 423.000 Abfragen pro Sekunde, und Profilanalysen zeigten, dass fast die gesamte Zeit in der Verarbeitung durch MD5 verbracht wurde.

Durch Ersetzung mit FNV-1a, einem schnellen nicht-kryptografischen Hash-Funktion, stieg die Durchsatzleistung um etwa das 15-fache. Der Ausgleich hingegen brach zusammen: die Standardabweichung der Last pro Node stieg von 10,4 % auf 30,7 %.

Durch Auswertung der berechneten Positionen einiger virtueller Nodes wird die Ursache offensichtlich:

cache-01#0 → 4037809751
cache-01#1 → 4021032132
cache-01#2 → 4071364989
cache-01#3 → 4054587370
cache-01#4 → 3970699275

Sie alle liegen in einem schmalen Bereich um 4,0 Milliarden. FNV-1a weist ein schwaches Avalanche-Ergebnisverhalten auf, was bedeutet, dass ähnliche Eingaben ähnliche Ausgaben ergeben. Die Namen der virtuellen Nodes unterscheiden sich nur durch einen Suffix, wodurch die 160 Punkte anstelle dessen, verteilt über den Kreis zu liegen, in einem engen Cluster pro Node angesammelt werden. Die Optimierung für Geschwindigkeit brachte problematisch Problem 3 wieder zurück.

Die Lösung besteht darin, die Ausgabe von FNV durch einen Bit-Mixing-Finalisator zu leiten, der letzte Schritt von MurmurHash3:

function fnv1a(str) {
  let h = 0x811c9dc5;
  for (let i = 0; i < str.length; i++) {
    h ^= str.charCodeAt(i);
    h = Math.imul(h, 0x01000193);
  }
  return h >>> 0;
}

// Scrambles the bits so near-identical inputs land far apart.
function fmix32(h) {
  h ^= h >>> 16;
  h = Math.imul(h, 0x85ebca6b);
  h ^= h >>> 13;
  h = Math.imul(h, 0xc2b2ae35);
  h ^= h >>> 16;
  return h >>> 0;
}

export const defaultHash = (str) => fmix32(fnv1a(str));

fmix32 wechselt zwischen Verschiebungen, XOR-Operationen und Multiplikationen, sodass eine Veränderung in einem Eingabebit sich auf alle Ausgabebits ausbreitet. Math.imul führt eine echte 32-Bit-Integer-Multiplikation durch, und >>>> 0 wandelt das Ergebnis wieder in eine ungesehene 32-Bit-Zahl um, die in einen Uint32Array passt. Durch etwa zehn zusätzliche Operationen liefert die kombinierte Hash-Funktion ein besseres Gleichgewicht als MD5 und ist dabei etwa 13-mal schneller.

Die daraus abgeleitete Lektion gilt weit über das Hashing hinaus: Wenn Sie ein Komponente durch eine schnellere ersetzen, müssen Sie die Eigenschaft messen, die Sie nicht optimiert haben. FNV-1a ist ein durchaus brauchbarer Hash; er eignet sich einfach nicht für diese Aufgabe, und in seiner Beschreibung wird Sie nichts davor warnen.

Umwandlung des Rings in einen Cache-Client

Ein Ring für sich allein ist kein Cache-Client. Produktionscode muss mit Fehlern umgehen können, und der Ring bietet eine klare Strategie: Übergang zum nächsten Knoten im Uhrzeigersinn.

Der untenstehende Wrapper nimmt ein Verzeichnis benannter Clients, erstellt daraus den Ring und versucht bei jeder get-Aktion nacheinander die ersten failoverDepth Eigentümer. Ein Knoten, der Fehler auslöst, wird für downtimeMs deaktiviert und aus dem Ring entfernt, bevor er nach Ablauf der Strafzeit wieder hinzugefügt wird.

export class ShardedCache {
  #ring; #clients; #down = new Map();

  constructor(clients, { replicas = 160, failoverDepth = 2, downtimeMs = 10_000 } = {}) {
    this.#clients = new Map(Object.entries(clients));
    this.#ring = new ConsistentHashRing({ replicas });
    for (const id of this.#clients.keys()) this.#ring.addNode(id);
    this.failoverDepth = failoverDepth;
    this.downtimeMs = downtimeMs;
  }

  #markDown(id) {
    this.#down.set(id, Date.now() + this.downtimeMs);
    this.#ring.removeNode(id);
  }

  #reviveExpired() {
    const now = Date.now();
    for (const [id, until] of this.#down) {
      if (now >= until) { this.#down.delete(id); this.#ring.addNode(id); }
    }
  }

  async get(key) {
    this.#reviveExpired();
    for (const id of this.#ring.getReplicas(key, this.failoverDepth)) {
      try {
        return { value: await this.#clients.get(id).get(key), node: id };
      } catch {
        this.#markDown(id);
      }
    }
    return { value: null, node: null, allDown: true };
  }
}

Die Ausführung gegen vier simulierte Knoten, die jeweils 10.000 Schlüssel enthielten, wobei anschließend einer von ihnen abgeschaltet wurde, ergab Folgendes:

Keys per node:  cache-01 2213 | cache-02 2462 | cache-03 2923 | cache-04 2402

Killing cache-02...
  served from cache: 7538
  cache misses:      2462
  hard failures:     0
  healthy nodes:     cache-01, cache-03, cache-04

  24.6% of traffic became a miss.

Diese Zahl ist Ihr Kapazitätsplan. In einer Flotte mit vier Knoten führt ein Ausfall dazu, dass 24,6 % mehr Lesevorgänge auf die Datenbank fallen, während eine Flotte mit acht Knoten diesen Wert auf 12,5 % begrenzt. Wenn die Datenbank den Betrag 1/N der gleichzeitigen Lesebelastung nicht bewältigen kann, liegt das eigentliche Problem in der Kapazität der Datenbank und nicht im Caching – wobei konsistente Hashing-Methoden dieses Problem von fatal zu erkennbar gemacht haben. Beachten Sie außerdem, dass es keine harten Ausfälle gab: Anfragen nach Schlüsseln des ausgefallenen Knotens wurden an den nächsten Knoten weitergeleitet und zu gewöhnlichen Fehlern.

Produktionsprobleme, die der Algorithmus nicht abdeckt

Mehrere Probleme liegen außerhalb des Algorithmus und schaden Ihnen in jedem Fall.

Abweichungen in der Ringversion

Der Ring ist nur dann nützlich, wenn jeder Client dasselbe Ergebnis berechnet. Führen Sie die Änderung der Mitgliedschaft schrittweise ein – für einige Minuten sieht die eine Hälfte der Flotte 8 Knoten, während die andere Hälfte 9 sieht. Die beiden Gruppen sind sich bei etwa 11 % der Schlüssel uneinig. Für einen Cache bedeutet das einen kurzen Rückgang der Trefferquote; für alles, was Schreibvorgänge zulässt, bedeutet es unterschiedliche Daten. Versionieren Sie den Mitgliedschaftssatz, verteilen Sie ihn über einen einzigen Kanal und geben Sie die Version in Ihren Metriken an, damit Sie Abweichungen beobachten können statt nur zu vermuten.

Zu voreiliges Ausschließen von Knoten

Nehmen Sie keinen Knoten nach einem Timeout entfernt. Ein schwankender Knoten, der ständig wieder hinzukommt und geht, verursacht einen „Churn-Sturm“, da jede Übertragung 1/N der Schlüssel verschiebt. Erfordern Sie mehrere aufeinanderfolgende Fehler innerhalb eines Zeitfensters, bevor der Knoten entfernt wird, und nehmen Sie ihn vorsichtig wieder auf. Das obige Beispiel verwendet eine feste Strafe von zehn Sekunden; Produktionscode sollte exponentielles Backoff sowie eine Gesundheitsprüfung vor der Wiederaufnahme eines Knotens verwenden.

Hot Keys sind nicht ausbalanciert

Konsistentes Hashing verteilt Schlüssel gleichmäßig, sagt aber nichts über Anfragen. Wenn ein einzelnes Produkt plötzlich extrem beliebt wird, befindet sich dessen Eintrag auf einem Server, und dieser Server überhitzt – obwohl der Ring seine Aufgabe erfüllt. Es gibt zwei Lösungsansätze: einen kleinen In-Process-Cache vor dem Ring für die am häufigsten genutzten Schlüssel oder die von Google-Forschern entwickelte Variante des konsistenten Hashings mit begrenzten Lasten, die festlegt, wie viel jedes Node übernehmen darf, und den Überschuss im Uhrzeigersinn weiterleitet.

Remapping ist keine Migration

Alles, was oben erwähnt wurde, geht davon aus, dass das Verlieren einer Schlüssel nur zu einem Cache-Miss führt. Wenn der Ring dauerhafte Daten routet, bedeutet „11 % der Schlüssel wurden verschoben“, dass 11 % der Daten physisch auf den neuen Knoten kopiert werden müssen, bevor sie dort gelesen werden können – in der Regel erfolgen dabei Lese- oder Schreibvorgänge an beiden Standorten, bis die Kopie abgeschlossen ist. Der Ring bestimmt lediglich was verschoben werden muss; er übernimmt selbst keine Aktivität zur Verschiebung. Genau deshalb verlassen sich Systeme wie Vitess auf einen Koordinator: Sie müssen die Migration steuern, nicht nur berechnen.

Die Anzahl der Replicas ist im Grunde dauerhaft

Wenn man replicas von 160 auf 500 ändert, verschiebt sich jede Ringposition und fast alle Schlüssel werden neu angeordnet – was genauso störend ist wie eine Änderung des Modulo-Werts. Betrachten Sie dies als einmalige Designentscheidung, die bereits vor dem Einsatz von Live-Daten getroffen wird; falls Sie unsicher sind, wählen Sie den höheren Wert.

Wann ein Ring das falsche Werkzeug ist

Ein Teil der ingenieurtechnischen Urteilsfähigkeit besteht darin, zu erkennen, wann die beeindruckende Lösung die falsche ist.

  • Verwenden Sie Rendezvous-Hashing. Es erfordert weniger Code, bietet besseres Gleichgewicht und es muss keine Anzahl an Kopien eingestellt werden. Der einzige Vorteil des Rings ist die Suchzeit von O(log N), was bei dieser Größe eher akademisch ist.
  • Austauschbare Behälter, bei denen nur die Anzahl sich ändert: Verwenden Sie Jump-Hashing – es benötigt zehn Zeilen Code und keinen Speicher.
  • Dauerhafte Daten, die kontrolliert bewegt werden müssen: Verwenden Sie einen Koordinator. Nur eine explizite Karte ermöglicht es, einen einzelnen Shard zu verschieben und gleichzeitig dessen Auswirkungen zu überwachen, was Hashing nicht leisten kann.
  • N, das tatsächlich niemals ändert: % N reicht aus. Bauen Sie keine Mechanismen für eine Veränderung, die nicht eintreten wird.

Ein Ring erlangt seine Bedeutung, wenn die Knoten benannt sind, heterogen, zahlreich und anfällig für Ausfälle. Das beschreibt eine Cache-Flotte nahezu perfekt, weshalb so viele verteilte Caches auf diesem Prinzip basieren.

Kernpunkte

  • Die Grundidee ist einfach: Man platziert Schlüssel und Server im selben Adressraum, sodass ein Mitgliedschaftswechsel nur einen bestimmten Bereich beeinträchtigt und nicht alles.
  • Durch Modulo-Routing wird jedes Skalierungsereignis sowie jeder Knotenausfall zu einem nahezu vollständigen Cache-Löschen, wobei der Schaden mit der Größe des Clusters zunimmt.
  • Rendezvous-Hashing ist oft die bessere Wahl für kleine Flotten; man greift auf den Ring zurück, wenn Größe, Gewichtung sowie willkürliche Entfernung von Knoten eine Rolle spielen.
  • Ohne virtuelle Knoten kann ein Server mehrere Male seinen fairen Anteil tragen. Wählen Sie die Anzahl der Replicas im Voraus aus, denn ein späterer Wechsel verschiebt alles neu.
  • Ein schneller Hash mit schwacher Lawinenwirkung kann heimlich ein Ungleichgewicht wiederherstellen; messen Sie stets die Verteilung und nicht nur die Geschwindigkeit.
  • Durch den Verlust eines Knotens gelangen etwa 1/N der Leseanfragen an die Datenbank. Planen Sie die Größe der Datenbank entsprechend, handhaben Sie Änderungen in der Mitgliedschaft vorsichtig, entfernen Sie Knoten zurückhaltend und kümmern Sie sich separat um häufig genutzte Schlüssel.
  • Der Ring selbst besteht nur aus einigen Dutzend Zeilen Code. Die Ingenieursarbeit, die seine Zuverlässigkeit in der Produktion gewährleistet, ist alles, was ihn umgibt.

    Verwandte Artikel