Comment configurer des files d’attente de messages pour l’envoi asynchrone
La mise en file d’attente de messages est un outil intéressant qui permet de créer des sites ou des services Web évolutifs. Une file d’attente de messages permet aux applications de communiquer de manière asynchrone en s’envoyant des messages.
Dans les grandes lignes, la mise en file d’attente de messages est assez simple. Un processus, appelé le Producteur, publie des messages dans une file d’attente, où ils sont stockés jusqu’à ce qu’un processus Consommateur soit prêt à les consommer.
La publication de messages vers un courtier de messages est une opération très rapide, que nous pouvons exploiter pour accélérer nos services Web. Nous pouvons déléguer certaines tâches à des processus en arrière-plan en faisant en sorte que notre service Web publie plutôt un message dans une file d’attente. Nous pouvons ensuite utiliser des processus consommateurs en arrière-plan pour consommer les messages et exécuter les tâches déléguées selon les besoins.
Dans ce guide, nous utiliserons RabbitMQ comme courtier de messages. Nous allons intégrer RabbitMQ dans un exemple d’application Flask pour différer l’envoi d’emails vers un processus en arrière-plan. L’application, Flaskr, vous semblera peut-être familière si vous avez déjà terminé le tutoriel Flask.
Examinons notre fonction de vue d’inscription :
@app.route('/signup', methods=['GET', 'POST'])
def signup():
error = None
if request.method == 'POST':
email = request.form['email']
password = request.form['password']
if not (email or password):
return signup_error('Email Address and Password are required.')
db = get_db()
c = db.cursor()
c.execute('SELECT * FROM users WHERE email=?;', (email,))
if c.fetchone():
return signup_error('Email Addres already has an account.')
c.execute('INSERT INTO users (email, password) VALUES (?, ?);',
(email, pbkdf2_sha256.hash(password)))
db.commit()
send_welcome_email(email)
flash('Account Created')
return redirect(url_for('login'))
else:
return render_template('signup.html')
def send_welcome_email(address):
res = requests.post(
"https://api.mailgun.net/v3/{}/messages".format(app.config['DOMAIN']),
auth=("api", MAILGUN_API_KEY),
data={"from": "Flaskr ".format(app.config['DOMAIN']),
"to": [address],
"subject": "Welcome to Flaskr!",
"text": "Welcome to Flaskr, your account is now active!"}
)
if res.status_code != 200:
# Something terrible happened <img class="emoji" alt="🙁" src="https://s.w.org/images/core/emoji/12.0.0-1/svg/1f641.svg">
raise MailgunError("{}-{}".format(res.status_code, res.reason))
def signup_error(error):
return render_template('signup.html', error=error)
Actuellement, la fonction de vue vérifie la base de données pour s’assurer qu’aucun compte n’existe déjà avec cette adresse email, puis ajoute le nouvel utilisateur à la base de données. La vue envoie ensuite immédiatement un email de bienvenue à l’utilisateur avant de terminer la requête.
Cette vue présente quelques problèmes que nous pouvons améliorer grâce à une file d’attente de messages. Il n’est pas nécessaire de faire attendre l’utilisateur pendant l’envoi de l’email de bienvenue. Il serait préférable de pouvoir répondre à l’utilisateur dès que possible. Nous voulons également un moyen de réessayer d’envoyer l’email de bienvenue si, pour une raison quelconque, nous n’y parvenons pas du premier coup.
Tout d’abord, exécutons une instance de RabbitMQ dans un conteneur Docker :
docker run -d --hostname my-rabbit -p 4369:4369 -p 5672:5672
-p 35197:35197 --name rabbitmq rabbitmq:3
Ajoutons maintenant un peu de code pour initialiser une file d’attente au démarrage de notre application :
def connect_queue():
if not hasattr(g, 'rabbitmq'):
g.rabbitmq = pika.BlockingConnection(
pika.ConnectionParameters(app.config['RABBITMQ_HOST'])
)
return g.rabbitmq
def get_welcome_queue():
if not hasattr(g, 'welcome_queue'):
conn = connect_queue()
channel = conn.channel()
channel.queue_declare(queue='welcome_queue', durable=True)
channel.queue_bind(exchange='amq.direct', queue='welcome_queue')
g.welcome_queue = channel
return g.welcome_queue
@app.teardown_appcontext
def close_queue(error)
if hasattr(g, 'rabbitmq'):
g.rabbitmq.close()
Nous déclarons une file d’attente appelée « welcome_queue » que nous utiliserons pour envoyer des messages aux workers avec l’adresse email à laquelle l’email de bienvenue doit être envoyé.
Ensuite, nous pouvons mettre à jour notre vue d’inscription pour publier un message dans la file d’attente au lieu d’envoyer l’email de bienvenue.
@app.route('/signup', methods=['GET', 'POST'])
def signup():
error = None
if request.method == 'POST':
email = request.form['email']
password = request.form['password']
if not (email or password):
return signup_error('Email Address and Password are required.')
db = get_db()
c = db.cursor()
c.execute('SELECT * FROM users WHERE email=?;', (email,))
if c.fetchone():
return signup_error('Email Addres already has an account.')
c.execute('INSERT INTO users (email, password) VALUES (?, ?);',
(email, pbkdf2_sha256.hash(password)))
db.commit()
q = get_welcome_queue()
q.basic_publish(
exchange='amq.direct',
routing_key='welcome_queue',
body=email,
properties=pika.BasicProperties(
delivery_mode=_DELIVERY_MODE_PERSISTENT
)
)
flash('Account Created')
return redirect(url_for('login'))
else:
return render_template('signup.html')
Nous pouvons maintenant écrire notre code de worker. L’avantage d’utiliser une file d’attente de messages est que le processus worker peut s’exécuter n’importe où. Il peut s’exécuter sur des serveurs dédiés aux workers, ou bien aux côtés de votre service Web. Voici à quoi ressemblerait un simple script de worker :
import pika
import requests
# Configuration
DOMAIN = 'example.com'
MAILGUN_API_KEY = 'YOUR_MAILGUN_API_KEY'
RABBITMQ_HOST = 'localhost'
connection = pika.BlockingConnection(
pika.ConnectionParameters(host=RABBITMQ_HOST)
)
channel = connection.channel()
channel.queue_declare(queue='welcome_queue', durable=True)
class Error(Exception):
pass
class MailgunError(Error):
def __init__(self, message):
self.message = message
def send_welcome_message(ch, method, properties, body):
address = body.decode('UTF-8')
print("Sending welcome email to {}".format(address))
res = requests.post(
"https://api.mailgun.net/v3/{}/messages".format(DOMAIN),
auth=("api", MAILGUN_API_KEY),
data={"from": "Flaskr ".format(DOMAIN),
"to": [address],
"subject": "Welcome to Flaskr!",
"text": "Welcome to Flaskr, your account is now active!"}
)
ch.basic_ack(delivery_tag=method.delivery_tag)
if res.status_code != 200:
# Something terrible happened :-O
raise MailgunError("{}-{}".format(res.status_code, res.reason))
channel.basic_consume(send_welcome_message, queue='welcome_queue')
channel.start_consuming()
Nous avons presque terminé ! Nous avons réussi à dissocier l’envoi de notre email de bienvenue du service Web. Il ne nous reste plus qu’à réessayer d’envoyer l’email de bienvenue au cas où son envoi échouerait pour une raison quelconque. Une façon intelligente d’ajouter une logique de nouvelle tentative à votre application consiste à utiliser une file d’attente de lettres mortes.
Nous ajouterons une deuxième file d’attente, la « retry_queue », que nous utiliserons pour placer temporairement un message si nous rencontrons une erreur d’envoi d’email. Les messages de la file d’attente de nouvelles tentatives auront une date d’expiration dans le futur. Nous allons configurer RabbitMQ de sorte que, lorsqu’un message expire, il soit replacé dans notre file d’attente de bienvenue, prêt à être récupéré par un worker pour tenter d’envoyer à nouveau l’email :
retry_channel = connection.channel()
retry_channel.queue_declare(
queue='retry_queue',
durable=True,
arguments={
'x-message-ttl': RETRY_DELAY_MS,
'x-dead-letter-exchange': 'amq.direct',
'x-dead-letter-routing-key': 'welcome_queue'
}
)
def send_welcome_message(ch, method, properties, body):
address = body.decode('UTF-8')
print("Sending welcome email to {}".format(address))
res = requests.post(
"https://api.mailgun.net/v3/{}/messages".format(DOMAIN),
auth=("api", MAILGUN_API_KEY),
data={"from": "Flaskr ".format(DOMAIN),
"to": [address],
"subject": "Welcome to Flaskr!",
"text": "Welcome to Flaskr, your account is now active!"}
)
ch.basic_ack(delivery_tag=method.delivery_tag)
if res.status_code != 200:
print("Error sending to {}. {} {}. Retrying...".format(
address, res.status_code, res.reason
))
retry_channel.basic_publish(
exchange='',
routing_key='retry_queue',
body=address,
properties=pika.BasicProperties(
delivery_mode=_DELIVERY_MODE_PERSISTENT
)
)
Et c’est tout ! Nous avons introduit avec succès l’envoi asynchrone dans notre application à l’aide de RabbitMQ. Pour télécharger le dépôt de code source complet de l’exemple, consultez mon dépôt Github.
Découvrez d’autres guides comme celui-ci en vous inscrivant au blog. Vous recevrez une mise à jour chaque semaine avec les derniers articles de l’équipe Mailgun. Et si vous n’utilisez pas encore Mailgun pour envoyer, recevoir et suivre les emails de votre application, inscrivez-vous ci-dessous !