Përshëndetje, Habr! Sot do të ndërtojmë një sistem që do të përpunojë flukse mesazhesh Apache Kafka duke përdorur Spark Streaming dhe do të regjistrojë rezultatin e përpunimit në një bazë të dhënash në AWS RDS.
Imagjinoni se një institucion kreditor na kërkon të përpunojmë transaksionet në hyrje "në kohë reale" për të gjithë filialet e tij. Kjo mund të bëhet me qëllim llogaritjen e shpejtë të pozicionit të hapur në monedhë për thesar, kufij ose rezultat financiar për transaksionet etj.
Si ta realizojmë këtë rast pa përdorur magji dhe fjalë magjike - lexoni më poshtë! Le të fillojmë!

Hyrje
Pa diskutim, përpunimi i një sasie të madhe të dhënash në kohë reale ofron mundësi të mëdha për përdorim në sistemet moderne. Një nga kombinimet më të njohura për këtë është tandem Apache Kafka dhe Spark Streaming, ku Kafka krijon një rrjedhë paketash të mesazheve hyrëse, ndërsa Spark Streaming përpunon këto paketa nëpërmjet një intervali të caktuar kohor.
PĂ«r tĂ« rritur disponueshmĂ«rinĂ« e aplikacionit, do tĂ« pĂ«rdorim piketat e kontrollit â checkpointet. Me kĂ«tĂ« mekanizĂ«m, kur moduli Spark Streaming ka nevojĂ« tĂ« rikuperojĂ« tĂ« dhĂ«nat e humbura, do t'i duhet vetĂ«m tĂ« kthehet te pika e fundit e kontrollit dhe tĂ« vazhdojĂ« llogaritjet nga aty.
Arkitektura e sistemit në zhvillim

