Aller au contenu

La file de tâches

Ce document décrit l'enfilement et le traitement des tâches de fond.

Le fichier de code correspondant est forge_mvc_jobs/queue.py.

1. Le modèle

La file est une table jobs.
On enfile une tâche depuis le code web, un process worker séparé la traite.
Pas de broker, pas de runtime async : le serveur reste synchrone (WSGI).
La réservation d'une tâche est atomique : le worker choisit une candidate, puis la réserve sous garde status='pending'.
Deux workers qui visent la même ligne ne peuvent pas gagner tous les deux, le second voyant zéro ligne affectée.
Plusieurs workers peuvent donc tourner en parallèle.

2. Enfiler (enqueue)

def enqueue(task, payload=None, *, queue="default", max_attempts=1, available_in=0, db=None) -> int

enqueue ajoute une tâche task avec sa charge utile (sérialisée en JSON) et renvoie son identifiant.
max_attempts borne les tentatives, available_in retarde la disponibilité de N secondes.
Lève JobError si task est vide, max_attempts < 1, ou si payload n'est pas sérialisable en JSON.

from forge_mvc_jobs import enqueue

enqueue("email.envoi", {"to": "eleve@exemple.fr"}, max_attempts=3)

3. Traiter (drain, run_worker, process_one)

def process_one(handlers, *, queue="default", db=None) -> bool
def drain(handlers, *, queue="default", max_jobs=None, db=None, stop=None) -> int
def run_worker(handlers, *, queue="default", poll_interval=1.0, db=None, stop=None) -> None

handlers est un dictionnaire {nom_de_tache: fonction} que l'application construit explicitement.
Chaque fonction reçoit la charge utile (un dict).

  • process_one réserve et exécute une tâche ; renvoie True si une tâche a été traitée.
  • drain traite toutes les tâches disponibles en une passe et renvoie le nombre traité.
  • run_worker boucle (vide la file, puis attend si vide).
    L'application la lance depuis son propre script worker, jamais depuis la requête HTTP.

L'arrêt propre (stop)

stop est une fonction sans argument, consultée entre deux tâches, qui demande au worker de s'arrêter.

Elle sert à répondre à un SIGTERM, celui que systemd envoie pour arrêter un service.

import signal

arret = False

def _demander_l_arret(signum, frame):
    global arret
    arret = True

signal.signal(signal.SIGTERM, _demander_l_arret)
run_worker({"email.envoi": envoyer_email}, stop=lambda: arret)

La tâche en cours va à son terme.
L'interrompre ne serait qu'un autre nom pour l'interruption brutale, et laisserait la moitié d'un envoi fait.

Elle n'était consultée qu'une fois la file vidée

stop ne l'était qu'entre deux passes, ce qui la rendait sans effet sur une file chargée.

Mesuré, un worker recevant l'ordre d'arrêt après trois tâches en traitait cinquante avant de le remarquer.
Sous systemd, TimeoutStopSec expirait au bout de quatre-vingt-dix secondes et le worker était tué au milieu d'une tâche, laissée à jobs:reclaim (JOBS-WORKER-GRACEFUL-STOP-001).

Un déploiement se fait justement quand la file est pleine, et c'est le seul moment où ce défaut se voyait.

from forge_mvc_jobs import run_worker

def envoyer_email(payload):
    ...

run_worker({"email.envoi": envoyer_email})

4. Reprise sur échec

Si un gestionnaire lève une exception, la tâche est re-mise en file tant que attempts < max_attempts, sinon marquée failed avec last_error.
Une tâche dont le nom n'a aucun gestionnaire enregistré est marquée failed.

5. Inspecter (pending_count, get_job)

def pending_count(*, queue="default", db=None) -> int
def get_job(job_id, *, db=None) -> Job | None

Job expose id, queue, task, status (pending/running/done/failed), attempts, max_attempts, last_error.

6. Limite V1

Une tâche réservée (running) dont le worker meurt reste running : il n'y a pas de reprise automatique par délai de visibilité dans cette version.

7. Voir aussi