Exemple de configuration simple de Celery

Sommaire- Producteur

  • Configuration du consommateur

Exemple de fichier de configuration Celery

Exemple simple Celery

arbre  -I 'containerd|vminit|__pycache__'
.
# app.py appartient au producteur
├── app.py
# celery_app pour configurer le consommateur et les informations de file d'attente
└── celery_app
# config.py informations de configuration
    ├── config.py
# __init__.py fichier d'initialisation de l'instance celery
    ├── __init__.py
# task1, task2 pour enregistrer les fonctions
    ├── task1.py
    └── task2.py


Producteru

app.py

from celery_app import tache1, tache2

# Appeler la méthode delay(), la tâche est soumise de manière asynchrone à la file d'attente de messages,
# le worker celery récupère la tâche en arrière-plan et l'exécute
# Retourne un objet AsyncResult, qui peut être utilisé pour suivre l'état et le résultat de la tâche
# Le processus d'exécution asynchrone des tâches : 1. Appeler delay() pour soumettre la tâche à la file d'attente de messages,
# 2. Retourner AsyncResult pour interroger les résultats
t2 = tache2.calcul.delay(3,4)
print(t2)

t1 = tache1.multiplication.delay(2,10)
print(t1)



Configuration du cnosommateur

celery_app/config.py

# BROKER_URL est l'adresse du courtier de messages (broker)
URL_BROKER = "redis://127.0.0.1:6379/2"
# CELERY_RESULT_BACKEND est le backend pour les résultats des tâches
CELERY_RESULT_BACKEND = "redis://127.0.0.1:6379/3"
# CELERY_TIMEZONE définit le fuseau horaire, utilisé pour les tâches planifiées
CELERY_TIMEZONE = "Asia/Shanghai"
# UTC

# Celery peut automatiquement charger et enregistrer les modules de tâches au démarrage
IMPORTS_CELERY = (
    'celery_app.tache1',
    'celery_app.tache2',
)



celery_app/init.py

from celery import Celery

# Créer une instance d'application celery, nommée demo, pour distinguer différentes instances celery
# Cette instance est l'objet central de celery, responsable de l'ordonnancement des tâches,
# du routage des tâches et du chargement de la configuration
application = Celery('demo')

# Charger la configuration via l'instance Celery
application.config_from_object('celery_app.config')


celery_app/task1.py

# Importer l'instance app de Celery
from celery_app import application
import time

# application.task transforme une fonction ordinaire en tâche Celery,
# cette tâche peut être exécutée via le mécanisme asynchrone de Celery
@application.task
def multiplication(x,y):
    time.sleep(5)
    return x * y 


celery_app/task2.py

from celery_app import application
import time

@application.task
def calcul(x,y):
    time.sleep(20)
    return x / y


Démarrage du worker celery

celery -A celery_app worker -l INFO


Vérification des journaux de démarrage celery

[configuration]
# app: demo indique que le nom de l'instance celery en cours d'exécution est demo
.> app:         demo:0x7f0014e6f3d0
# transport représente le courtier de messages (broker),
# responsable de la transmission du producteur au consommateur
.> transport:   redis://127.0.0.1:6379/2
# results représente le backend des résultats, utilisé pour stocker les résultats d'exécution des tâches
.> results:     redis://127.0.0.1:6379/3
# concurrency indique que le worker celery peut traiter deux tâches simultanément
.> concurrency: 2 (prefork)
# indique que la surveillance des événements peut être activée avec -E,--events
.> task events: OFF (enable -E to monitor tasks in this worker)

# représente les files d'attente actuellement écoutées par le worker celery
[files d'attente]
# celery représente le nom de la file d'attente, par défaut Celery utilise une file d'attente nommée celery
# pour recevoir et traiter les tâches
# exchange est l'échangeur de messages, de type direct,
# un échangeur direct signifie que les messages seront routés directement
# vers les files d'attente avec une clé de routage correspondante
# key=celery est la clé de routage, indiquant que seuls les messages avec la clé de routage celery
# seront routés vers cette file d'attente, par défaut toutes les tâches sont envoyées
# à la file d'attente celery via cette clé de routage
.> celery           exchange=celery(direct) key=celery

# représente toutes les tâches déjà chargées et enregistrées au démarrage du worker celery
[tâches]
# représente 2 tâches déjà chargées
. celery_app.tache1.multiplication
. celery_app.tache2.calcul



Invocation du producteur

python app.py 
# retourne l'identifiant de la tâche
#49c5e568-7fb2-49ab-8be7-c5ee6124011f
#71b170fc-09a9-476e-be91-567dd78b163c


Résultats retournés

# format du journal
#[timestamp: INFO/MainProcess] Task celery_app.taches1.multiplication[task_id] received
# réception de la tâche
[2024-10-19 09:06:35,842: INFO/MainProcess] Task celery_app.tache1.multiplication[414dc010-f7f7-473d-83dc-282d4dfc4cca] received
# tâche réussie
# chaque tâche est exécutée par un processus worker celery différent (ForkPoolWorker).
#[timestamp: INFO/ForkPoolWorker-2] Task celery_app.tache1.multiplication[task_id] état en temps écoulé: résultat
[2024-10-19 09:06:40,844: INFO/ForkPoolWorker-2] Task celery_app.tache1.multiplication[414dc010-f7f7-473d-83dc-282d4dfc4cca] succeeded in 5.0014993995428085s: 20


Étiquettes: Celery Redis Python tâches asynchrones file d'attente

Publié le 26 août à 04h14