Email

Cómo configurar colas de mensajes para el envío asíncrono

La gestión de colas de mensajes es una herramienta interesante que puede ayudarnos a crear sitios web o servicios web escalables. Una cola de mensajes permite a las aplicaciones comunicarse de forma asíncrona mediante el envío de mensajes entre sí.
Imagen para Cómo configurar colas de mensajes para el envío asíncrono

La gestión de colas de mensajes es una herramienta interesante que puede ayudarnos a crear sitios web o servicios web escalables. Una cola de mensajes permite a las aplicaciones comunicarse de forma asíncrona mediante el envío de mensajes entre sí.

A grandes rasgos, la gestión de colas de mensajes es bastante sencilla. Un proceso, llamado Productor, publica mensajes en una cola, donde se almacenan hasta que un proceso Consumidor está listo para consumirlos.

Publicar mensajes en un bróker de mensajes es una operación muy rápida y podemos aprovecharlo para acelerar nuestros servicios web. En su lugar, podemos delegar algunas tareas a procesos en segundo plano haciendo que nuestro servicio web publique un mensaje en una cola. A continuación, podemos hacer que procesos consumidores en segundo plano consuman los mensajes y realicen las tareas delegadas según sea necesario.

En esta guía utilizaremos RabbitMQ como bróker de mensajes. Integraremos RabbitMQ en una aplicación Flask de ejemplo para posponer el envío del email a un proceso en segundo plano. La aplicación, Flaskr, puede que te resulte familiar si ya has completado el Tutorial de Flask.

Echemos un vistazo a nuestra función de vista de registro:

                                

                                    @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)
                                
                            

Actualmente, la función de vista comprueba la base de datos para asegurarse de que no existe ya una cuenta con esa dirección de email y, a continuación, añade el nuevo usuario a la base de datos. A continuación, la vista envía inmediatamente un mensaje de email de bienvenida al usuario antes de finalizar la solicitud.

Hay un par de problemas con esta vista que podemos mejorar con una cola de mensajes. No tenemos que hacer esperar al usuario por una respuesta mientras enviamos el email de bienvenida. Sería mejor si pudiéramos responder al usuario lo antes posible. También queremos una forma de reintentar el envío del email de bienvenida si por alguna razón no podemos hacerlo al primer intento.

En primer lugar, vamos a ejecutar una instancia de RabbitMQ en un contenedor de Docker:

                                

                                        docker run -d --hostname my-rabbit -p 4369:4369 -p 5672:5672  
     -p 35197:35197 --name rabbitmq rabbitmq:3
                                
                            

Ahora vamos a añadir algo de código para inicializar una cola cuando se inicie nuestra aplicación:

                                

                                    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()
                                
                            

Estamos declarando una cola llamada “welcome_queue” que utilizaremos para enviar mensajes a los workers con la dirección de email a la que se debe enviar el email de bienvenida.

A continuación, podemos actualizar nuestra vista de registro para que publique un mensaje en la cola en lugar de enviar el email de bienvenida.

                                

                                    @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')
                                
                            

Ahora podemos escribir nuestro código de worker. Lo mejor de usar una cola de mensajes es que el proceso worker puede ejecutarse en cualquier lugar. Puede ejecutarse en servidores dedicados para workers o quizá junto a tu servicio web. Este sería el aspecto de un script worker sencillo:

                                

                                    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()  
                                
                            

¡Ya casi hemos terminado! Hemos desvinculado con éxito el envío de nuestro email de bienvenida del servicio web. Ahora lo único que falta es que volvamos a intentar enviar el email de bienvenida en caso de que falle el envío por algún motivo. Una forma inteligente de añadir una lógica de reintento a tu aplicación es utilizar una cola de mensajes muertos (dead letter queue).

Añadiremos una segunda cola, “retry_queue”, que usaremos para colocar temporalmente un mensaje si nos encontramos con un error de envío del email. Los mensajes de la cola de reintento tendrán una fecha de caducidad en el futuro. Configuraremos RabbitMQ de tal forma que cuando un mensaje caduque volverá a nuestra cola de bienvenida para que un worker lo recoja e intente enviar el email de nuevo:

                                

                                    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
            )
        )
                                
                            

¡Y ya está! Hemos introducido con éxito el envío asíncrono en nuestra aplicación utilizando RabbitMQ. Para descargar el repositorio completo del código fuente con el ejemplo práctico, echa un vistazo a mi repo de GitHub.

Recibe más guías como esta suscribiéndote al blog. Recibirás una actualización cada semana con las últimas publicaciones del equipo de Mailgun. Y si aún no utilizas Mailgun para enviar, recibir y hacer un seguimiento de los emails de tu aplicación, ¡deberías registrarte a continuación!