Komponentët e përdorur:
- â Ă«shtĂ« njĂ« sistem i shpĂ«rndarĂ« pĂ«r shkĂ«mbimin e mesazheve me publikim dhe abonim. PĂ«rshtatet pĂ«r konsumimin e mesazheve nĂ« mĂ«nyrĂ« tĂ« autonomĂ«, si dhe nĂ« mĂ«nyrĂ« online. PĂ«r tĂ« parandalimeve tĂ« humbjes sĂ« tĂ« dhĂ«nave, mesazhet e Kafka ruhen nĂ« disk dhe janĂ« tĂ« replikura brenda klasterit. Sistemi Kafka Ă«shtĂ« ndĂ«rtuar mbi shĂ«rbimin e sinkronizimit ZooKeeper;
- â komponenti Spark pĂ«r pĂ«rpunimin e tĂ« dhĂ«nave nĂ« rrjedhĂ«. Moduli Spark Streaming Ă«shtĂ« ndĂ«rtuar duke pĂ«rdorur arkitekturĂ«n e "mikropaketimeve" (micro-batch architecture), ku rrjedha e tĂ« dhĂ«nave interpretohet si njĂ« seri tĂ« vazhdueshme paketash 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 nĂ«pĂ«rmjet intervaleve tĂ« rregullta tĂ« kohĂ«s. NĂ« fillim tĂ« çdo intervali kohor krijohet njĂ« paketĂ« e re, dhe çdo tĂ« dhĂ«nĂ« qĂ« ka ardhur gjatĂ« kĂ«tij intervali pĂ«rfshihet nĂ« paketĂ«. NĂ« fund tĂ« intervalit, rritja e paketĂ«s ndalon. MadhĂ«sia e intervalit pĂ«rcaktohet nga njĂ« parameter qĂ« quhet intervali i paketimit (batch interval);
- â bashkon pĂ«rpunimin relacional me programimin funksional tĂ« Spark. TĂ« dhĂ«nat e strukturuara i referohen tĂ« dhĂ«nave 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, pĂ«r shkak tĂ« informacionit pĂ«r skemĂ«n, ai mund tĂ« nxjerrĂ« nĂ« mĂ«nyrĂ« efikase vetĂ«m fushat e nevojshme tĂ« regjistrimeve, siç ofron gjithashtu interface-t e API-sĂ« pĂ«r DataFrame;
- â Ă«shtĂ« njĂ« bazĂ« tĂ« dhĂ«nash relacional nĂ« cloud relativisht e pĂ«rballueshme, njĂ« shĂ«rbim nĂ« internet qĂ« thjeshton konfigurimin, funksionimin dhe shkallĂ«zimin, administruar direkt nga Amazon.
Instalimi dhe nisja e serverit Kafka
Para përdorimit të drejtpërdrejtë të Kafka, duhet të siguroheni që Java është e instaluar, pasi për funksionimin përdoret JVM:
sudo apt-get update
sudo apt-get install default-jre
java -version
Të krijojmë një përdorues të ri për punë me Kafka:
sudo useradd kafka -m
sudo passwd kafka
sudo adduser kafka sudo
Më pas shkarkoni shpërndarjen 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"Zhvilloni 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 tjetër është opsional. Arsyet është se cilësimet e paracaktuara nuk lejojnë shfrytëzimin e plotë të të gjitha mundësive të Apache Kafka. Për shembull, për të fshirë një temë, kategori, grup, në të cilat mund të publikohen mesazhe. Për ta ndërruar këtë, do të redaktojmë skedarin e konfigurimit:
vim ~/kafka/config/server.propertiesShtoni në fund të skedarit sa vijon:
delete.topic.enable = truePara të nisni serverin Kafka, së pari duhet të filloni serverin ZooKeeper, do të përdorim një skenar ndihmës që vjen së bashku me shpërndarjen e Kafka:
Cd ~/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties
Pasi ZooKeeper të nisë me sukses, në një terminal tjetër nisim serverin Kafka:
bin/kafka-server-start.sh config/server.propertiesTë krijojmë një temë të re me emrin Transaction:
bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic transactionTë sigurohemi që tema me numrin e duhur të ndarjeve dhe replikave është krijuar:
bin/kafka-topics.sh --describe --zookeeper localhost:2181 
Do ta injorojmĂ« testimin e prodhuesit dhe konsumatorit pĂ«r temĂ«n e krijuar rishtazi. MĂ« shumĂ« detaje rreth mĂ«nyrĂ«s se si mund tĂ« testoni dĂ«rgimin dhe marrjen e mesazheve janĂ« shkruar nĂ« dokumentacionin zyrtar â . Tani kalojmĂ« nĂ« shkruarjen e prodhuesit nĂ« Python duke pĂ«rdorur KafkaProducer API.
Shkrimi i prodhuesit
Prodhuesi do të gjenerojë të dhëna të rastësishme - 100 mesazhe çdo sekondë. Me të dhëna të rastësishme do të kuptojmë një fjalor me tre fusha:
- DegĂ« â emri i pikĂ«s sĂ« shitjes sĂ« institucionit financiar;
- Monedha â monedha e transaksionit;
- Shuma â shuma e transaksionit. Shuma do tĂ« jetĂ« njĂ« numĂ«r pozitiv nĂ«se Ă«shtĂ« blerje e valutĂ«s nga Banka, dhe negativ nĂ«se Ă«shtĂ« shitje.
Kodi për prodhuesin 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ë temën që na nevojitet, 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ë temë:
{}, 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 e fillojmë skriptin, marrim mesazhe të tilla në terminal:

