Nachrichten-Warteschlangen für den asynchronen Versand einrichten
Die Nachrichten-Warteschlange ist ein interessantes Tool, das uns bei der Erstellung skalierbarer Websites oder Web-Dienste unterstützt. Eine Nachrichten-Warteschlange ermöglicht Anwendungen die asynchrone Kommunikation, indem sie Nachrichten untereinander austauschen.
Im Grunde genommen sind Nachrichten-Warteschlangen ziemlich einfach. Ein Prozess, der sogenannte Producer, veröffentlicht Nachrichten in einer Warteschlange, wo sie gespeichert werden, bis ein Consumer-Prozess bereit ist, diese abzurufen.
Das Veröffentlichen von Nachrichten an einen Nachrichten-Broker ist ein sehr schneller Vorgang, den wir nutzen können, um unsere Web-Dienste zu beschleunigen. Wir können einige Aufgaben an Hintergrundprozesse delegieren, indem unser Web-Dienst stattdessen eine Nachricht an eine Warteschlange sendet. Anschließend können Consumer-Prozesse im Hintergrund diese Nachrichten abrufen und die delegierten Aufgaben bei Bedarf ausführen.
In diesem Leitfaden verwenden wir RabbitMQ als Nachrichten-Broker. Wir integrieren RabbitMQ in eine Flask-Beispielanwendung, um den E-Mail-Versand an einen Hintergrundprozess auszulagern. Die Anwendung Flaskr kommt Ihnen vielleicht bekannt vor, wenn Sie bereits Folgendes abgeschlossen haben: Flask-Tutorial.
Sehen wir uns unsere View-Funktion für die Anmeldung an:
@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)
Aktuell prüft die View-Funktion in der Datenbank, ob bereits ein Konto mit dieser E-Mail-Adresse existiert, und fügt dann den neuen Nutzer zur Datenbank hinzu. Anschließend sendet die Ansicht sofort eine Willkommens-E-Mail an den Nutzer, bevor die Anfrage abgeschlossen wird.
Bei dieser Ansicht gibt es einige Punkte, die wir mit einer Nachrichten-Warteschlange verbessern können. Wir müssen den Nutzer nicht auf eine Antwort warten lassen, während wir die Willkommens-E-Mail versenden. Es wäre besser, wenn wir dem Nutzer so schnell wie möglich antworten könnten. Zudem benötigen wir eine Möglichkeit, den Versand der Willkommens-E-Mail zu wiederholen, falls er beim ersten Versuch aus irgendeinem Grund fehlschlägt.
Führen wir zunächst eine Instanz von RabbitMQ in einem Docker-Container aus:
docker run -d --hostname my-rabbit -p 4369:4369 -p 5672:5672
-p 35197:35197 --name rabbitmq rabbitmq:3
Fügen wir nun etwas Code hinzu, um beim Start unserer Anwendung eine Warteschlange zu initialisieren:
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()
Wir deklarieren eine Warteschlange namens „welcome_queue“, über die wir Nachrichten mit der E-Mail-Adresse, an die die Willkommens-E-Mail gesendet werden soll, an Worker senden.
Als Nächstes können wir unsere Anmeldeansicht so aktualisieren, dass sie eine Nachricht in der Warteschlange veröffentlicht, anstatt die Willkommens-E-Mail direkt zu versenden.
@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')
Jetzt können wir unseren Worker-Code schreiben. Das Besondere an einer Nachrichten-Warteschlange ist, dass der Worker-Prozess überall ausgeführt werden kann. Er kann auf dedizierten Worker-Servern oder auch parallel zu Ihrem Web-Dienst laufen. So würde ein einfaches Worker-Skript aussehen:
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()
Wir sind fast fertig. Wir haben den Versand unserer Willkommens-E-Mail erfolgreich vom Web-Dienst entkoppelt. Jetzt fehlt nur noch, dass wir den Versand der Willkommens-E-Mail wiederholen können, falls dieser aus irgendeinem Grund fehlschlägt. Eine clevere Methode, um eine Wiederholungslogik in Ihre Anwendung einzubauen, ist die Verwendung einer Dead-Letter-Warteschlange.
Wir fügen eine zweite Warteschlange hinzu, die „retry_queue“, in der wir eine Nachricht vorübergehend ablegen, falls ein Fehler beim E-Mail-Versand auftritt. Die Nachrichten in der Wiederholungs-Warteschlange erhalten ein Ablaufdatum in der Zukunft. Wir konfigurieren RabbitMQ so, dass eine Nachricht nach Ablauf dieser Frist zurück in unsere Willkommens-Warteschlange verschoben wird. Dort kann sie von einem Worker abgerufen werden, um den E-Mail-Versand erneut zu versuchen:
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
)
)
Und das war’s. Wir haben den asynchronen Versand mithilfe von RabbitMQ erfolgreich in unsere Anwendung integriert. Um das Quellcode-Repository mit dem vollständigen und funktionierenden Beispiel herunterzuladen, besuchen Sie mein GitHub-Repo.
Abonnieren Sie den Blog, um weitere Leitfäden wie diesen zu erhalten. Sie erhalten jede Woche ein Update mit den neuesten Beiträgen des Mailgun-Teams. Und falls Sie Mailgun noch nicht für den Versand, den Empfang und das Tracking der E-Mails Ihrer Anwendung nutzen, melden Sie sich einfach unten an.