29. pepsi-dispatch¶
Der einzige langlebige Koordinator, der Nachrichten durch die Stages treibt.
29.1. Rolle¶
pepsi-dispatch ist der einzige Prozess, der Stage-Programme startet. Betreiben Sie genau einen pro System, als pepsi-dispatch serve (das einzige Unterkommando; die globalen Schalter -c/-L/-v stehen davor). Er beansprucht neu eingereihte Datensätze und reicht jeden an einen untätigen Worker-Prozess der Stage des Datensatzes weiter, reiht pausierte Nachrichten erneut ein, wenn ihre Wiederholungszeit fällig ist, und erholt sich von Abstürzen und Timeouts. Er verarbeitet selbst keine Nachrichteninhalte. Als root gestartet, wechselt er auf das Dienstkonto pepsi, das die Worker erben. Referenz: pepsi-dispatch(1).
29.2. Funktionen¶
Ereignisgesteuertes, gestapeltes Beanspruchen. Lauscht mit
LISTENauf den Kanälenworkqueueundconfig_changed, mit einem periodischenPOLL_INTERVAL-Heartbeat. Bei jedem Aufwachen arbeitet er seine ausstehenden internen Ereignisse ab (sodass die Kapazitäten pro Stage aktuell sind) und beansprucht dann Arbeit für jede Stage mit freier Kapazität in einemUPDATE … RETURNING— jede Stage bis zu ihrer eigenen verbleibenden Kapazität (einLATERAL … ORDER BY workqueue_id LIMIT cappro Stage, das über den partiellenpending-Index bei der Obergrenze stoppt). Da genau ein Dispatcher seriell und als einzigerpending``→``running-Schreiber beansprucht, benötigt das Beanspruchen keinFOR UPDATE SKIP LOCKED.Persistente Worker mit Pipelining. Jede Stage führt
PROGRAM worker(aufPATHgesucht, sofern nicht absolut) als langlebige Prozesse aus, die sich einmal mit der Datenbank verbinden, Nachrichten-IDs von der Standardeingabe lesen und pro Nachricht der Reihe nach mit einem einzeiligen Status (0bei Erfolg) antworten. Einem Worker werden bis zuQUEUE_LIMITIDs auf einmal übergeben (in einem einzigen Puffer auf seine stdin geschrieben); er verarbeitet sie streng seriell, aber die nächste ID steht bereits an, was einen Round-Trip zum Koordinator pro Nachricht spart. Die Kapazität einer Stage an Nachrichten in Bearbeitung beträgt daherQUEUE_LIMIT × PARALLELISM, und weil die Anzahl der Worker unverändert bleibt, kostet das keine zusätzlichen Datenbankverbindungen.Elastische Pools pro Stage. Worker werden bei Bedarf bis zum
PARALLELISMeiner Stage gestartet; die Arbeit wird auf möglichst wenige Worker gepackt, und überschüssige Worker werden nachWORKER_IDLE_TIMEOUTgestoppt — sodass eine ruhige Stage gegen null Worker schrumpft und eine ausgelastete bis zu ihrer Obergrenze wächst. NachMAX_MESSAGESerneuert ein Worker seinen Kindprozess (sobald dessen Pipeline geleert ist); eine Stage, derenPROGRAMnicht starten kann, wird mit einer kurzen Abkühlung zurückgehalten, statt ununterbrochen neu gestartet zu werden.Das Weiterschalten ist der Schreibvorgang der Stage. Eine Stage, die eine Nachricht weiterschaltet, setzt den Datensatz selbst auf
pendingbei der nächsten Stage, in derselben Aktualisierung, diestageverschiebt; der Dispatcher schreibt einen Datensatz bei Erfolg nie um. Die Statusmeldung des Workers bringt den Koordinator zu seiner nächsten Beanspruchung, die den Datensatz für den Pool der nächsten Stage aufgreift (ein Weiterschalten braucht daher keine Benachrichtigung). Ein Worker kann auch einen schnellen Nachfolger in denselben Durchlauf fusionieren (siehe pepsi-dispatch(1), Stage-Fusion).Einplanung von Wiederholungen. Schläft bis zum
timeoutdes frühestenpaused-Datensatzes (oder bis zum fälligen Einsammeln eines Workers, einer Start-Abkühlung oder dem Heartbeat), kippt dann fällige Datensätze zurück aufpendingund scannt erneut.Absturz- und Watchdog-Wiederherstellung. Ein Worker, der einen Wert ungleich null meldet, lässt diese eine Nachricht fehlschlagen und überlebt; ein abgestürzter Kindprozess lässt seine vorderste Nachricht fehlschlagen, und eine vorderste Nachricht, die nicht innerhalb von
MAX_RUNTIMEbeantwortet wird (gemessen ab dem Zeitpunkt, an dem sie an die Spitze gelangte), führt zutimeout. In den Absturz- und Timeout-Fällen liefen die übrigen an den Worker übergebenen Nachrichten nie und werden erneut alspendingeingereiht, statt verloren zu gehen. Gründe werden unter einemstatus = 'running'-Schutz instate.dispatch_errorgeschrieben, sodass eine Stage, die den Datensatz bereits umgeschaltet hat, nie überschrieben wird.Sicherheit bei Start und Herunterfahren. Setzt verwaiste
running-Datensätze (von einem vorherigen Dispatcher) beim Start aufpendingzurück; beiSIGINT/SIGTERMstoppt er die Worker und setzt deren Datensätze (sowie alle beanspruchten, aber nicht zugewiesenen) zurück.Isolation. Höchstens ein Schreiber berührt je einen gegebenen Datensatz der Warteschlange, sodass der gemeinsame Pool mit
READ COMMITTEDläuft; Warteschlangenoperationen sind dennoch in einen Wiederholungshelfer gehüllt, der einen vorübergehenden40001/40P01-Fehler als günstige Versicherung erneut ausführt.Konfigurationsweitergabe. Übergibt
CONFIG_FILEüber-c(vor dem Unterkommandoworker) an die Worker, sodass sie dieselbe Konfiguration laden.Neuladen der Pipeline im laufenden Betrieb. Auf eine
config_changed-Benachrichtigung hin baut er seine eigene Stage-Tabelle aus der Datenbank-Überlagerung neu auf und mustert jeden Worker aus — untätige sofort, beschäftigte, sobald ihre aktuelle Nachricht zurückgemeldet ist —, sodass eine[stage-*]-Änderung ohne Neustart wirksam wird. Ein neu geladener Graph, der sich nicht parsen lässt, beendet den Prozess, statt nach einer Pipeline zu routen, die die Worker nicht mehr teilen. Seine eigenen[pepsi-dispatch]-Optionen werden nicht neu geladen.Statistiken. Sammelt Durchsatz, Kills/Timeouts, Abstürze und Verarbeitungszeit pro Stage sowie die globalen Stage-/Nachrichten-Gesamtwerte im Speicher und schreibt die Deltas in einer einzigen Transaktion in die Tabellen
stage_stats/dispatch_stats: alleSTATS_INTERVAL, immer wenn die Pipeline in den Leerlauf geht (höchstens einmal pro Sekunde, damit die Zahlen eines beendeten Ansturms nicht auf den nächsten Takt warten müssen) und beim Herunterfahren; ohne Deltas wird nichts geschrieben. Er ist der einzige Schreiber dieser Tabellen — Stages, die in den Worker-Durchlauf einer anderen Stage fusioniert wurden, werden auf der Statuszeile dieses Durchlaufs zurückgemeldet, statt vom Worker geschrieben zu werden, was verhindert, dass alle Worker einer Stage um deren einzigen Datensatz konkurrieren. Diese werden über/metricsvon pepsi-httpd exportiert.
29.3. Konfiguration¶
[pepsi-dispatch]: MAX_RUNTIME (Standardwert 300 s), WORKER_IDLE_TIMEOUT (Standardwert 5 s), POLL_INTERVAL (Standardwert 30 s), STATS_INTERVAL (Standardwert 60 s), DB_POOL_SIZE (Standardwert 8) und CONFIG_FILE (kein Standardwert; ohne dieses fallen die Worker auf die voreingestellte Konfigurationssuche zurück). Die Worker-Pools pro Stage werden durch PARALLELISM (Standardwert 4), QUEUE_LIMIT (Standardwert 4, die gleichzeitig an einen Worker übergebenen Nachrichten) und MAX_MESSAGES (Standardwert 1000) in jedem [stage-<name>]-Abschnitt dimensioniert. Siehe pepsi.conf(5) und Konfiguration.