Kjo do tĂ« thotĂ« se gjithçka funksionon siç e donim â prodhuesi gjeneron dhe dĂ«rgon mesazhe nĂ« temĂ«n qĂ« na nevojitet.
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 pĂ«rpunimin e klasterĂ«ve.
Sa i përket performancës, Spark tejkalon implementimet më të njohura të modelit MapReduce, duke ofruar gjithashtu mbështetje për një gamë më të gjerë llojesh përpunimi, duke përfshirë kërkesat interaktive dhe përpunimin e flukseve. Shpejtësia luan një rol kritik në përpunimin e sasisë të mëdha të të dhënave, pasi shpejtësia lejon të punojmë në mënyrë interaktive, pa pritur minuta ose orë. Një nga avantazhet më të rëndësishme të Spark që ofron një shpejtësi kaq të lartë është aftësia për të kryer llogaritje në kujtesë.
Ky framework është shkruar në Scala, kështu që duhet të instaloni atë si hapin e parë:
sudo apt-get install scalaShkarkoni distribucionin e Spark nga faqja zyrtare:
wget "http://mirror.linux-ia64.org/apache/spark/spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz"Shkëputni arkivin:
sudo tar xvf spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz -C /usr/local/sparkShtoni rrugën për në Spark në skedarin bash:
vim ~/.bashrcShtoni përmes redaktuesit këto linja:
SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH
Ecjirim komandën më poshtë pasi të bëjmë ndryshimet në bashrc:
source ~/.bashrcZhvillimi i AWS PostgreSQL
Tani na mbetet të zhvillojmë 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.
HymĂ« nĂ« konsolĂ«n AWS â> AWS RDS â> Databases â> Krijo bazĂ«n e tĂ« dhĂ«nave:

Zgjidhim PostgreSQL dhe klikojmë butonin Next:

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

MĂ« pas, vendosim njĂ« shenjĂ« nĂ« bllokun Free Tier, dhe pas kĂ«saj do tĂ« na ofrohet automatikisht njĂ« instancĂ« e klasĂ«s t2.micro â ndonĂ«se e dobĂ«t, Ă«shtĂ« falas dhe e pĂ«rshtatshme pĂ«r qĂ«llimin tonĂ«:

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

Në faqe tjetër ndodhen parametrat që përcaktojnë 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, i cili do të lejojë qasjen nga jashtë në serverin tonë DB përmes portit 5432 (PostgreSQL).
TĂ« kalojmĂ« nĂ« njĂ« dritare tĂ« veçantĂ« shfletuesi nĂ« konsolĂ«n AWS nĂ« seksionin VPC Dashboard â> Security Groups â> Krijo grupin e sigurisĂ«:

Vendosim emrin pĂ«r grupin e sigurisĂ« â PostgreSQL, pĂ«rshkrimin, specifikojmĂ« se cila VPC duhet tĂ« asociohet me kĂ«tĂ« grup dhe klikojmĂ« butonin Krijo:

Plotësojmë për grupin e sapokrijuar rregullat hyrëse Inbound rules për portin 5432, siç është treguar në imazhin më poshtë. Portin mund ta zgjedhim drejtpërdrejt nga lista e tipave duke mos e specifikuar manualisht.
Në fakt, vlera ::/0 tregon aksesin e trafikut hyrës për serverin nga e gjithë bota, që në mënyrë kanonike nuk është krejtësisht e saktë, por për qëllime të këtij shembulli, do ta lejojmë këtë qasje:

Kthehemi nĂ« faqen e shfletuesit, ku kemi hapur "Configure advanced settings" dhe zgjedhim nĂ« seksionin VPC security groups â> Zgjidh grupet ekzistuese tĂ« sigurisĂ« VPC â> PostgreSQL:

MĂ« pas, nĂ« seksionin Database options â> Emri i databazĂ«s â> vendosim emrin â habrDB.
Parametrat e tjerĂ«, pĂ«rveç ndoshta çaktivizimit tĂ« backup-it (backup retention period â 0 days), monitorimit dhe Performance Insights, mund t'i lĂ«mĂ« sipas parazgjedhjeve. Klikoni butonin Create database:

