Cześć, Habr! Dziś zbudujemy system, który za pomocą Spark Streaming będzie przetwarzał strumienie wiadomości Apache Kafka i zapisywał wyniki przetwarzania w chmurze AWS RDS.
Załóżmy, że jakaś instytucja kredytowa stawia przed nami zadanie przetwarzania przychodzących transakcji "na żywo" we wszystkich swoich oddziałach. Może to być zrobione w celu szybkiego obliczenia otwartej pozycji walutowej dla skarbnika, limitów lub wyniku finansowego z transakcji itd.
Jak zrealizować ten przypadek bez użycia magii i czarodziejskich zaklęć — czytaj dalej! Zaczynamy!

Wprowadzenie
Oczywiście przetwarzanie dużej ilości danych w czasie rzeczywistym stwarza szerokie możliwości wykorzystania w nowoczesnych systemach. Jedną z najpopularniejszych kombinacji do tego celu jest tandem Apache Kafka i Spark Streaming, gdzie Kafka tworzy strumień pakietów przychodzących wiadomości, a Spark Streaming przetwarza te pakiety w określonych odstępach czasu.
Aby zwiększyć odporność aplikacji na awarie, użyjemy punktów kontrolnych — checkpointów. Dzięki temu mechanizmowi, gdy moduł Spark Streaming będzie musiał przywrócić utracone dane, wystarczy, że wróci do ostatniego punktu kontrolnego i wznowi obliczenia z tego miejsca.
Architektura opracowywanego systemu

