Sharding, Hot Partitions und CAP

Track Konzepte · Datenbanken verteilen · ca. 40 Min.

Was du aus Teil a brauchst

Du weißt, was Replikation ist (ein Leader, Follower mit Replication Lag) und dass Kopien nie ganz gleichzeitig aktuell sind. Das ist die Grundlage für die Frage, was du bei einer Netzwerkpartition aufgibst. Siehe Teil 1.

Worum es geht

Replikation hilft bei Ausfall und Leselast, aber nicht, wenn die Daten nicht mehr auf eine Maschine passen oder eine Maschine alle Schreibvorgänge nicht schafft. Dann teilst du die Daten auf (Sharding, auch Partitionierung, partitioning). Das kostet: Manche Abfragen werden teuer, und ein schlechter Schlüssel erzeugt eine Hot Partition. Dazu kommt das CAP-Theorem, das beschreibt, was du bei einer Netzwerkpartition aufgibst.

Am Ende von Teil 2 kannst du einen Sharding-Schlüssel nach Last bewerten, ausrechnen, wie viele Keys beim Rebalancing wandern, und erklären, was CAP wirklich sagt. Wie in Teil 1 simulierst du alles mit Code ohne echte Server.

Von JS/TS her gedacht

Du kennst Verteilte Datenbank
Eintrag landet im “falschen” Cache-Knoten Sharding-Schlüssel falsch gewählt
array.length % n als Bucket Hash-Sharding mit % n
Ein Mikroservice, den alle ansprechen und der umfällt Hot Partition

Konzept

Schritt 1: Sharding

Beim Sharding (Partitionierung) hat jeder Shard nur einen Teil der Daten. Die Wahl des Sharding-Schlüssels (shard key) entscheidet, wie gleichmäßig Last und Daten verteilt sind. Zwei Standard-Strategien:

  • Hash-Sharding: shard = hash(key) % n. Gleichmäßig, aber Bereichsabfragen (“alle Bestellungen von März”) treffen alle Shards.
  • Range-Sharding: Bereiche des Schlüssels (z. B. nach Datum). Bereichsabfragen sind billig, aber es entstehen leicht Hot Partitions.

Eine Hot Partition ist ein Shard, der viel mehr Last bekommt als die anderen. Dann hilft das Hinzufügen von Maschinen nicht, denn der heiße Shard bleibt ein einzelner Engpass. Typische Ursachen: Zeit als Schlüssel (alle neuen Schreibvorgänge landen im jüngsten Bereich), ein Großkunde als Schlüssel (alle seine Daten in einem Shard), ein sehr beliebter Key.

Beide Strategien im Vergleich. Das Hashing nutzt zlib.crc32, weil Pythons hash() für Strings bei jedem Start anders rechnet:

Der Datumsschlüssel sieht über den Verlauf gleichmäßig aus (7, 8, 8, 7), aber die Last von heute liegt zu 100 % in einem Shard. Merke: Gleichmäßige Daten sind nicht gleichmäßige Last.

Zwei weitere Punkte.

Rebalancing: Wechselst du von % 4 auf % 5, ändert sich für fast alle Keys der Shard. Hier die Rechnung über 10 000 Keys:

Es wandern 8046 von 10 000, also etwa 80 %. Das ist der Grund für Consistent Hashing (Keys und Knoten auf einem Ring, beim Hinzufügen wandert nur ein kleiner Teil). Dazu mehr in einer späteren Lektion.

Cross-Shard-Queries: Eine Abfrage, die den Sharding-Schlüssel nicht enthält, muss alle Shards fragen und die Ergebnisse zusammenführen (Scatter-Gather). Das ist langsamer und fällt aus, sobald ein Shard nicht antwortet. Joins und Transaktionen über Shards hinweg sind noch teurer.

Schritt 2: CAP richtig gelesen

Das CAP-Theorem wird oft als “wähle zwei von drei” verkürzt. Das führt in die Irre. Genauer:

  • Consistency: Jede Lesung sieht den letzten Schreibvorgang (hier im Sinn von Linearisierbarkeit, linearizability: Das System verhält sich, als gäbe es nur eine einzige Kopie der Daten).
  • Availability: Jede Anfrage an einen funktionierenden Knoten bekommt eine Antwort.
  • Partition Tolerance: Das System arbeitet weiter, obwohl das Netzwerk Knoten voneinander trennt.