Përpunuesi i rrjedhave
Hapi i fundit do të jetë zhvillimi i një Spark-jobi, i cili do të përpunojë të dhëna të reja çdo dy sekonda, të ardhura nga Kafka dhe do të ruajë rezultatin në bazën e të dhënave.
Siç u përmend më sipër, pikëkontrolli (checkpoints) është mekanizmi kryesor në SparkStreaming që duhet të konfigurohet për të siguruar qëndrushmërinë. Do të përdorim pikëkontrollin dhe, në rast se procedura dështon, moduli Spark Streaming për rikuperimin e të dhënave të humbura do të duhet thjesht të kthehet në pikëkontrollin e fundit dhe të vazhdojë llogaritjet nga aty.
Pika e kontrollit mund të aktivizohet duke vendosur katalogun në një sistem të besueshëm, të qëndrueshëm për skedarët (p.sh., HDFS, S3 etj.), ku do të ruhen informatat e pikëkontrollit. Kjo bëhet, për shembull:
streamingContext.checkpoint(checkpointDirectory)Në shembullin tonë, do të përdorim qasjen e mëposhtme, domethënë, nëse checkpointDirectory ekziston, konteksti do të rikrijohet nga të dhënat e pikëkontrollit. Nëse katalogu nuk ekziston (dmth., ekzekutohet për herë të parë), do të thërritet funksioni functionToCreateContext për të krijuar një kontekst të ri dhe për të konfiguruar DStreams:
nga pyspark.streaming import StreamingContext
context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)
Krijojmë një objekt DirectStream për të lidhur me temën «transaction» përmes metodës createDirectStream nga biblioteka 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})
Parse 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 tekstin e pyetjes dhe ekzekutimi i saj përmes Spark SQL:
sql_query = get_sql_query()
testResultDataFrame = spark.sql(sql_query)
testResultDataFrame.show(n=5)
Dhe më pas ruajmë të dhënat e grumbulluara në një tavëll në AWS RDS. Për të ruajtur rezultatet e grumbullimit 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 disa fjalë mbi konfigurimin e lidhjes me AWS RDS. Përdoruesi dhe fjalëkalimi për këtë u krijuan në hapin "Zhvillimi i AWS PostgreSQL". Si url për serverin e bazës së të dhënave, duhet të përdoret Endpoint-i që shfaqet në seksionin Connectivity & security:
PĂ«r tĂ« siguruar njĂ« lidhje tĂ« saktĂ« midis Spark dhe Kafka, duhet tĂ« ekzekutojmĂ« punĂ«n pĂ«rmes spark-submit duke pĂ«rdorur artifactin spark-streaming-kafka-0-8_2.11. SĂ« bashku, do tĂ« aplikojmĂ« gjithashtu artifactin pĂ«r ndĂ«rveprimin me bazĂ«n e tĂ« dhĂ«nave PostgreSQL, qĂ« do tâi kalojmĂ« pĂ«rmes âpackages.
Për fleksibilitetin e skriptit, do të nxjerrim gjithashtu emrin e serverit të mesazheve dhe temën nga e cila duam të marrim të dhënat si parametra hyrës.
Pra, erdhi koha të nisni dhe të kontrolloni 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Ă« doli me sukses! Siç shihet nĂ« imazhin mĂ« poshtĂ« â gjatĂ« punĂ«s sĂ« aplikacionit, rezultatet e reja tĂ« agregimit shfaqen çdo 2 sekonda, sepse kemi vendosur intervalin e paketimit tĂ« barabartĂ« me 2 sekonda, kur krijuam objekti StreamingContext:

Tani, bëjmë një kërkesë të thjeshtë në bazën e të dhënave për të kontrolluar nëse ka regjistrime në tabelë. transaction_flow:

Përfundimi
Në këtë artikull u shqyrtua një shembull i përpunimit të informacionit në kohë reale 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 të rrjedhës dhe aplikacioneve që veprojnë në shkallë reale është e papërfshirë.
Kodi burimor i plotë mund ta gjeni në depunimin tim në .
Me kënaqësi jam i gatshëm të diskutoj këtë artikull, po pres komentet tuaja dhe shpresoj për kritika konstruktive nga të gjithë lexuesit e interesuar.
Ju uroj sukses!
Ps. Fillimisht ishte planifikuar të përdorej një DB lokale PostgreSQL, por duke marrë parasysh dashurinë time për AWS, vendosa ta zhvendos bazën e të dhënave në cloud. Në artikullin e ardhshëm në këtë temë, do të tregoj se si të realizojmë gjithë sistemin e përshkruar më sipër në AWS me ndihmën e AWS Kinesis dhe AWS EMR. Qëndroni të informuar!
Burimi: habr.com

