Apache Kafka und Datenstreaming mit Spark Streaming

Hallo, Habr! Heute werden wir ein System bauen, das mit Hilfe von Spark Streaming Streams von Apache Kafka verarbeitet und die Verarbeitungsresultate in der Cloud-Datenbank AWS RDS speichert.

Stellen wir uns vor, dass eine Kreditorganisation uns die Aufgabe stellt, eingehende Transaktionen "in Echtzeit" ĂŒber alle ihre Filialen hinweg zu verarbeiten. Dies kann zum Zweck der zeitnahen Berechnung offener Devisenpositionen fĂŒr die Schatzkammer, von Limits oder von finanziellen Ergebnissen aus den GeschĂ€ften usw. geschehen.

Wie man diesen Anwendungsfall ohne Magie und Zauberformeln umsetzt – lesen Sie weiter! Auf geht's!

Apache Kafka und Datenstreaming mit Spark Streaming
(Bildquelle)

EinfĂŒhrung

Zweifellos bietet die Verarbeitung großer Datenmengen in Echtzeit umfangreiche Möglichkeiten fĂŒr die Nutzung in modernen Systemen. Eine der beliebtesten Kombinationen dafĂŒr ist das Tandem aus Apache Kafka und Spark Streaming, bei dem Kafka einen Strom von eingehenden Nachrichtenpaketen erzeugt und Spark Streaming diese Pakete ĂŒber einen festgelegten Zeitintervall verarbeitet.

Um die Ausfallsicherheit der Anwendung zu erhöhen, werden wir Kontrollpunkte – Checkpoints – verwenden. Mithilfe dieses Mechanismus muss das Spark Streaming-Modul, wenn es verlorene Daten wiederherstellen muss, nur zur letzten Kontrollstelle zurĂŒckkehren und die Berechnungen von dort aus fortsetzen.

Architektur des entwickelten Systems

Apache Kafka und Datenstreaming mit Spark Streaming

Verwendete Komponenten:

  • Apache Kafka – ist ein verteiltes Nachrichtenaustauschsystem mit Veröffentlichung und Abonnement. Es eignet sich sowohl fĂŒr den autonomen als auch fĂŒr den Online-Verbrauch von Nachrichten. Um Datenverlust zu verhindern, werden Nachrichten in Kafka auf der Festplatte gespeichert und innerhalb des Clusters repliziert. Das Kafka-System ist auf dem Synchronisierungsdienst ZooKeeper aufgebaut;
  • Apache Spark Streaming — Spark-Komponente zur Verarbeitung von Streaming-Daten. Das Spark Streaming-Modul ist auf einer "Mikrobatch-Architektur" aufgebaut, bei der der Datenstrom als kontinuierliche Folge kleiner Datenpakete interpretiert wird. Spark Streaming empfĂ€ngt Daten aus verschiedenen Quellen und fasst sie in kleinen Paketen zusammen. Neue Pakete werden in regelmĂ€ĂŸigen AbstĂ€nden erstellt. Zu Beginn jedes Zeitintervalls wird ein neues Paket erstellt, und alle in diesem Intervall empfangenen Daten werden in dieses Paket aufgenommen. Am Ende des Intervalls wird das Wachstum des Pakets gestoppt. Die GrĂ¶ĂŸe des Intervalls wird durch einen Parameter bestimmt, der als Batch-Intervall bezeichnet wird;
  • Apache Spark SQL — verbindet relationale Verarbeitung mit funktionalem Programmieren in Spark. Unter strukturierten Daten versteht man Daten, die ein Schema haben, das heißt, einen einheitlichen Satz von Feldern fĂŒr alle DatensĂ€tze. Spark SQL unterstĂŒtzt die Eingabe aus vielen Quellen strukturierter Daten und kann dank der vorhandenen Schema-Informationen effizient nur die erforderlichen Felder der DatensĂ€tze extrahieren; es bietet auch API-Schnittstellen fĂŒr DataFrames;
  • AWS RDS — ist eine vergleichsweise kostengĂŒnstige Cloud-relationalen Datenbank, ein Webdienst, der die Einrichtung, den Betrieb und das Scaling vereinfacht und direkt von Amazon verwaltet wird.

Installation und Inbetriebnahme des Kafka-Servers

Vor der direkten Nutzung von Kafka muss sichergestellt werden, dass Java vorhanden ist, da die JVM verwendet wird:

sudo apt-get update 
sudo apt-get install default-jre
java -version

Erstellen wir einen neuen Benutzer fĂŒr die Arbeit mit Kafka:

sudo useradd kafka -m
sudo passwd kafka
sudo adduser kafka sudo

Laden wir nun das Distributionspaket von der offiziellen Apache Kafka-Website herunter:

wget -P /YOUR_PATH "http://apache-mirror.rbc.ru/pub/apache/kafka/2.2.0/kafka_2.12-2.2.0.tgz"

Entpacken wir das heruntergeladene Archiv:

tar -xvzf /YOUR_PATH/kafka_2.12-2.2.0.tgz
ln -s /YOUR_PATH/kafka_2.12-2.2.0 kafka

Der nÀchste Schritt ist optional. Tatsache ist, dass die Standardeinstellungen es nicht erlauben, alle Funktionen von Apache Kafka voll auszunutzen. Zum Beispiel das Löschen von Themen, Kategorien oder Gruppen, auf die Nachrichten veröffentlicht werden können. Um dies zu Àndern, bearbeiten wir die Konfigurationsdatei:

vim ~/kafka/config/server.properties

FĂŒgen Sie am Ende der Datei Folgendes hinzu:

delete.topic.enable = true

Vor dem Start des Kafka-Servers muss der ZooKeeper-Server gestartet werden. Wir verwenden ein Hilfsskript, das mit dem Kafka-Distributionspaket geliefert wird:

Cd ~/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties

Nachdem ZooKeeper erfolgreich gestartet wurde, starten wir in einem separaten Terminal den Kafka-Server:

bin/kafka-server-start.sh config/server.properties

Lass uns ein neues Thema mit dem Namen Transaction erstellen:

bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic transaction

Stellen wir sicher, dass das Thema mit der erforderlichen Anzahl an Partitionen und Replikationen erstellt wurde:

bin/kafka-topics.sh --describe --zookeeper localhost:2181

Apache Kafka und Datenstreaming mit Spark Streaming

Wir lassen die Testphasen fĂŒr den Produzenten und Konsumenten des neu erstellten Themas aus. Detaillierte Informationen dazu, wie man das Senden und Empfangen von Nachrichten testen kann, sind in der offiziellen Dokumentation zu finden — Senden Sie einige Nachrichten. Jetzt gehen wir zur Erstellung des Produzenten in Python unter Verwendung der KafkaProducer-API ĂŒber.

Schreiben des Produzenten

Der Produzent wird zufĂ€llige Daten generieren — 100 Nachrichten pro Sekunde. Unter zufĂ€lligen Daten verstehen wir ein Dictionary, das aus drei Feldern besteht:

  • Branch — der Name des Verkaufsstandortes der Kreditorganisation;
  • Currency — die WĂ€hrung des GeschĂ€fts;
  • Amount — der Betrag des GeschĂ€fts. Der Betrag wird eine positive Zahl sein, wenn es sich um den Kauf von WĂ€hrung durch die Bank handelt, und negativ — wenn es sich um den Verkauf handelt.

Der Code fĂŒr den Produzenten sieht folgendermaßen aus:

from numpy.random import choice, randint

def get_random_value():
    new_dict = {}

    branch_list = ['Kazan', 'SPB', 'Novosibirsk', 'Surgut']
    currency_list = ['RUB', 'USD', 'EUR', 'GBP']

    new_dict['branch'] = choice(branch_list)
    new_dict['currency'] = choice(currency_list)
    new_dict['amount'] = randint(-100, 100)

    return new_dict

Anschließend senden wir die Nachricht mit der Methode send an den Server, in das gewĂŒnschte Thema, im JSON-Format:

from kafka import KafkaProducer    

producer = KafkaProducer(bootstrap_servers=['localhost:9092'],
                             value_serializer=lambda x:dumps(x).encode('utf-8'),
                             compression_type='gzip')
my_topic = 'transaction'
data = get_random_value()

try:
    future = producer.send(topic = my_topic, value = data)
    record_metadata = future.get(timeout=10)
    
    print('--> Die Nachricht wurde an ein Thema gesendet: 
            {}, Partition: {}, Offset: {}' 
            .format(record_metadata.topic,
                record_metadata.partition,
                record_metadata.offset ))   
                             
except Exception as e:
    print('--> Es scheint ein Fehler aufgetreten zu sein: {}'.format(e))

finally:
    producer.flush()

