Përshëndetje, Habr! Sot do të ndërtojmë një sistem që do të përpunojë rrjedhat e mesazheve Apache Kafka me ndihmën e Spark Streaming dhe do të regjistrojë rezultatet e përpunimit në databazën në re AWS RDS.
Të imagjinojmë se një organizatë krediti na jep detyrën e përpunimit të transaksioneve hyrëse "në kohë reale" për të gjitha degët e saj. Kjo mund të bëhet me qëllim të llogaritjes së menjëhershme të pozitat e monedhës të hapura për thesarin, kufijtë ose rezultatet financiare nga transaksionet, etj.
Si ta realizojmĂ« kĂ«tĂ« rast pa pĂ«rdorur magji dhe formula tĂ« çuditshme â lexoni mĂ« poshtĂ«! Le tĂ« fillojmĂ«!

Hyrje
Pa dyshim, pĂ«rpunimi i njĂ« sĂ«rĂ« tĂ« madhe tĂ« dhĂ«nash nĂ« kohĂ« reale ofron mundĂ«si tĂ« gjera pĂ«r tâu pĂ«rdorur nĂ« sistemet moderne. NjĂ« nga kombinimet mĂ« tĂ« njohura pĂ«r kĂ«tĂ« Ă«shtĂ« tandemi Apache Kafka dhe Spark Streaming, ku Kafka krijon njĂ« lum fluksi mesazhesh tĂ« ardhura, ndĂ«rsa Spark Streaming i pĂ«rpunon kĂ«to paketa pĂ«rmes njĂ« intervali tĂ« caktuar tĂ« kohĂ«s.
PĂ«r tĂ« rritur qĂ«ndrueshmĂ«rinĂ« e aplikacionit, do tĂ« pĂ«rdorim pikĂ« kontrolli â checkpoints. Me ndihmĂ«n e kĂ«tij mekanizmi, kur moduli Spark Streaming ka nevojĂ« tĂ« rikuperojĂ« tĂ« dhĂ«nat e humbura, i nevojitet vetĂ«m tĂ« kthehet te pika e fundit e kontrollit dhe tĂ« vazhdojĂ« llogaritjet nga ajo.
Arkitektura e sistemit në zhvillim