In einem verteilten System kommt eine Netzwerkpartition vor, du kannst sie nicht abwählen. Die eigentliche Frage ist daher: Was gibst du auf, solange die Partition dauert?

  • CP: Du gibst Verfügbarkeit auf. Die Seite ohne Mehrheit lehnt Anfragen ab (Fehler), dafür gibt es keine veralteten oder widersprüchlichen Daten.
  • AP: Du gibst Konsistenz auf. Beide Seiten antworten weiter, können aber auseinanderlaufen und müssen später zusammengeführt werden (hier: Last Write Wins, der spätere Zeitstempel gewinnt und die andere Änderung geht verloren).

Ohne Partition kannst du beides haben. Hier der Simulator mit zwei Knoten: A ist die Mehrheit, B die Minderheit.

Szenario: v1 wird geschrieben, dann fällt das Netz zwischen A und B aus, danach schreibt jemand v2 an A, und jemand liest bei B. Zum Schluss heilt das Netz.

Ergebnis: CP liefert (True, 'FEHLER', 'v2'). B antwortet während der Partition nicht, danach stimmt alles. AP liefert (True, 'v1', 'v2'). B antwortet immer, aber während der Partition mit dem veralteten Wert v1. Beides ist eine bewusste Entscheidung, keine Panne. Wie sich das bei Schreibvorgängen auf beiden Seiten auswirkt, rechnest du in Übung 3 aus.

Falle

  • Den Sharding-Schlüssel nach Datenmenge statt nach Last wählen. Ein Zeitstempel verteilt Historie gleichmäßig und die Gegenwart auf einen Shard.

  • CAP als dauerhafte Wahl lesen. Die Wahl betrifft nur die Zeit der Partition. Im Normalbetrieb gibt es andere Abwägungen (z. B. Latenz gegen Konsistenz).

  • Modulo-Sharding und dann Shards hinzufügen. % n auf % n+1 bewegt die meisten Keys. Plane Rebalancing ein (Consistent Hashing, eine spätere Lektion).

Übungen

Übung 1: Sharding-Schlüssel wählen, Hot Partition berechnen (ca. 12 Min.)

Ein Webshop will seine Bestellungen auf 4 Shards verteilen. Gemessen wird die Schreiblast von heute. Du hast die Liste events mit 1000 Schreibvorgängen als Tupel (bestell_id, kunde_id, tag). Alle sind von Tag 30. Kunde 7 ist ein Großkunde. Das Dictionary STRATEGIEN enthält drei Funktionen f(event, n), die den Shard (0 bis n-1) liefern: kunde_hash, bestell_hash, datum_bereich.

  1. Schreibe lastverteilung(events, shard_fn, n): Liste mit n Zahlen, wie viele Events pro Shard landen.
  2. Setze bester auf den Namen der Strategie, bei der der am stärksten belastete Shard am wenigsten Last trägt. Probiere es aus: Lass dir die Verteilung der drei Strategien ausgeben, bevor du dich entscheidest. (Abfragen nach Kunde ignorierst du in dieser Aufgabe. Die Folgen davon stehen im Lösungstext.)

lastverteilung ist eine Liste [0] * n, die du pro Event am Index shard_fn(event, n) erhöhst. Für bester zählt der Maximalwert der Liste, nicht der Durchschnitt (der ist immer 250). Wo sitzt der Großkunde?

def lastverteilung(events, shard_fn, n):
    zaehler = [0] * n
    for e in events:
        zaehler[shard_fn(e, n)] += 1
    return zaehler

bester = "bestell_hash"

lastverteilung, bester

Jeder Schlüssel hat einen Preis. kunde_hash verteilt die Kunden, aber der Großkunde (35 % der Events) sitzt komplett in einem Shard: eine Hot Partition. datum_bereich legt alle heutigen Schreibvorgänge in den jüngsten Bereich: eine noch heißere Hot Partition. bestell_hash verteilt die Schreiblast gleichmäßig. Der Preis: “Alle Bestellungen von Kunde X” liegt über alle Shards verstreut und braucht Scatter-Gather. Wenn diese Abfrage wichtiger ist als gleichmäßiges Schreiben, kann kunde_hash mit einem Plan für den Großkunden trotzdem die bessere Wahl sein. Das ist ein Trade-off, und du entscheidest ihn anhand deiner Zugriffsmuster.

Übung 2: Rebalancing, wie viele Keys wandern? (ca. 8 Min.)

Ein Dienst verteilt Sitzungs-IDs mit hash_shard(sitzung, n) auf n Shards (dieselbe Funktion wie im Text, der Setup-Code liegt unsichtbar bereit). Jetzt soll ein Shard dazukommen, und du willst wissen, wie viel Daten umziehen muss.

  1. Schreibe wanderer(keys, n_alt, n_neu): die Anzahl der Keys, deren Shard sich beim Wechsel von n_alt auf n_neu ändert.
  2. Der Dienst läuft auf 4 Shards und soll auf 5, 6 oder 8 wachsen. Probiere es mit der Liste sitzungen (6000 IDs) aus und setze beste_n auf die Shard-Zahl, bei der am wenigsten Keys wandern.

