Como configurar filas de mensagens para envio assíncrono
O enfileiramento de mensagens é uma ferramenta interessante que pode nos ajudar a criar sites ou serviços web escaláveis. Uma fila de mensagens permite que os aplicativos se comuniquem de forma assíncrona, enviando mensagens entre si.
Em alto nível, o enfileiramento de mensagens é bem simples. Um processo, chamado Produtor, publica mensagens em uma fila, onde elas são armazenadas até que um processo Consumidor esteja pronto para consumi-las.
A publicação de mensagens em um message broker é uma operação muito rápida, e podemos aproveitar isso para acelerar nossos serviços web. Podemos delegar algumas tarefas a processos em segundo plano, fazendo com que nosso serviço web publique uma mensagem em uma fila. Assim, podemos ter processos consumidores em segundo plano que consomem as mensagens e executam as tarefas delegadas conforme necessário.
Neste guia, usaremos RabbitMQ como message broker. Vamos integrar o RabbitMQ a um aplicativo Flask de exemplo para adiar o envio de e-mails para um processo em segundo plano. O aplicativo, Flaskr, pode parecer familiar se você já concluiu o Tutorial do Flask.
Vamos dar uma olhada na nossa função de visualização de cadastro:
@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)
No momento, a função de visualização verifica o banco de dados para garantir que não exista uma conta com esse endereço de e-mail e, em seguida, adiciona o novo usuário ao banco de dados. A visualização então envia imediatamente uma mensagem de e-mail de boas-vindas ao usuário antes de concluir a solicitação.
Existem alguns problemas com essa visualização que podemos melhorar com uma fila de mensagens. Não precisamos fazer o usuário esperar por uma resposta enquanto enviamos o e-mail de boas-vindas. Seria melhor se pudéssemos responder ao usuário o mais rápido possível. Também queremos uma maneira de tentar enviar o e-mail de boas-vindas novamente se, por algum motivo, não conseguirmos na primeira tentativa.
Primeiro, vamos executar uma instância do RabbitMQ em um contêiner Docker:
docker run -d --hostname my-rabbit -p 4369:4369 -p 5672:5672
-p 35197:35197 --name rabbitmq rabbitmq:3
Agora, vamos adicionar um pouco de código para inicializar uma fila quando nosso aplicativo iniciar:
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 uma fila chamada “welcome_queue”, que usaremos para enviar mensagens aos workers com o endereço de e-mail para o qual o e-mail de boas-vindas precisa ser enviado.
Em seguida, podemos atualizar nossa visualização de cadastro para publicar uma mensagem na fila, em vez de enviar o e-mail de boas-vindas.
@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')
Agora podemos escrever o código do nosso worker. A parte legal de usar uma fila de mensagens é que o processo do worker pode ser executado em qualquer lugar. Ele pode ser executado em servidores de workers dedicados ou até mesmo junto com seu serviço web. É assim que seria um script simples 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()
Estamos quase terminando! Desacoplamos com sucesso o envio do nosso e-mail de boas-vindas do serviço web. Agora, a única coisa que falta é tentarmos enviar novamente o e-mail de boas-vindas, caso o envio falhe por algum motivo. Uma maneira inteligente de adicionar uma lógica de nova tentativa ao seu aplicativo é usar uma fila de mensagens não entregues (dead letter queue).
Adicionaremos uma segunda fila, a “retry_queue”, que usaremos para colocar uma mensagem temporariamente se encontrarmos um erro de envio de e-mail. As mensagens na fila de novas tentativas terão uma data de expiração no futuro. Configuraremos o RabbitMQ de forma que, quando uma mensagem expirar, ela seja colocada de volta em nossa fila de boas-vindas, pronta para ser coletada por um worker para tentar enviar o e-mail de novo:
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
)
)
Simples assim! Introduzimos com sucesso o envio assíncrono em nosso aplicativo usando o RabbitMQ. Para baixar o repositório completo com o código-fonte de exemplo em funcionamento, confira o meu repositório no GitHub.
Receba mais guias como este assinando o blog. Você receberá uma atualização toda semana com as publicações mais recentes da equipe da Mailgun. E se você ainda não usa a Mailgun para enviar, receber e monitorar os e-mails do seu aplicativo, cadastre-se abaixo!