Files
presentations/2026/celery/slides/02-progressbar.md
T

5.8 KiB

Les barres de progression dans Celery

  • Intro
  • Communication Worker->Django
  • Communication Django->Client
  • Accès concurrent

Intro

Liste des courses

  • Progression en temps réel
  • Barre de progression propre à chaque tâche instanciée

Contraintes

  • Tâches possiblement longues
  • Éviter le polling (fonctionner sur évènements)
  • Pas de résidus en cas de plantage

Situation initiale

  • Fonctionnalité non existante dans Celery
  • Solutions sur le net : progression par comptage de micro-tâches
  • Une bibliothèque existante sur le net : celery-progress

Complexité supplémentaire

<script type="text/mermaid"> %%{init: {'theme': 'light', 'themeVariables': { 'darkMode': false }}}%% sequenceDiagram actor User participant Django participant Broker@{ "type" : "database" } participant Worker participant Postgres@{ "type" : "database" } User->>Django: Click button Django->>Broker: Send message : schedule task Broker->>Worker: Acquire task activate Worker Worker->>Postgres: Save results deactivate Worker User ->>Django: Query task result Django->>Postgres: Query task result Postgres-->>Django: task result Django-->>User: task result </script>
<script type="text/mermaid"> sequenceDiagram actor User participant Django participant Broker@{ "type" : "database" } participant Worker participant Postgres@{ "type" : "database" } User->>Django: Click button Django->>Broker: Send message : schedule task Broker->>Worker: Acquire task activate Worker rect rgb(255, 224, 224) loop Worker->>Worker: Work Worker->>Broker: Progress event Broker->>Django: Progress event Django->>User: Progress event end end Worker->>Postgres: Save results deactivate Worker User ->>Django: Query task result Django->>Postgres: Query task result Postgres-->>Django: task result Django-->>User: task result </script>
<script type="text/mermaid"> sequenceDiagram actor User participant Django participant Worker participant Postgres@{ "type" : "database" } User->>Django: Click button Django->>()Worker: Send message : schedule task activate Worker rect rgb(255, 224, 224) loop Worker->>Worker: Work Worker->>()Django: Progress event Django->>User: Progress event end end Worker->>Postgres: Save results deactivate Worker User ->>Django: Query task result Django->>Postgres: Query task result Postgres-->>Django: task result Django-->>User: task result </script>

Communication worker -> server

  • Coté Worker
  • Coté Django
  • Accès concurrent

Coté Worker

https://docs.celeryq.dev/en/main/reference/celery.events.html

The worker has the ability to send a message whenever some event happens. These events are then captured by tools like Flower, and celery events to monitor the cluster.

Coté Worker

  • celery.events.EventDispatcher 📢
  • celery.events.EventReceiver 👂

Coté Worker

class PtfTask(AbortableTask):
    _progress_dispatcher = None

    @property
    def progress_dispatcher(self):
        # Singleton
        ...
        
    def update_state(self, task_id, state, meta, **kwargs):
        self.progress_dispatcher.send(
            "ptf-task-progress",
            data=meta,
            uuid=task_id or self.request.id,
        )
        return super().update_state(task_id, state, meta, **kwargs)

On envoie un évènement ptf-task-progress quand on met à jour l'état de la tâche

Coté Worker

class PtfTask(AbortableTask):
    _progress_dispatcher = None

    @property
    def progress_dispatcher(self):
        # Singleton
        ...
        
    def update_state(self, task_id, state, meta, **kwargs):
        self.progress_dispatcher.send(
            "ptf-task-progress",
            data=meta,
            uuid=task_id or self.request.id,
        )
        return super().update_state(task_id, state, meta, **kwargs)

Point important: on peut faire avancer la progression de n'importe quelle tâche tant qu'on a son ID!

Coté Worker

@shared_task(name="crawler.tasks.toto", base=PtfTask, bind=True)
def print_toto(self, toto = "toto"):   
    meta = {"current": 0, "total": len(toto)}
    # Initialize progress bar at 0
    self.update_state(
        meta=meta,
        state=states.STARTED,
    )
    
    for index, letter in enumerate(toto):
        # Do work
        print(letter)
        
        # Update progress bar
        meta["current"] = index
        self.update_state(
            meta=meta,
            state=states.STARTED,
        )

Coté Django

from celery import Celery, current_app

def send_sse_message(event):
    ...
    
handlers = {
    "ptf-task-progress": send_sse_message,
}

receiver = current_app.events.Receiver(
    channel=current_app.connection_for_read(), app=current_app, handlers=handlers
)

On envoie un message au client quand on reçoit un évènement ptf-task-progress

Coté Django

Django 4.2 Workaround