Komponentët e përdorur:
- â Ă«shtĂ« njĂ« sistem i shpĂ«rndarĂ« i dĂ«rgimit tĂ« mesazheve me publikim dhe abonim. PĂ«rshtatet si pĂ«r konsum autonom ashtu edhe pĂ«r konsum online tĂ« mesazheve. PĂ«r tĂ« parandaluar humbjen e tĂ« dhĂ«nave, mesazhet e Kafka ruhen nĂ« disk dhe replikohen brenda klasterit. Sistemi Kafka Ă«shtĂ« ndĂ«rtuar mbi shĂ«rbimin e sinkronicitetit ZooKeeper;
- â komponenti Spark pĂ«r pĂ«rpunimin e tĂ« dhĂ«nave nĂ« kohĂ« reale. Moduli Spark Streaming Ă«shtĂ« ndĂ«rtuar me pĂ«rdorimin e arkitekturĂ«s "mikropaketash" (micro-batch architecture), kur rrjedha e tĂ« dhĂ«nave interpretohet si njĂ« vazhdimĂ«si e pandĂ«rprerĂ« e paketimeve tĂ« vogla tĂ« dhĂ«nash. Spark Streaming merr tĂ« dhĂ«na nga burime tĂ« ndryshme dhe i bashkon ato nĂ« paketa tĂ« vogla. Paketat e reja krijohen pĂ«rmes intervaleve tĂ« rregullta tĂ« kohĂ«s. NĂ« fillim tĂ« çdo intervali kohor krijohet njĂ« paketĂ« e re, dhe çdo tĂ« dhĂ«nĂ« e pranuar gjatĂ« kĂ«tij intervali pĂ«rfshihet nĂ« paketĂ«. NĂ« fund tĂ« intervalit, rritja e paketĂ«s ndalet. MadhĂ«sia e intervalit pĂ«rcaktohet nga njĂ« parametĂ«r i quajtur intervali i paketimit (batch interval);
- â bashkon pĂ«rpunimin relacional me programimin funksional nĂ« Spark. TĂ« dhĂ«nat e strukturuara referojnĂ« nĂ« tĂ« dhĂ«na qĂ« kanĂ« njĂ« skemĂ«, domethĂ«nĂ« njĂ« grup tĂ« vetĂ«m fushash pĂ«r tĂ« gjitha regjistrimet. Spark SQL mbĂ«shtet hyrjen nga shumĂ« burime tĂ« tĂ« dhĂ«nave tĂ« strukturuara, dhe falĂ« informacionit mbi skemĂ«n, ai mund tĂ« nxjerrĂ« nĂ« mĂ«nyrĂ« efikase vetĂ«m fushat e nevojshme tĂ« regjistrimeve, si dhe siguron interface API tĂ« DataFrame;
- â Ă«shtĂ« njĂ« bazĂ« tĂ« dhĂ«nash relacionalĂ« cloud relativisht e lirĂ«, njĂ« shĂ«rbim nĂ« internet qĂ« thjeshtĂ«son konfigurimin, operimin dhe shkallĂ«zimin, i menaxhuar drejtpĂ«rdrejt nga Amazon.
Instalimi dhe nisja e serverit Kafka
Para përdorimit të drejtpërdrejtë të Kafka, është e nevojshme të sigurohemi për praninë e Java, pasi është e nevojshme JVM:
sudo apt-get update
sudo apt-get install default-jre
java -version
Të krijojmë një përdorues të ri për punën me Kafka:
sudo useradd kafka -m
sudo passwd kafka
sudo adduser kafka sudo
Më pas shkarkojmë distributorin nga faqja zyrtare e Apache Kafka:
wget -P /YOUR_PATH "http://apache-mirror.rbc.ru/pub/apache/kafka/2.2.0/kafka_2.12-2.2.0.tgz"Zbërthejmë arkivën e shkarkuar:
tar -xvzf /YOUR_PATH/kafka_2.12-2.2.0.tgz
ln -s /YOUR_PATH/kafka_2.12-2.2.0 kafka
Hapi i ardhshëm është opsional. Kjo është për shkak se konfigurimet e parazgjedhura nuk lejojnë shfrytëzimin e plotë të të gjitha mundësive të Apache Kafka. Për shembull, fshirjen e një teme, kategorie, grupi, në të cilat mund të publikohen mesazhe. Për ta ndryshuar këtë, redaktoni skedarin e konfigurimit:
vim ~/kafka/config/server.propertiesShtoni në fund të skedarit sa vijon:
delete.topic.enable = truePara të nisni serverin Kafka, duhet të filloni serverin ZooKeeper, do të përdorim një skenar ndihmës që vjen me distribucionin e Kafka:
Cd ~/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties
Pasi ZooKeeper të ketë filluar me sukses, në një terminal tjetër nisim serverin Kafka:
bin/kafka-server-start.sh config/server.propertiesDo të krijojmë një topic të ri me emrin Transaction:
bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic transactionDo të sigurohemi që topic-i me numrin e duhur të pjesëve dhe replikimit të jetë krijuar:
bin/kafka-topics.sh --describe --zookeeper localhost:2181 
Do ta kalojmĂ« momentin e testimit tĂ« produesit dhe konsumatorit pĂ«r topic-in e sapokrijuar. MĂ« shumĂ« informacion mbi si tĂ« testoni dĂ«rgimin dhe pranimin e mesazheve mund tĂ« gjeni nĂ« dokumentacionin zyrtar â . NĂ« tĂ« vĂ«rtetĂ«, ne kalojmĂ« te shkruarja e produesit nĂ« Python duke pĂ«rdorur API-nĂ« KafkaProducer.
Shkrimi i produesit
Produesit do tĂ« gjenerojĂ« tĂ« dhĂ«na tĂ« rastit â 100 mesazhe çdo sekondĂ«. TĂ« dhĂ«nat e rastit do tĂ« nĂ«nkuptojnĂ« njĂ« fjalor tĂ« pĂ«rbĂ«rĂ« nga tre fusha:
- DegĂ« â emri i pikĂ«s sĂ« shitjes sĂ« organizatĂ«s kreditore;
- Currency â valuta e transaksionit;
- Amount â shuma e transaksionit. Shuma do tĂ« jetĂ« njĂ« numĂ«r pozitiv nĂ«se Ă«shtĂ« njĂ« blerje valute nga Banka, dhe negativ nĂ«se Ă«shtĂ« shitje.
Kodi për produesin duket si më poshtë:
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
Më pas, duke përdorur metodën send, dërgojmë mesazhin në server, në topic-in e nevojshëm, në formatin 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('--> Mesazhi është dërguar në një topic:
{}, pjesa: {}, offset: {}'
.format(record_metadata.topic,
record_metadata.partition,
record_metadata.offset ))
except Exception as e:
print('--> Duket se ndodhi një gabim: {}'.format(e))
finally:
producer.flush()
Kur ekzekutojmë skriptin, marrim në terminal këto mesazhe:

Kjo do tĂ« thotĂ« qĂ« gjithçka funksionon siç doja â produesit gjeneron dhe dĂ«rgon mesazhe nĂ« topic-in e dĂ«shiruar.
Hapi tjetër do të jetë instalimi i Spark dhe përpunimi i këtij fluksi mesazhesh.
Instalimi i Apache Spark
Apache Spark është një platformë universale dhe me performancë të lartë për llogaritjet në klaster.
Në performancë, Spark e kalon implementimet më të njohura të modelit MapReduce, duke ofruar gjithashtu mbështetje për një gamë më të gjerë llojesh llogaritjesh, duke përfshirë kërkesa interaktive dhe përpunim në kohë reale. Shpejtësia luan një rol të rëndësishëm në përpunimin e volumit të madh të të dhënave, pasi shpejtësia lejon të punosh në mënyrë interaktive, pa pritur minuta ose orë. Një nga përfitimet kryesore të Spark që garanton këtë shpejtësi të lartë është capabiliteti për të kryer llogaritjet në memorie.
Ky kornizë është shkruar në Scala, kështu që duhet të instalojmë atë së pari:
sudo apt-get install scalaShkarko nga faqja zyrtare shpërndarjen e Spark:
wget "http://mirror.linux-ia64.org/apache/spark/spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz"Shkrijmë arkivin:
sudo tar xvf spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz -C /usr/local/sparkShtojmë rrugën për në Spark në skedarin bash:
vim ~/ .bashrcShtojmë përmes editorit këto rreshta:
SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH
Ekzekutojmë komandën më poshtë pas editimit të bashrc:
source ~/ .bashrcZhvillimi i AWS PostgreSQL
Kemi mbetur për të zhvilluar bazën e të dhënave, ku do të ngarkojmë informacionin e përpunuar nga rrjedhat. Për këtë, do të përdorim shërbimin AWS RDS.
Hyr në konsolën AWS -> AWS RDS -> Databases -> Krijo bazën e të dhënave:

Zgjidh PostgreSQL dhe klikoni butonin Next:

Duke marrë parasysh se ky shembull shqyrtohet ekskluzivisht për qëllime edukative, do të përdorim një server falas "në minimalet" (Free Tier):

MĂ« pas, vendosim njĂ« cek nĂ« bllokun Free Tier, dhe pas kĂ«saj automatikisht do tĂ« na ofrohet njĂ« instancĂ« e klasĂ«s t2.micro â edhe pse e dobĂ«t, Ă«shtĂ« falas dhe pĂ«rfundimisht do tĂ« pĂ«rshtatet pĂ«r nevojĂ«n tonĂ«:

Pasojnë disa gjëra shumë të rëndësishme: emri i instancës së DB, emri i përdoruesit kryesor dhe fjalëkalimi i tij. Do ta quajmë instancën: myHabrTest, përdoruesit kryesor: habr, fjalëkalimi: habr12345 dhe klikojmë butonin Next:

Në faqen tjetër ndodhen parametrat që përgjigjen për aksesin e serverit tonë DB nga jashtë (Public accessibility) dhe aksesin e porteve:

Le të krijojmë një konfigurim të ri për grupin e sigurisë VPC, që do të lejojë qasjen nga jashtë në serverin tonë DB përmes portit 5432 (PostgreSQL).
Të kalojmë në një dritare të veçantë të shfletuesit në konsolën AWS në seksionin VPC Dashboard -> Security Groups -> Krijo grupin e sigurisë:

Caktojmë emrin për grupin e sigurisë - PostgreSQL, përshkrimin, tregojmë se në cilën VPC duhet të asociohet ky grup dhe klikojmë butonin Create:

Plotsojmë për grupin e sapokrijuar rregullat hyrëse për portin 5432, siç tregohet në imazhin më poshtë. Portin mund ta lëmë bosh dhe të zgjedhim PostgreSQL nga lista e rënëse Type.
Rreptësisht, vlera ::/0 nënkupton disponueshmërinë e trafikut hyrës për serverin nga e gjithë bota, që nuk është krejtësisht e saktë, por për qëllim ilustrimi do të lejojmë një qasje të tillë:

Kthehemi në faqen e shfletuesit, ku kemi hapur «Configure advanced settings» dhe zgjedhim në seksionin VPC security groups -> Choose existing VPC security groups -> PostgreSQL:

Më pas, në seksionin Database options -> Database name -> caktojmë emrin - habrDB.
Parametrat e tjerë, përveç ndoshta çaktivizimit të kopjes rezervë (backup retention period - 0 days), monitorimit dhe Performance Insights, mund t'i lëmë në parametrat e paracaktuar. Klikojmë butonin Create database:

Menaxheri i flukseve
Hapi përfundimtar do të jetë zhvillimi i një Spark-job, i cili çdo dy sekonda do të procesojë të dhëna të reja që vijnë nga Kafka dhe do të ruajë rezultatin në bazën e të dhënave.
Siç u përmend më parë, pikëmbështetjet (checkpoints) janë mekanizmi kryesor në SparkStreaming, i cili duhet të konfigurohet për të siguruar qëndrueshmëri. Ne do të përdorim pikëmbështetjet dhe, nëse ndodh një rënie e procedurës, moduli Spark Streaming për rikuperimin e të dhënave të humbura do të duhet vetëm të kthehet në pikën e fundit të mbështetjes dhe të rinisë llogaritjet prej saj.
Pika e mbështetjes mund të aktivizohet duke vendosur katalogun në një sistem skedarësh të qëndrueshëm dhe të besueshëm (p.sh., HDFS, S3 etj.), në të cilin do të ruhet informacioni i pikëmbështetjes. Kjo bëhet, për shembull:
streamingContext.checkpoint(checkpointDirectory)Në shembullin tonë do të përdorim qasjen e mëposhtme, dmth, nëse checkpointDirectory ekziston, konteksti do të rishpallet nga të dhënat e pikëmbështetjes. Nëse katalogu nuk ekziston (dmth po ekzekutohet për herë të parë), atëherë thirret funksioni functionToCreateContext për krijimin e një konteksti të ri dhe për konfigurimin e DStreams:
from pyspark.streaming import StreamingContext
context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)
Krijojmë një objekt DirectStream me qëllim lidhjen me temën «transaction» nëpërmjet metodës createDirectStream të bibliotekës KafkaUtils:
nga 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})
Po kemi të dhënat hyrëse në formatin JSON:
rowRdd = rdd.map(lambda w: Row(branch=w['branch'],
currency=w['currency'],
amount=w['amount']))
testDataFrame = spark.createDataFrame(rowRdd)
testDataFrame.createOrReplaceTempView("treasury_stream")
Duke përdorur Spark SQL, bëjmë një grupim të thjeshtë dhe e nxjerrim rezultatin në konsolë:
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
Marrja e tekstit të kërkesës dhe ekzekutimi i saj përmes Spark SQL:
sql_query = get_sql_query()
testResultDataFrame = spark.sql(sql_query)
testResultDataFrame.show(n=5)
Dhe pastaj ruajmë të dhënat e agreguara në një tabelë në AWS RDS. Për të ruajtur rezultatet e agregatës në tabelën e bazës së të dhënave, do të përdorim metodën write të objektit 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()
Disa fjalë për konfigurimin e lidhjes me AWS RDS. Përdoruesin dhe fjalëkalimin e tij i kemi krijuar në hapin "Zhvillimi i AWS PostgreSQL". Si url për serverin e bazës së të dhënave duhet të përdorim Endpoint, i cili shfaqet në seksionin Connectivity & security:
PĂ«r njĂ« lidhje tĂ« saktĂ« mes Spark dhe Kafka, duhet tĂ« ekzekutojmĂ« punĂ«n pĂ«rmes spark-submit duke pĂ«rdorur artefaktin spark-streaming-kafka-0-8_2.11. PĂ«r mĂ« tepĂ«r, do tĂ« pĂ«rdorim gjithashtu artefaktin pĂ«r ndĂ«rveprimin me bazĂ«n e tĂ« dhĂ«nave PostgreSQL, tĂ« cilat do t'i kalojmĂ« pĂ«rmes âpackages.
Për fleksibilitetin e skriptit, do të nxjerrim si parametra hyrës gjithashtu emrin e serverit të mesazheve dhe temën nga e cila dëshirojmë të marrim të dhëna.
Kështu, erdhi koha të nisemi dhe të verifikojmë funksionimin e sistemit:
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
E gjithĂ« ka shkuar mirĂ«! Siç shihet nĂ« imazhin mĂ« poshtĂ« â gjatĂ« procesit tĂ« punĂ«s sĂ« aplikacionit rezultatet e reja tĂ« agregimit shfaqen çdo 2 sekonda, sepse vendosĂ«m intervalin e paketimit tĂ« barabartĂ« me 2 sekonda kur krijuam objektin StreamingContext:

Më pas, bëjmë një kërkesë të thjeshtë në bazën e të dhënave për të kontrolluar praninë e regjistrimeve në tabelën transaction_flow:

Përfundim
Në këtë artikull u shqyrtua një shembull i përpunimit të informacionit në rrjedhë me përdorimin e Spark Streaming në lidhje me Apache Kafka dhe PostgreSQL. Me rritjen e volumit të të dhënave nga burime të ndryshme, vlera praktike e Spark Streaming për krijimin e aplikacioneve në rrjedhë dhe aplikacioneve që veprojnë në kohë reale është e vështirë të nënvizohet.
Kodi burimor i plotë mund ta gjeni në depot time në .
Me kënaqësi jam i gatshëm të diskutoj për këtë artikull, pres komentet tuaja, si dhe shpresoj për kritikë konstruktive nga të gjithë lexuesit e përkushtuar.
Ju uroj suksese!
Ps. Fillimisht ishte planifikuar të përdorej një DB lokale PostgreSQL, por duke marrë parasysh dashurinë time për AWS, vendosa ta çoj bazën e të dhënave në cloud. Në artikullin e ardhshëm në këtë temë do të tregoj si ta zbatoj sistemin e përshkruar më sipër në AWS duke përdorur AWS Kinesis dhe AWS EMR. Qëndroni të lidhur për lajme!
Burimi: habr.com

