Zum Inhalt springen
Deutsch

Scatter/Gather

Router.scatterGatherFirstCompleted sendet jede Nachricht an alle seine Routees und antwortet dem Aufrufer mit der ersten Antwort. Die Verlierer laufen weiter; ihre Antworten werden verworfen.

import { ActorSystem, Router, ScatterGatherOptions, Actor } from 'actor-ts';
class Replica extends Actor<string> {
override onReceive(key: string): void {
this.sender.toNullable()?.tell(`${this.self.path.name}:${key}`);
}
}
const system = ActorSystem.create('demo');
const hedgeOptions = ScatterGatherOptions.create().withTimeoutMs(250);
const replicas = system.spawn(
Router.scatterGatherFirstCompleted(3, Replica, hedgeOptions),
'replicas',
);
const value = await replicas.ask<string>('user-42');

Drei Replicas werden parallel gefragt; wer zuerst antwortet, bestimmt den Wert.

Das ist das Hedged-Request-Pattern. Es erhöht nicht den Durchsatz — es senkt den Tail. N Routees machen jeweils die ganze Arbeit, du zahlst also das N-Fache und bekommst das Minimum aus N Latenzen statt einer zufälligen.

Dieser Handel lohnt, wenn der langsame Tail deutlich langsamer ist als der Median:

  • Redundante Replica-Reads — drei Caches, die erste Antwort gewinnt.
  • Multi-Source-Lookup — mehrere Indizes, nimm wer antwortet.
  • Spekulative Ausführung — derselbe Job auf N Workern, das erste Ergebnis zählt.

Sonst ist es schlicht N-fache Verschwendung. Um unterschiedliche Arbeit über einen Pool zu verteilen, nimm Router.roundRobin oder Router.smallestMailbox.

Der Router muss die Routee-Antworten abfangen, um einen Gewinner zu küren — er braucht also ein Ziel für die gewinnende Antwort. Zwei Formen funktionieren:

// 1. ask — der Normalfall.
const value = await replicas.ask<string>('user-42');
// 2. tell mit explizitem Sender — innerhalb eines Actors.
replicas.tell('user-42', this.self);

Ein nacktes replicas.tell('user-42') hat kein Antwortziel. Es wird nichts gestreut, die Nachricht wird verworfen, und der Router loggt eine Warnung — N-fach auszufächern für eine Antwort, die niemand entgegennimmt, ist genau die Verschwendung, für die das Pattern etwas kaufen soll. Es wirft nicht: Ein Wurf würde den Router über die Supervision fehlschlagen lassen, und der Restart risse jeden anderen laufenden Scatter mit.

Die Antwort wird dem Routee zugeschrieben, der sie produziert hat, nicht dem Router. Ein Aufrufer, der this.sender liest, sieht das, was er auch gesehen hätte, wenn er diesen Routee direkt gefragt hätte — und für Hedged Reads beantwortet das die Frage, die das Pattern aufwirft: welche Replica gewonnen hat.

Jeder Fehlschlag rejected mit einem AggregateError. Sein errors-Array hält einen Fehler pro Routee in Scatter-Reihenfolge, und die Message benennt die Ursache:

Situationmessage sagterrors enthält
Niemand antwortete rechtzeitignone of N routees replied within …msN AskTimeoutError
Jeder Routee schlug fehlall N routees failedwas jeder Routee zurückgab
Alle Routees sind gestopptno routees left to scatter toleer
Router stoppte mitten im Scatterstopped while the scatter was still openleer
try {
await replicas.ask<string>('user-42');
} catch (e) {
const aggregate = e as AggregateError;
console.error(aggregate.message, aggregate.errors);
}

Ein Typ für jeden Fehlschlag, mit Absicht: Ein Aufrufer, der auf den Fehlertyp verzweigen müsste, um zu erfahren, wie viele Routees fehlschlugen, erführe nichts, was das errors-Array nicht schon trägt. Um gezielt auf “zu langsam” zu prüfen, sieh dir die Mitglieder an:

const allTimedOut = aggregate.errors.every((e) => e instanceof AskTimeoutError);

AggregateError übersteht die Cluster-Leitung — der eingebaute Serializer kodiert seine Mitglieds-Fehler mit — ein Scatter hinter einem ClusterRouter schlägt auf dem aufrufenden Knoten also genauso fehl.

Ein Routee, der wirft statt zu antworten, produziert hier keinen Fehler direkt: Sein Supervisor restartet ihn, und der Ask des Routers für diesen Routee läuft einfach in den Timeout. Um einen Fehler als Fehler zu melden, antworte mit einem Error — das rejected den Ask sofort, und der Scatter kann zum nächstschnellsten Routee weiterziehen, statt die Uhr abzuwarten.