Zähle die Keys, bei denen hash_shard(k, n_alt) und hash_shard(k, n_neu) verschieden sind. Für beste_n vergleichst du die drei Zahlen, die deine Funktion für 4 auf 5, 6 und 8 liefert.

def wanderer(keys, n_alt, n_neu):
    return sum(hash_shard(k, n_alt) != hash_shard(k, n_neu) for k in keys)

beste_n = 8

wanderer, beste_n

Bei 4 auf 8 bleibt jeder Key, dessen Hash modulo 8 kleiner als 4 ist, auf seinem Shard: etwa die Hälfte wandert. Bei 5 und 6 bleiben viel weniger Keys, wo sie sind. Auch die beste Zahl bewegt noch rund die Hälfte aller Daten. Genau das löst Consistent Hashing (später).

Übung 3: CAP vorhersagen (ca. 10 Min.)

Der Simulator ZweiKnoten aus Schritt 2 liegt bereit (A ist die Mehrheit, B die Minderheit, lies gibt bei Nichterreichbarkeit "FEHLER" zurück, schreibe gibt True oder False zurück). Dieses Szenario läuft nacheinander ab:

c = ZweiKnoten(modus)
c.schreibe("A", "stand", "v1")      # Netz in Ordnung
c.trennen()                         # Netzwerkpartition beginnt
r1 = c.schreibe("A", "stand", "v2") # Client 1 schreibt bei A
r2 = c.schreibe("B", "stand", "v3") # Client 2 schreibt bei B
a  = c.lies("A", "stand")
b  = c.lies("B", "stand")
c.heilen()                          # Netz kommt zurück, Last Write Wins
nachher = c.lies("B", "stand")

Trage für CP und für AP je das Tupel (r1, r2, a, b, nachher) ein. Der Check führt das Szenario selbst aus.

CP lehnt auf der Minderheitsseite ab, AP antwortet überall. Bei AP: Wessen Zeitstempel ist der spätere, v2 oder v3? Und was passiert beim Heilen mit dem früheren Wert?

antwort = {
    "CP": (True, False, "v2", "FEHLER", "v2"),
    "AP": (True, True, "v2", "v3", "v3"),
}
antwort

CP: B lehnt Schreiben und Lesen ab. Es gibt nur eine Wahrheit (v2), aber ein Client bekam Fehler. AP: Beide Seiten nehmen Schreibvorgänge an, A liest v2, B liest v3, die Antworten widersprechen sich. Beim Heilen gewinnt der spätere Zeitstempel, v3. Die Änderung v2 ist verloren, und beide Clients dachten, ihr Schreiben sei erfolgreich. Genau das ist der Preis von AP, und deshalb braucht AP eine Strategie zur Konfliktauflösung.

Merksatz

Sharding und CAP kaufen Skalierung und Verfügbarkeit mit einfachen Abfragen und Konsistenz. Du wählst den Sharding-Schlüssel nach der Last (nicht nach der Datenmenge), und bei einer Netzwerkpartition entscheidest du bewusst: Fehler anzeigen (CP) oder veraltete Antworten riskieren (AP).

Prüfstein

Dein Webshop hat eine Hot Partition, weil ein Großkunde 35 % der Schreibvorgänge erzeugt. Warum hilft ein weiterer Shard nicht, welche Schlüsselwahl verteilt die Last besser, und was bezahlst du dafür bei Abfragen?

Zurück zu Teil 1.


Quelle: quellen/konzeptuebersicht-software-grundlagen.md, Abschnitt “5. Datenbanken” (Verteilung: Sharding, CAP-Theorem, Eventual Consistency); quellen/konzeptuebersicht-software-fortgeschritten.docx, Schicht 5 “Partitionierung” (Sharding-Strategien, Hot Partitions, Rebalancing, Cross-Shard-Queries). Beide Quellen sind Stichwortlisten. Über die Quelle hinaus (allgemeines Fachwissen): die Beschreibung von Hash- und Range-Sharding sowie der Rechnung zum Rebalancing, die Präzisierung des CAP-Theorems (Wahl nur während einer Partition) und Last Write Wins. Consistent Hashing wird nur erwähnt (eigene spätere Lektion). Alle Zahlen und Ausgaben im Text stammen aus ausgeführtem Code (Python 3). Konkrete Standardwerte oder Verhalten einzelner Datenbankprodukte werden nicht genannt (bitte prüfen, falls ergänzt).