Używane komponenty:
- — to rozproszony system wymiany wiadomości z publikacją i subskrypcją. Nadaje się zarówno do autonomicznego, jak i do online'owego przetwarzania wiadomości. Aby zapobiec utracie danych, wiadomości Kafka są zapisywane na dysku i replikowane w obrębie klastra. System Kafka zbudowany jest na bazie usługi synchronizacji ZooKeeper;
- — komponent Spark do przetwarzania strumieni danych. Moduł Spark Streaming zbudowany jest w oparciu o architekturę „mikropakietową”, w której strumień danych interpretowany jest jako nieprzerwany ciąg małych pakietów danych. Spark Streaming przyjmuje dane z różnych źródeł i łączy je w małe pakiety. Nowe pakiety tworzone są w regularnych odstępach czasu. Na początku każdego przedziału czasu tworzony jest nowy pakiet, a wszelkie dane, które wpłynęły w tym czasie, są do niego włączane. Na końcu przedziału wzrost pakietu jest wstrzymywany. Rozmiar przedziału definiowany jest przez parametr zwany przedziałem pakietowania;
- — łączy relacyjne przetwarzanie z funkcjonalnym programowaniem Spark. Pod danymi strukturalnymi rozumie się dane o schemacie, czyli jednolity zestaw pól dla wszystkich zapisów. Spark SQL obsługuje wprowadzanie danych z różnych źródeł danych strukturalnych i, dzięki posiadaniu informacji o schemacie, może efektywnie wydobywać tylko niezbędne pola zapisów oraz zapewnia interfejsy API DataFrame;
- — to stosunkowo niedroga chmurowa relacyjna baza danych, usługa internetowa, która upraszcza konfigurację, obsługę i skalowanie, zarządzana bezpośrednio przez Amazon.
Instalacja i uruchomienie serwera Kafka
Przed bezpośrednim użyciem Kafka, należy upewnić się, że Java jest zainstalowana, ponieważ do pracy używana jest JVM:
sudo apt-get update
sudo apt-get install default-jre
java -version
Utworzymy nowego użytkownika do pracy z Kafka:
sudo useradd kafka -m
sudo passwd kafka
sudo adduser kafka sudo
Następnie pobieramy dystrybucję z oficjalnej strony Apache Kafka:
wget -P /YOUR_PATH "http://apache-mirror.rbc.ru/pub/apache/kafka/2.2.0/kafka_2.12-2.2.0.tgz"Rozpakowujemy pobrany archiwum:
tar -xvzf /YOUR_PATH/kafka_2.12-2.2.0.tgz
ln -s /YOUR_PATH/kafka_2.12-2.2.0 kafka
Następny krok jest opcjonalny. Rzecz w tym, że domyślne ustawienia nie pozwalają w pełni wykorzystać wszystkich możliwości Apache Kafka. Na przykład, usuwanie tematów, kategorii, grup, do których mogą być publikowane wiadomości. Aby to zmienić, edytujemy plik konfiguracyjny:
vim ~/kafka/config/server.propertiesDodaj na końcu pliku następujące:
delete.topic.enable = truePrzed uruchomieniem serwera Kafka, należy uruchomić serwer ZooKeeper. Skorzystamy z pomocniczego skryptu, który jest dostarczany z dystrybucją Kafka:
Cd ~/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties
Po pomyślnym uruchomieniu ZooKeeper, w osobnym terminalu uruchamiamy serwer Kafka:
bin/kafka-server-start.sh config/server.propertiesUtwórzmy nowy temat o nazwie Transaction:
bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic transactionUpewnijmy się, że temat został utworzony z odpowiednią liczbą partycji i replikacją:
bin/kafka-topics.sh --describe --zookeeper localhost:2181 
Pomińmy momenty testowania producenta i konsumenta dla nowo utworzonego tematu. Szczegółowe informacje o tym, jak testować wysyłanie i odbieranie wiadomości, znajdują się w oficjalnej dokumentacji — . A my przechodzimy do pisania producenta w Pythonie z wykorzystaniem API KafkaProducer.
Pisanie producenta
Producent będzie generował losowe dane — po 100 wiadomości na każdą sekundę. Pod losowymi danymi rozumiemy słownik składający się z trzech pól:
- Oddział — nazwa punktu sprzedaży instytucji kredytowej;
- Currency — waluta transakcji;
- Amount — kwota transakcji. Kwota będzie liczbą dodatnią, jeśli jest to zakup waluty przez Bank, a ujemną — jeśli sprzedaż.
Kod dla producenta wygląda następująco:
from numpy.random import choice, randint
def get_random_value():
new_dict = {}
branch_list = ['Kazań', 'SPB', 'Nowosybirsk', '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
Następnie, używając metody send, wysyłamy wiadomość na serwer, do potrzebnego nam tematu, w formacie JSON:
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('--> Wiadomość została wysłana do tematu:
{}, partycja: {}, offset: {}'
.format(record_metadata.topic,
record_metadata.partition,
record_metadata.offset ))
except Exception as e:
print('--> Wygląda na to, że wystąpił błąd: {}'.format(e))
finally:
producer.flush()
Podczas uruchamiania skryptu otrzymujemy w terminalu następujące komunikaty:

To oznacza, że wszystko działa tak, jak chcieliśmy — producent generuje i wysyła wiadomości do potrzebnego nam tematu.
Kolejnym krokiem będzie zainstalowanie Sparka i przetwarzanie tego strumienia wiadomości.
Instalacja Apache Spark
Apache Spark — to uniwersalna i wydajna platforma obliczeniowa w chmurze.
Pod względem wydajności Spark przewyższa popularne realizacje modelu MapReduce, równocześnie zapewniając wsparcie dla szerszego zakresu typów obliczeń, w tym zapytań interaktywnych i przetwarzania strumieniowego. Szybkość odgrywa kluczową rolę w przetwarzaniu dużych zbiorów danych, ponieważ to właśnie szybkość pozwala na pracę w trybie interaktywnym, bez czekania minutami czy godzinami. Jedną z najważniejszych zalet Sparka, która zapewnia tak wysoką szybkość, jest zdolność wykonywania obliczeń w pamięci.
Niniejszy framework jest napisany w języku Scala, dlatego najpierw należy go zainstalować:
sudo apt-get install scalaPobieramy dystrybucję Sparka z oficjalnej strony:
wget "http://mirror.linux-ia64.org/apache/spark/spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz"Rozpakowujemy archiwum:
sudo tar xvf spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz -C /usr/local/sparkDodajemy ścieżkę do Sparka w pliku bash:
vim ~/.bashrcDodajemy następujące linie przez edytor:
SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH
Wykonujemy poniższe polecenie po wprowadzeniu zmian w bashrc:
source ~/.bashrcRozwój AWS PostgreSQL
Teraz pozostało rozwinąć bazę danych, do której będziemy zrzucać przetworzone informacje ze strumieni. Do tego użyjemy usługi AWS RDS.
Przechodzimy do konsoli AWS —> AWS RDS —> Bazy danych —> Utwórz bazę danych:

Wybieramy PostgreSQL i klikamy przycisk Dalej:

Ponieważ ten przykład jest omawiany wyłącznie w celach edukacyjnych, będziemy korzystać z bezpłatnego serwera „minimalnego” (Free Tier):

Następnie zaznaczamy pole w bloku Free Tier, a po tym automatycznie zaproponowany zostanie nam instancja klasy t2.micro — choć słaba, to darmowa i w zupełności wystarczająca do naszych zadań:

Następnie mamy bardzo ważne rzeczy: nazwa instancji Bazy Danych, nazwa użytkownika głównego i jego hasło. Nazwijmy instancję: myHabrTest, użytkownik główny: habr, hasło: habr12345 i klikamy przycisk Dalej:

Na następnej stronie znajdują się parametry dotyczące dostępności naszego serwera Bazy Danych z zewnątrz (Public accessibility) oraz dostępność portów:

Utwórzmy nową konfigurację dla grupy zabezpieczeń VPC, która pozwoli na zewnętrzny dostęp do naszego serwera Bazy Danych przez port 5432 (PostgreSQL).
Przechodzimy w osobnym oknie przeglądarki do konsoli AWS w sekcji VPC Dashboard —> Grupy zabezpieczeń —> Utwórz grupę zabezpieczeń:

Nadajemy nazwę dla grupy zabezpieczeń — PostgreSQL, opis, określamy, z którą VPC ta grupa powinna być skojarzona i klikamy przycisk Utwórz:

Wypełniamy dla świeżo utworzonej grupy zasady przychodzące dla portu 5432, jak pokazano na poniższym obrazku. Port możemy również wybrać z rozwijanego listy Typ, a nie wpisywać ręcznie.
Mówiąc ściślej, wartość ::/0 oznacza dostępność przychodzącego ruchu dla serwera z całego świata, co kanonicznie nie jest do końca poprawne, ale dla celów tego przykładu pozwolimy sobie na takie podejście:

Wracamy do strony przeglądarki, gdzie mamy otwarte „Skonfiguruj zaawansowane ustawienia” i wybieramy w sekcji grupy zabezpieczeń VPC —> Wybierz istniejące grupy zabezpieczeń VPC —> PostgreSQL:

Następnie, w sekcji Opcje bazy danych —> Nazwa bazy danych —> nadajemy nazwę — habrDB.
Pozostałe parametry, z wyjątkiem wyłączenia tworzenia kopii zapasowych (czas przechowywania kopii zapasowych — 0 dni), monitorowania i Analiz Wydajności, możemy pozostawić domyślne. Klikamy przycisk Utwórz bazę danych:

Obsługiwacz strumieni
Końcowym etapem będzie opracowanie zadania Spark, które co dwie sekundy będzie przetwarzać nowe dane przychodzące z Kafki i zapisywać wyniki w bazie danych.
Jak wspomniano powyżej, punkty kontrolne (checkpoints) są głównym mechanizmem w SparkStreaming, który musi być skonfigurowany, aby zapewnić odporność na awarie. Będziemy używać punktów kontrolnych, a w przypadku awarii procedury, moduł Spark Streaming do odzyskania utraconych danych będzie musiał jedynie wrócić do ostatniego punktu kontrolnego i wznowić obliczenia od niego.
Punkt kontrolny można włączyć, ustawiając katalog w odpornej na awarie, niezawodnej systemie plików (np. HDFS, S3 itd.), w którym zostaną zapisane informacje o punkcie kontrolnym. Robi się to za pomocą na przykład:
streamingContext.checkpoint(checkpointDirectory)W naszym przykładzie zastosujemy następujące podejście, a mianowicie, jeśli checkpointDirectory istnieje, to kontekst będzie odtworzony z danych punktu kontrolnego. Jeśli katalog nie istnieje (tj. wykonujemy to po raz pierwszy), wywoływana jest funkcja functionToCreateContext do utworzenia nowego kontekstu i skonfigurowania DStreams:
from pyspark.streaming import StreamingContext
context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)
Tworzymy obiekt DirectStream w celu połączenia się z tematem „transaction” za pomocą metody createDirectStream biblioteki KafkaUtils:
from 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})
Parsimy dane wejściowe w formacie JSON:
rowRdd = rdd.map(lambda w: Row(branch=w['branch'],
currency=w['currency'],
amount=w['amount']))
testDataFrame = spark.createDataFrame(rowRdd)
testDataFrame.createOrReplaceTempView("treasury_stream")
Używając Spark SQL, wykonujemy prostą grupowanie i wyświetlamy wyniki w konsoli:
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
Pobieramy tekst zapytania i uruchamiamy go przez Spark SQL:
sql_query = get_sql_query()
testResultDataFrame = spark.sql(sql_query)
testResultDataFrame.show(n=5)
A następnie zapisujemy uzyskane zsumowane dane do tabeli w AWS RDS. Aby zapisać wyniki agregacji w tabeli bazy danych, użyjemy metody write obiektu DataFrame:
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()
Kilka słów na temat konfiguracji połączenia z AWS RDS. Użytkownika i hasło do niego stworzyliśmy na etapie „Wdrażania AWS PostgreSQL”. Jako url serwera baz danych należy użyć Endpointu, który wyświetla się w sekcji Connectivity & security:
W celu poprawnej integracji Spark i Kafka, należy uruchomić zadanie za pomocą spark-submit z wykorzystaniem artefaktu spark-streaming-kafka-0-8_2.11. Dodatkowo zastosujemy również artefakt do interakcji z bazą danych PostgreSQL, które przekażemy przez —packages.
Aby zwiększyć elastyczność skryptu, wyodrębnimy również nazwę serwera wiadomości i temat, z którego chcemy uzyskać dane.
Nadszedł czas na uruchomienie i sprawdzenie działania systemu:
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
Udało się! Jak widać na poniższym obrazku — podczas pracy aplikacji nowe wyniki agregacji są wyświetlane co 2 sekundy, ponieważ ustawiliśmy interwał pakietowania na 2 sekundy podczas tworzenia obiektu StreamingContext:

Następnie wykonujemy proste zapytanie do bazy danych, aby sprawdzić obecność rekordów w tabeli transaction_flow:

Podsumowanie
W tym artykule omówiono przykład przetwarzania strumieniowego danych z wykorzystaniem Spark Streaming w połączeniu z Apache Kafka i PostgreSQL. Wraz ze wzrostem ilości danych z różnych źródeł trudno przecenić praktyczną wartość Spark Streaming w tworzeniu aplikacji strumieniowych i aplikacji działających w czasie rzeczywistym.
Pełny kod źródłowy można znaleźć w moim repozytorium na .
Chętnie omówię ten artykuł, czekam na Wasze komentarze oraz liczę na konstruktywną krytykę wszystkich niezainteresowanych czytelników.
Życzę powodzenia!
Ps. Początkowo planowano użycie lokalnej bazy danych PostgreSQL, ale biorąc pod uwagę moją miłość do AWS, postanowiłem przenieść bazę danych do chmury. W następnym artykule na ten temat pokażę, jak zrealizować w całości powyżej opisaną system w AWS za pomocą AWS Kinesis i AWS EMR. Śledźcie nasze nowości!
Źródło: habr.com