withTimeoutMs ist die Deadline für einen ganzen Scatter — Akka nennt denselben Knopf within. Er wird zum Timeout auf dem ask jedes Routees, ist also pro Routee durchgesetzt, und der Scatter schlägt fehl, sobald der letzte abgelaufen ist. Default: 4_500 — knapp unter den 5_000 von ActorRef.ask.

Der Abstand ist der Punkt. Der Router kann erst berichten, welche Routees gescheitert sind, wenn seine eigene Deadline verstrichen ist und er deren Fehler eingesammelt hat. Wäre sein Budget so groß wie das des Aufrufers, hätte dessen ask längst mit AskTimeoutError aufgegeben. Was auch immer du hier setzt: Halte das ask-Timeout des Aufrufers darüber — sonst bekommst du dein eigenes Timeout statt des AggregateError, für den es diesen Router gibt.

const hedgeOptions = ScatterGatherOptions.create().withTimeoutMs(250);

Ein einfaches Objekt geht auch — { timeoutMs: 250 } — und beide durchlaufen dieselbe Validierung: Ein gesetztes timeoutMs muss eine positive endliche Zahl sein, geprüft am Router.scatterGatherFirstCompleted(...)-Aufruf, damit der Stack noch auf dich zeigt.

Die Dimensionierung ist der halbe Punkt. Setze ihn unter das Ask-Timeout des Aufrufers: Eine hängende Replica soll einen Bruchteil des Budgets des Aufrufers kosten, nicht alles — und der Aufrufer muss noch zuhören, wenn der Router meldet, welche Routees fehlschlugen.

Beachte: Laufende Routee-Asks werden nicht abgebrochen, wenn ein Scatter früh endet — ihre Reply-Refs laufen einfach auf ihren eigenen Timern ab. Ein langes timeoutMs hält also ein paar Timer über die Antwort hinaus am Leben.

postStop lässt jeden noch offenen Scatter fehlschlagen, statt seinen Aufrufer erst dann merken zu lassen, dass der Router weg ist, wenn der Ask irgendwann in den Timeout läuft. Die Routees sind zu diesem Zeitpunkt schon gestoppt, ihre Antworten können also nicht mehr ankommen; timeoutMs abzuwarten machte den Shutdown nur so langsam wie die längste konfigurierte Deadline.

Ein Restart macht dasselbe, aus demselben Grund — preRestart fällt auf postStop zurück, und die Routees, die ein restartender Router gefragt hat, werden gerade abgerissen und neu gespawnt.

Der Fan-Out wird abgefeuert, ohne ihn im Handler zu awaiten, und die Continuation antwortet außerhalb des Mailbox-Turns. Das ist kein Detail: Die Runtime awaitet den Handler eines Actors, bevor sie die nächste Nachricht entnimmt — ein async onReceive um Promise.any(...) würde also die gesamte Mailbox des Routers für die Dauer eines Scatters parken. Ein einziger hängender Routee hielte dann jeden anderen Aufrufer für das volle timeoutMs auf, und gleichzeitige Asks liefen nacheinander statt parallel.

So wie es geschrieben ist, überlappen sich Hunderte Scatters frei.

Zwei Familien, emittiert wenn die Metrics-Extension aktiviert ist:

MetrikArtLabels
router_scatter_gather_resolved_totalCounteroutcome
router_scatter_gather_latency_secondsHistogram

outcome ist eines von first, timeout, all-failed, stopped, no-reply-target. timeout und all-failed sind getrennt, weil sie unterschiedliche Reaktionen verlangen: Das erste sagt, die Routees sind zu langsam für das konfigurierte Budget, das zweite, dass sie kaputt sind.

Das Latenz-Histogram beobachtet die Zeit vom Streuen bis zur Antwort an den Aufrufer, ist also die Latenz des Gewinners — die Verteilung, die du bewegen willst. Für no-reply-target, wo kein Scatter lief, wird nichts beobachtet.

Keine der beiden Familien trägt einen Router-Pfad. Bei mehreren Scatter/Gather-Routern in einem System aggregieren die Serien; der Actor-Pfad in den Log-Zeilen unterscheidet sie.

Wie der restliche lokale Router ist auch dieser ein Pool — er spawnt und besitzt seine Routees. Es gibt keine Group-Variante, die an bereits existierende Actors streut; siehe Pool vs Group dafür, warum der lokale Router pool-förmig ist, und Cluster-Router für das group-förmige Cluster-Äquivalent.

  • Router — die fünf One-Shot-Strategie- Factories, neben denen diese steht.
  • Strategien — wie die anderen Router einen Routee auswählen.
  • Circuit Breaker — das andere Tail-Latenz-Werkzeug: eine fehlschlagende Dependency gar nicht mehr aufrufen, statt sie N-mal aufzurufen.