29. pepsi-dispatch

L’unique coordinateur de longue durée qui fait avancer les messages à travers les étapes.

29.1. Rôle

pepsi-dispatch est le seul processus qui démarre les programmes d’étape. Exécutez-en exactement un par système, sous la forme pepsi-dispatch serve (l’unique sous-commande ; les indicateurs globaux -c/-L/-v viennent avant elle). Il revendique les lignes nouvellement mises en file et transmet chacune à un processus worker inactif de l’étape de la ligne, remet en file les messages en pause lorsque leur heure de nouvelle tentative est échue, et récupère après les plantages et les timeouts. Il ne traite lui-même aucun contenu de message. Démarré en tant que root, il abandonne ses privilèges au profit du compte de service pepsi, dont les workers héritent. Référence : pepsi-dispatch(1).

29.2. Fonctionnalités

  • Revendication pilotée par événements, par lots. Fait un LISTEN sur les canaux workqueue et config_changed avec un battement de cœur POLL_INTERVAL périodique. À chaque réveil, il draine ses événements internes en attente (de sorte que les capacités par étape soient à jour) puis revendique le travail pour chaque étape disposant de capacité libre en un seul UPDATE … RETURNING — chaque étape jusqu’à sa propre capacité restante (un LATERAL … ORDER BY workqueue_id LIMIT cap par étape qui s’arrête au plafond via l’index partiel pending). Avec exactement un dispatcher revendiquant en série comme unique écrivain pending``→``running, la revendication n’a besoin d’aucun FOR UPDATE SKIP LOCKED.

  • Workers persistants et pipelinés. Chaque étape exécute PROGRAM worker (trouvé sur le PATH sauf s’il est absolu) sous forme de processus de longue durée qui se connectent une fois à la base de données et lisent les identifiants de message sur l’entrée standard, répondant par un statut d’une ligne (0 en cas de succès) par message, dans l’ordre. Un worker reçoit jusqu’à QUEUE_LIMIT identifiants à la fois (écrits sur son entrée standard en un seul tampon) ; il les traite strictement en série, mais l’identifiant suivant est déjà en file, ce qui supprime un aller-retour avec le coordinateur par message. La capacité en vol d’une étape est donc QUEUE_LIMIT × PARALLELISM, et comme le nombre de workers est inchangé, cela ne coûte aucune connexion supplémentaire à la base de données.

  • Pools élastiques par étape. Les workers sont démarrés à la demande jusqu’au PARALLELISM de l’étape ; le travail est tassé sur le moins de workers possible, et les workers en surplus sont arrêtés après WORKER_IDLE_TIMEOUT — de sorte qu’une étape calme retombe vers zéro worker et qu’une étape chargée croît jusqu’à son plafond. Après MAX_MESSAGES messages, un worker recycle son processus enfant (une fois son pipeline drainé) ; une étape dont le PROGRAM ne peut pas démarrer est mise en attente avec un court délai de refroidissement plutôt que de tourner en boucle.

  • Avancer est l’écriture de l’étape. Une étape qui fait avancer un message positionne elle-même la ligne à pending à l’étape suivante, dans la même mise à jour qui déplace stage ; le dispatcher ne réécrit jamais une ligne en cas de succès. Le rapport d’état du worker ramène le coordinateur à sa revendication suivante, qui prend la ligne pour le pool de l’étape suivante (de sorte qu’une avancée n’a besoin d’aucune notification). Un worker peut aussi fusionner un successeur rapide dans la même passe (voir pepsi-dispatch(1), Fusion d’étapes).

  • Ordonnancement des réessais. Dort jusqu’au timeout de la ligne paused la plus proche (ou jusqu’à une récupération de worker due, un refroidissement de lancement ou le battement de cœur), puis rebascule les lignes dues en pending et rebalaie.

  • Récupération après plantage et chien de garde. Un worker qui rapporte un statut non nul fait échouer ce seul message et survit ; un processus enfant planté fait échouer son message en tête de file, et un message en tête de file resté sans réponse pendant MAX_RUNTIME (compté depuis qu’il a atteint la tête) → timeout. Dans les cas de plantage et de timeout, les autres messages pipelinés du worker n’ont jamais été exécutés et sont remis en pending plutôt que perdus. Les raisons sont écrites dans state.dispatch_error sous une garde status = 'running', afin qu’une étape ayant déjà fait transiter la ligne ne soit jamais écrasée.

  • Sûreté au démarrage et à l’arrêt. Réinitialise en pending au démarrage les lignes running orphelines (laissées par un dispatcher précédent) ; sur SIGINT/SIGTERM, arrête les workers et réinitialise leurs lignes (ainsi que toute ligne revendiquée mais non assignée).

  • Isolation. Au plus un écrivain touche jamais une ligne de file donnée, de sorte que le pool partagé fonctionne en READ COMMITTED ; les opérations de file sont tout de même enveloppées dans un helper de réessai qui réexécute un échec transitoire 40001/40P01 comme assurance bon marché.

  • Propagation de la config. Passe CONFIG_FILE aux workers via -c (avant la sous-commande worker) afin qu’ils chargent la même configuration.

  • Rechargement à chaud du pipeline. Sur une notification config_changed, il reconstruit sa propre table d’étapes à partir de la surcouche en base de données et retire tous les workers — immédiatement ceux qui sont inactifs, et ceux qui sont occupés lorsque leur message en cours rend compte — de sorte qu’une modification de [stage-*] prenne effet sans redémarrage. Un graphe rechargé qui ne s’analyse pas met fin au processus plutôt que d’acheminer selon un pipeline que les workers ne partagent plus. Ses propres options [pepsi-dispatch] ne sont pas rechargées.

  • Statistiques. Accumule en mémoire le débit par étape, les kills/timeouts, les plantages et le temps de traitement, plus les totaux globaux d’étapes et de messages, et vide les deltas vers les tables stage_stats/dispatch_stats en une seule transaction à chaque STATS_INTERVAL, dès que le pipeline devient inactif (au plus une fois par seconde, afin que les chiffres d’une rafale terminée n’attendent pas le tic suivant) et à l’arrêt ; rien n’est écrit en l’absence de deltas. C’est l”unique écrivain de ces tables — les étapes fusionnées dans la passe de worker d’une autre étape sont rapportées sur la ligne de statut de cette passe au lieu d’être écrites par le worker, ce qui évite que tous les workers d’une étape ne se disputent l’unique ligne de cette étape. Ces valeurs sont exportées par le /metrics de pepsi-httpd.

29.3. Configuration

[pepsi-dispatch] : MAX_RUNTIME (par défaut 300 s), WORKER_IDLE_TIMEOUT (par défaut 5 s), POLL_INTERVAL (par défaut 30 s), STATS_INTERVAL (par défaut 60 s), DB_POOL_SIZE (par défaut 8) et CONFIG_FILE (sans valeur par défaut ; sans lui, les workers retombent sur la recherche de configuration par défaut). Les pools de workers par étape sont dimensionnés par PARALLELISM (par défaut 4), QUEUE_LIMIT (par défaut 4, le nombre de messages pipelinés vers un worker à la fois) et MAX_MESSAGES (par défaut 1000) dans chaque section [stage-<name>]. Voir pepsi.conf(5) et Configuration.

29.4. Voir aussi

Architecture, pepsi-queue, pepsi-ingress, pepsi-dispatch(1).