Beim AusfĂŒhren des Skripts erhalten wir die folgenden Nachrichten im Terminal:

Apache Kafka und Datenstreaming mit Spark Streaming

Das bedeutet, dass alles wie gewĂŒnscht funktioniert — der Produzent generiert und sendet Nachrichten in das gewĂŒnschte Thema.
Der nÀchste Schritt wird die Installation von Spark und die Verarbeitung dieses Nachrichtenstreams sein.

Installation von Apache Spark

Apache Spark — ist eine universelle und hochleistungsfĂ€hige Cluster-Computing-Plattform.

In Bezug auf die Leistung ĂŒbertrifft Spark die beliebten Implementierungen des MapReduce-Modells und bietet gleichzeitig UnterstĂŒtzung fĂŒr eine breitere Palette von Berechnungsarten, einschließlich interaktiver Abfragen und Streaming-Verarbeitung. Geschwindigkeit spielt eine wichtige Rolle bei der Verarbeitung großer Datenmengen, da gerade die Geschwindigkeit es ermöglicht, interaktiv zu arbeiten, ohne Minuten oder Stunden auf das Warten zu verwenden. Eines der wichtigsten Merkmale von Spark, das diese hohe Geschwindigkeit ermöglicht, ist die FĂ€higkeit, Berechnungen im Speicher durchzufĂŒhren.

Dieses Framework ist in Scala geschrieben, daher ist es notwendig, es zunÀchst zu installieren:

sudo apt-get install scala

Laden Sie das Spark-Distribution von der offiziellen Website herunter:

wget "http://mirror.linux-ia64.org/apache/spark/spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz"

Entpacken Sie das Archiv:

sudo tar xvf spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz -C /usr/local/spark

FĂŒgen Sie den Pfad zu Spark in die Bash-Datei ein:

vim ~/.bashrc

FĂŒgen Sie im Editor die folgenden Zeilen hinzu:

SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH

FĂŒhren Sie den folgenden Befehl nach den Änderungen in bashrc aus:

source ~/.bashrc

Bereitstellung von AWS PostgreSQL

Es bleibt, die Datenbank bereitzustellen, in die wir die verarbeiteten Informationen aus den Streams hochladen werden. Dazu verwenden wir den AWS RDS-Service.

Gehen Sie zur AWS-Konsole —> AWS RDS —> Datenbanken —> Datenbank erstellen:
Apache Kafka und Datenstreaming mit Spark Streaming

WÀhlen Sie PostgreSQL und klicken Sie auf die SchaltflÀche Weiter:
Apache Kafka und Datenstreaming mit Spark Streaming

Da dieses Beispiel ausschließlich zu Bildungszwecken behandelt wird, verwenden wir einen kostenlosen Server "auf minimalen Anforderungen" (Free Tier):
Apache Kafka und Datenstreaming mit Spark Streaming

Aktivieren Sie dann das KontrollkĂ€stchen im Block Free Tier, und danach wird uns automatisch eine Instanz der Klasse t2.micro angeboten — auch wenn sie schwach ist, ist sie kostenlos und eignet sich gut fĂŒr unsere Aufgabe:
Apache Kafka und Datenstreaming mit Spark Streaming

Als nÀchstes kommen sehr wichtige Dinge: der Name der DB-Instanz, der Name des Master-Benutzers und sein Passwort. Wir nennen die Instanz: myHabrTest, Master-Benutzer: habr, Passwort: habr12345 und klicken Sie auf die SchaltflÀche Weiter:
Apache Kafka und Datenstreaming mit Spark Streaming

Auf der nĂ€chsten Seite befinden sich die Parameter, die fĂŒr die externe ZugĂ€nglichkeit unseres DB-Servers (Öffentliche ZugĂ€nglichkeit) und die VerfĂŒgbarkeit der Ports verantwortlich sind:

Apache Kafka und Datenstreaming mit Spark Streaming

Lassen Sie uns eine neue Einstellung fĂŒr die VPC-Sicherheitsgruppe erstellen, die es ermöglicht, von außen ĂŒber Port 5432 (PostgreSQL) auf unseren DB-Server zuzugreifen.
Wechseln Sie in einem separaten Browserfenster zur AWS-Konsole zum Abschnitt VPC-Dashboard —> Sicherheitsgruppen —> Sicherheitsgruppe erstellen:
Apache Kafka und Datenstreaming mit Spark Streaming

