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ë 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ë!

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

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

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

Shtoni në fund të skedarit sa vijon:

delete.topic.enable = true

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

Të krijojmë një temë të re me emrin Transaction:

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

Të sigurohemi që tema me numrin e duhur të ndarjeve dhe replikave është 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 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 — DĂ«rgoni disa mesazhe. 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:

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

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 scala

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

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

vim ~/.bashrc

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

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

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

Duke qenë se ky shembull shqyrtohet ekskluzivisht për qëllime arsimore, 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Ă« 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Ă«:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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

Në faqe tjetër ndodhen parametrat që përcaktojnë 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, 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Ă«:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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

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:
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 —> Zgjidh grupet ekzistuese tĂ« sigurisĂ« VPC —> PostgreSQL:
Apache Kafka dhe përpunimi i të dhënave nëpërmjet Spark Streaming

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

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:

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

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:

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

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:

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

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

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

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