Email

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.
Bild für 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.