Wir geben den Namen fĂŒr die Sicherheitsgruppe an – PostgreSQL, die Beschreibung, an welche VPC diese Gruppe assoziiert werden soll und klicken auf die SchaltflĂ€che Erstellen:
Apache Kafka und Datenstreaming mit Spark Streaming

FĂŒr die neu erstellte Gruppe fĂŒllen wir die Eingangsregeln fĂŒr den Port 5432 aus, wie im Bild unten gezeigt. Den Port mĂŒssen wir manuell nicht angeben, sondern können PostgreSQL aus dem Dropdown-MenĂŒ Typ auswĂ€hlen.

Genauer gesagt bedeutet der Wert ::/0, dass eingehender Datenverkehr fĂŒr den Server aus der ganzen Welt verfĂŒgbar ist, was kanonisch nicht ganz korrekt ist, aber um das Beispiel zu erlĂ€utern, erlauben wir uns, diesen Ansatz zu verwenden:
Apache Kafka und Datenstreaming mit Spark Streaming

Wir kehren zur Browserseite zurĂŒck, wo wir "Erweiterte Einstellungen konfigurieren" geöffnet haben und wĂ€hlen im Abschnitt VPC-Sicherheitsgruppen –> Vorhandene VPC-Sicherheitsgruppen auswĂ€hlen –> PostgreSQL aus:
Apache Kafka und Datenstreaming mit Spark Streaming

Als nĂ€chstes geben wir im Abschnitt Datenbankoptionen –> Datenbankname –> den Namen an – habrDB.

Die anderen Parameter, mit Ausnahme der Deaktivierung der Sicherung (Backup-Aufbewahrungszeitraum – 0 Tage), Überwachung und Performance Insights, können wir auf die Standardwerte belassen. Wir klicken auf die SchaltflĂ€che Datenbank erstellen:
Apache Kafka und Datenstreaming mit Spark Streaming

Stream-Handler

Der abschließende Schritt besteht darin, eine Spark-Job zu entwickeln, die alle zwei Sekunden neue Daten verarbeitet, die von Kafka kommen, und das Ergebnis in die Datenbank eintrĂ€gt.

Wie bereits erwĂ€hnt, sind Checkpoints der Hauptmechanismus in Spark Streaming, der konfiguriert werden muss, um fehlertolerant zu sein. Wir werden Checkpoints verwenden, und im Falle eines Ausfalls des Verfahrens muss sich das Spark Streaming-Modul nur auf den letzten Checkpoint zurĂŒckziehen und die Berechnungen von dort aus fortsetzen.

Der Checkpoint kann aktiviert werden, indem ein Verzeichnis in einem fehlertoleranten, zuverlÀssigen Dateisystem (z. B. HDFS, S3 usw.) festgelegt wird, in dem die Checkpoint-Informationen gespeichert werden. Dies geschieht beispielsweise mit:

streamingContext.checkpoint(checkpointDirectory)

In unserem Beispiel verwenden wir den folgenden Ansatz, nĂ€mlich dass, wenn checkpointDirectory existiert, der Kontext aus den Daten des Checkpoints wiederhergestellt wird. Wenn das Verzeichnis nicht existiert (d. h. es wird zum ersten Mal ausgefĂŒhrt), wird die Funktion functionToCreateContext aufgerufen, um einen neuen Kontext zu erstellen und DStreams einzurichten:

from pyspark.streaming import StreamingContext

context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)

Wir erstellen ein DirectStream-Objekt, um eine Verbindung zum Thema "transaction" ĂŒber die Methode createDirectStream der Bibliothek KafkaUtils herzustellen:

von pyspark.streaming.kafka import KafkaUtils
    
sc = SparkContext(conf=conf)
ssc = StreamingContext(sc, 2)

broker_list = 'localhost:9092'
topic = 'transaction'

directKafkaStream = KafkaUtils.createDirectStream(ssc,
                                [topic],
                                {"metadata.broker.list": broker_list})

Wir parsen die eingehenden Daten im JSON-Format:

rowRdd = rdd.map(lambda w: Row(branch=w['branch'],
                                       currency=w['currency'],
                                       amount=w['amount']))
                                       
testDataFrame = spark.createDataFrame(rowRdd)
testDataFrame.createOrReplaceTempView("treasury_stream")

Mit Spark SQL fĂŒhren wir eine einfache Gruppierung durch und geben das Ergebnis in der Konsole aus:

select 
    from_unixtime(unix_timestamp()) as curr_time,
    t.branch                        as branch_name,
    t.currency                      as currency_code,
    sum(amount)                     as batch_value
from treasury_stream t
group by
    t.branch,
    t.currency

Wir erhalten den Text der Abfrage und fĂŒhren ihn ĂŒber Spark SQL aus:

sql_query = get_sql_query()
testResultDataFrame = spark.sql(sql_query)
testResultDataFrame.show(n=5)

Anschließend speichern wir die erhaltenen aggregierten Daten in einer Tabelle in AWS RDS. Um die Ergebnisse der Aggregation in einer Datenbanktabelle zu speichern, verwenden wir die Methode write des DataFrame-Objekts:

testResultDataFrame.write 
    .format("jdbc") 
    .mode("append") 
    .option("driver", 'org.postgresql.Driver') 
    .option("url","jdbc:postgresql://myhabrtest.ciny8bykwxeg.us-east-1.rds.amazonaws.com:5432/habrDB") 
    .option("dbtable", "transaction_flow") 
    .option("user", "habr") 
    .option("password", "habr12345") 
    .save()

Einige Worte zur Konfiguration der Verbindung zu AWS RDS. Den Benutzer und das Passwort dafĂŒr haben wir im Schritt „Bereitstellung von AWS PostgreSQL“ erstellt. Als URL des Datenbankservers sollten wir den Endpoint verwenden, der im Abschnitt Connectivity & security angezeigt wird:

Apache Kafka und Datenstreaming mit Spark Streaming

Um eine korrekte Verbindung zwischen Spark und Kafka herzustellen, sollten wir den Job ĂŒber spark-submit mit dem Artefakt starten spark-streaming-kafka-0-8_2.11. ZusĂ€tzlich verwenden wir auch ein Artefakt zur Interaktion mit der PostgreSQL-Datenbank, das wir ĂŒber —packages ĂŒbergeben werden.

FĂŒr die FlexibilitĂ€t des Skripts ziehen wir auch den Namen des Messaging-Servers und das Thema, aus dem wir Daten erhalten möchten, als Eingabeparameter in Betracht.

Nun ist es an der Zeit, das System zu starten und zu testen:

spark-submit 
--packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.0.2,
org.postgresql:postgresql:9.4.1207 
spark_job.py localhost:9092 transaction

Alles hat geklappt! Wie im Bild unten zu sehen ist, werden wĂ€hrend der AusfĂŒhrung der Anwendung alle 2 Sekunden neue Aggregationsergebnisse ausgegeben, da wir das Paketintervall auf 2 Sekunden eingestellt haben, als wir das StreamingContext-Objekt erstellt haben:

Apache Kafka und Datenstreaming mit Spark Streaming

Anschließend fĂŒhren wir eine einfache Abfrage an der Datenbank durch, um das Vorhandensein von EintrĂ€gen in der Tabelle zu ĂŒberprĂŒfen transaction_flow:

Apache Kafka und Datenstreaming mit Spark Streaming

Fazit

In diesem Artikel wurde ein Beispiel fĂŒr die Stream-Verarbeitung von Daten unter Verwendung von Spark Streaming in Kombination mit Apache Kafka und PostgreSQL betrachtet. Mit dem Anstieg der Datenmengen aus verschiedenen Quellen ist der praktische Wert von Spark Streaming fĂŒr die Erstellung von Stream-Anwendungen und Anwendungen, die in Echtzeit arbeiten, kaum zu ĂŒberschĂ€tzen.

Den vollstÀndigen Quellcode finden Sie in meinem Repository auf GitHub.

Ich freue mich darauf, diesen Artikel zu diskutieren. Ich warte auf Ihre Kommentare und hoffe auf konstruktive Kritik von allen interessierten Lesern.

Ich wĂŒnsche Ihnen viel Erfolg!

Ps. UrsprĂŒnglich war geplant, eine lokale PostgreSQL-Datenbank zu verwenden, aber in Anbetracht meiner Vorliebe fĂŒr AWS habe ich beschlossen, die Datenbank in die Cloud zu verlagern. In einem nĂ€chsten Artikel zu diesem Thema werde ich zeigen, wie man das oben beschriebene System vollstĂ€ndig in AWS mit Hilfe von AWS Kinesis und AWS EMR umsetzt. Bleiben Sie dran!

Quelle: habr.com

60GB SSD 8Gb DDR4