Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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Ă«!

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming
(Burimi i figurës)

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

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

Komponentët e përdorur:

  • Apache Kafka — Ă«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;
  • Apache Spark Streaming — 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);
  • Apache Spark SQL — 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;
  • AWS RDS — Ă«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.properties

Shtoni në fund të skedarit sa vijon:

delete.topic.enable = true

Para 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.properties

Do të krijojmë një topic të ri me emrin Transaction:

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

Do 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

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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 — DĂ«rgoni disa mesazhe. 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:

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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 scala

Shkarko 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/spark

Shtojmë rrugën për në Spark në skedarin bash:

vim ~/ .bashrc

Shtojmë 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 ~/ .bashrc

Zhvillimi 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:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

Zgjidh PostgreSQL dhe klikoni butonin Next:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

Duke marrë parasysh se ky shembull shqyrtohet ekskluzivisht për qëllime edukative, do të përdorim një server falas "në minimalet" (Free Tier):
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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Ă«:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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ë:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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ë:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

Kthehemi në faqen e shfletuesit, ku kemi hapur «Configure advanced settings» dhe zgjedhim në seksionin VPC security groups -> Choose existing VPC security groups -> PostgreSQL:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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:

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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:

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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:

Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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ë GitHub.

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

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster