Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Tere, Habr! TĂ€na ehitame sĂŒsteemi, mis kasutab Spark Streamingut Apache Kafka sĂ”numivoogude töötlemiseks ja kirjutab töötlemise tulemuse AWS RDS pilveandmebaasi.

Kujutage ette, et vahendite ettevĂ”te seab meile ĂŒlesande sisenevate tehingute töötlemiseks „reaalajas“ kĂ”igis oma filiaalides. Seda vĂ”ib teha, et kiiresti arvutada avatud valuuta positsiooni, limiteid vĂ”i finantstulemust tehingute osas jne.

Kuidas seda juhtumit ellu viia ilma maagia ja nĂ”iduse rakendamiseta — loe edasi! Alustame!

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga
(Pildi allikas)

Sissejuhatus

Ilmselt pakub suure hulga andmete töötlemine reaalajas kaasaegsetes sĂŒsteemides laialdasi vĂ”imalusi. Üheks populaarseks kombinatsiooniks on Apache Kafka ja Spark Streaming tandem, kus Kafka genereerib sisenevate sĂ”numite voogu ja Spark Streaming töötleb neid pakette mÀÀratud ajavahemiku jooksul.

Rakenduse seisukindluse suurendamiseks kasutame kontrollpunkte — checkpoints. Selle mehhanismi abil, kui Spark Streaming moodul vajab kadunud andmete taastamist, peab see lihtsalt naasma viimasest kontrollpunktist ja jĂ€tkama arvutusi sealt.

Arendatava sĂŒsteemi arhitektuur

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Kasutatavad komponendid:

  • Apache Kafka — see on hajustatud sĂ”numite jagamise sĂŒsteem, millel on avaldamine ja tellimine. Sobib nii iseseisvaks kui ka veebipĂ”hiseks sĂ”numite tarbimiseks. Andmete kaotuse vĂ€ltimiseks salvestatakse Kafka sĂ”numid kettale ja replikeeritakse klastris. Kafka sĂŒsteem on ĂŒles ehitatud ZooKeeper sĂŒnkroniseerimisteenuse peale;
  • Apache Spark Streaming — Spark komponent andmete voogude töötlemiseks. Spark Streaming moodul on ĂŒles ehitatud mikropartiidest koosneva arhitektuuri alusel, kus andmevoog tĂ”lgendatakse pideva vĂ€ikeste andmepakettide jadana. Spark Streaming vĂ”tab andmeid erinevatest allikatest ja ĂŒhendab need vĂ€ikesteks pakettideks. Uued paketid luuakse regulaarsete ajavahemike jĂ€rel. Iga ajavahemiku alguses luuakse uus pakett, ja kĂ”ik andmed, mis selle ajavahemiku jooksul saabuvad, lisatakse paketti. Ajavahemiku lĂ”pus paketise suurendamine lĂ”ppeb. Ajavahemiku suurus mÀÀratakse parameetriga nimega paketi intervall;
  • Apache Spark SQL — ĂŒhendab relatsioonilise töötlemise Spark'i funktsionaalse programmeerimisega. Struktureeritud andmed viitavad andmetele, millel on skeem, st kĂ”igi kirje jaoks ĂŒhtne vĂ€ljade kogum. Spark SQL toetab sisendit mitmest struktureeritud andmete allikast ja skeeminformatsiooni olemasolu tĂ”ttu suudab tĂ”husalt ekstraktsioonida ainult vajalikud vĂ€ljad, samuti pakub DataFrame API.
  • AWS RDS — see on suhteliselt odav pilvepĂ”hine relatsiooniline andmebaas, veebiteenus, mis lihtsustab seadistamist, haldamist ja skaleerimist ning mida haldab otse Amazon.

Kafka serveri paigaldamine ja kÀivitamine

Enne Kafka kasutamist peate veenduma, et Java on olemas, kuna JVM-i kasutatakse tööks:

sudo apt-get update 
sudo apt-get install default-jre
java -version

Loome uue kasutaja Kafka jaoks:

sudo useradd kafka -m
sudo passwd kafka
sudo adduser kafka sudo

SeejÀrel laadime alla distributsiooni Apache Kafka ametlikult veebisaidilt:

wget -P /YOUR_PATH "http://apache-mirror.rbc.ru/pub/apache/kafka/2.2.0/kafka_2.12-2.2.0.tgz"

Pakkime alla laetud arhiivi lahti:

tar -xvzf /YOUR_PATH/kafka_2.12-2.2.0.tgz
ln -s /YOUR_PATH/kafka_2.12-2.2.0 kafka

JÀrgmine samm on valikuline. Asjaolu on see, et vaikeseaded ei vÔimalda Apache Kafka kÔiki funktsioone tÀielikult kasutada. NÀiteks teema, kategooria vÔi grupi kustutamine, kuhu sÔnumeid avaldada. Selle muutmiseks redigeerime konfiguratsioonifaili:

vim ~/kafka/config/server.properties

Lisage faili lÔppu jÀrgmine:

delete.topic.enable = true

Enne Kafka serveri kÀivitamist tuleb kÀivitada ZooKeeperi server. Kasutame abiskripti, mis tuleb koos Kafka distributsiooniga:

Cd ~/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties

Kui ZooKeeper on edukalt kÀivitunud, kÀivitame Kafka serveri eraldi terminalis:

bin/kafka-server-start.sh config/server.properties

Loome uue teema nimega Transaction:

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

Veendume, et teema on loodud vajaliku partitsioonide arvu ja replikatsiooniga:

bin/kafka-topics.sh --describe --zookeeper localhost:2181

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

JĂ€tame vahele proovide testimise tootja ja tarbija jaoks uue teema. TĂ€psemalt, kuidas saata ja vastu vĂ”tta sĂ”numeid, on kirjas ametlikus dokumentatsioonis — Saada mĂ”ned sĂ”numid. Siiski, liigume edasi Pythonis tootja kirjutamise juurde, kasutades KafkaProducer API-t.

Tootja kirjutamine

Tootja genereerib juhuslikke andmeid — 100 sĂ”numit iga sekundi tagant. Juhuslike andmete all mĂ”istame kolme vĂ€lja koosnevat sĂ”nastikku:

  • Filiaal — krediidiasutuse mĂŒĂŒgikoha nimi;
  • Currency — tehingu valuuta;
  • Amount — tehingu summa. Summa on positiivne number, kui see on valuuta ostmine pangast, ja negatiivne, kui mĂŒĂŒk.

Produtsendi kood nÀeb vÀlja jÀrgmiselt:

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

SeejÀrel, kasutades meetodit send, saadame teate serverile, sobivasse teema, JSON formaadis:

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('--> Teade on saadetud teemasse: 
            {}, jagu: {}, offset: {}' 
            .format(record_metadata.topic,
                record_metadata.partition,
                record_metadata.offset ))   
                             
except Exception as e:
    print('--> NĂ€ib, et tekkis viga: {}'.format(e))

finally:
    producer.flush()

Skripti kÀivitamisel saame terminalis jÀrgmised sÔnumid:

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

See tĂ€hendab, et kĂ”ik töötab nii nagu soovisime — produtsent genereerib ja saadab sĂ”numeid sobivasse teema.
JÀrgmine samm on Spark'i installimine ja selle sÔnumivoogude töötlemine.

Apache Spark'i installimine

Apache Spark on universaalne ja suure jÔudlusega klastritehnoloogia arvutusalus.

Spark ĂŒletab oma jĂ”udluses tuntud MapReduce'i rakendused, pakkudes samal ajal toe laiema hulga arvutusliikide, sealhulgas interaktiivsete pĂ€ringute ja voogedastuse jaoks. Kiirus mĂ€ngib suurt rolli suurte andmemahtude töötlemisel, kuna just kiirus vĂ”imaldab töötada interaktiivses reĆŸiimis, raiskamata minuteid vĂ”i tunde ootamisele. Üks Spark'i peamisi eeliseid, mis tagab sellise kĂ”rge kiirus, on vĂ”ime teostada arvutusi mĂ€lus.

See raamistik on kirjutatud Scala keeles, seega tuleb see esmalt installida:

sudo apt-get install scala

Laadime ametlikult alla Spark'i distributsiooni:

wget "http://mirror.linux-ia64.org/apache/spark/spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz"

Pakkime arhivi lahti:

sudo tar xvf spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz -C /usr/local/spark

Lisame Spark'i tee bash-faili:

vim ~/.bashrc

Lisame redigeerijaga jÀrgmised read:

SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH

KÀivitage alljÀrgnev kÀsk pÀrast muutmisi bashrc-s:

source ~/.bashrc

AWS PostgreSQLi ĂŒles seadmine

JĂ€rgmiseks peame ĂŒles seadma andmebaasi, kuhu laadime voogudest töödeldud teabe. Selleks kasutame AWS RDS teenust.

Logige sisse AWS konsooli —> AWS RDS —> Andmebaasid —> Looge andmebaas:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Valige PostgreSQL ja klÔpsake nuppu JÀrgmine:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Arvestades, et see nÀide on mÔeldud ainult hariduslikel eesmÀrkidel, kasutame tasuta serverit 'minimaalsel tasemel' (Free Tier):
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

SeejĂ€rel mĂ€rkige Free Tier kast ja pĂ€rast seda pakutakse automaatselt instantsi klassist t2.micro — kuigi see on nĂ”rk, on see tasuta ja sobib meie vajadustele:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

SeejÀrel on vÀga olulised viisid: andmebaasi instantsi nimi, master-kasutaja nimi ja selle parool. Nimetame instantsi: myHabrTest, master-kasutaja: habr, parool: habr12345 ja klÔpsake nuppu JÀrgmine:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

JÀrgmises lehes on parameetrid, mis vastutavad meie andmebaasi serveri vÀlise kÀttesaadavuse (Public accessibility) ja portsede kÀttesaadavuse eest:

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Loo uus seadistus VPC turvagruppi, mis vĂ”imaldab vĂ€ljastpoolt meie andmebaasi serveriga ĂŒhendust luua sadamast 5432 (PostgreSQL).
Avame eraldi brauseriaknas AWS konsoolis VPC juhtpaneelis ja valime turvagruppide alt Create security group:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Kasutame turvagrupi jaoks nime — PostgreSQL, lisame kirjelduse, valime, mille VPC-ga see grupp peab olema seotud, ja vajutame nuppu Create:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

TĂ€idame uue grupi Inbound rules reeglid portidele 5432, nagu on nĂ€idatud alloleval pildil. Porti ei pea kĂ€sitsi mÀÀrama, vaid valige PostgreSQL rippmenĂŒĂŒst Type.

TĂ€psemalt öeldes tĂ€hendab vÀÀrtus ::/0, et sisenev liiklus on serverile saadaval ĂŒle kogu maailma, mis pole klassikaliselt tĂ€iesti Ă”ige, kuid selle nĂ€ite analĂŒĂŒsimiseks lubame endale sellise lĂ€henemise:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Naaseme brauseri lehele, kus on avatud „Configure advanced settings“ ja valime VPC turvagruppide alt — Choose existing VPC security groups — PostgreSQL:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

SeejĂ€rel Database options — Database name — lisame nime — habrDB.

ÜlejÀÀnud parameetreid, vĂ€lja arvatud varundamise vĂ€ljalĂŒlitamine (backup retention period — 0 days), jĂ€lgimise ja Performance Insights osas, vĂ”ime jĂ€tta vaikimisi. Vajutame nuppu Create database:
Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Voogude töötleja

Viimane etapp on Spark-töötluse arendamine, mis töötleb iga kahe sekundi jÀrel uusi andmeid, mis saabuvad Kafka'st, ja talletab tulemuse andmebaasi.

Nagu eespool mainitud, on kontrollpunktid (checkpoints) pĂ”hiline mehhanism SparkStreaming'is, mis tuleb seadistada, et tagada sĂŒsteemi vastupidavus. Kasutame kontrollpunkte ja juhul, kui protseduur tĂ”rkub, vĂ”ib Spark Streaming mooduli abil kadunud andmed taastada, lihtsalt naastes viimasele kontrollpunktile ja jĂ€tkates arvutusi sealt.

Kontrollpunkti saab aktiveerida, mÀÀrates katalooge vastupidavasse, usaldusvÀÀrsesse failisĂŒsteemi (nt HDFS, S3 jne), kuhu salvestatakse kontrollpunkti teave. Seda saab teha nĂ€iteks:

streamingContext.checkpoint(checkpointDirectory)

Meie nÀites rakendame jÀrgmist lÀhenemist, nimelt, kui checkpointDirectory eksisteerib, rekonstrueeritakse kontekst kontrollpunkti andmete pÔhjal. Kui kataloog ei eksisteeri (st see toimub esmakordselt), siis kutsutakse vÀlja funktsioon functionToCreateContext uue konteksti loomiseks ja DStreamide seadistamiseks:

import pyspark.streaming.StreamingContext

context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)

Loome objekti DirectStream, et ĂŒhendada teema 'transaction' kasutades KafkaUtils raamatukogu meetodit createDirectStream:

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})

Parseerime sissetulevaid andmeid JSON formaadis:

rowRdd = rdd.map(lambda w: Row(branch=w['branch'],
                                       currency=w['currency'],
                                       amount=w['amount']))
                                       
testDataFrame = spark.createDataFrame(rowRdd)
testDataFrame.createOrReplaceTempView("treasury_stream")

Kasutades Spark SQL, teeme lihtsa grupeerimise ja prindime tulemuse konsooli:

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

KĂŒsimuse teksti saamine ja selle KĂ€itamine Spark SQL-is:

sql_query = get_sql_query()
testResultDataFrame = spark.sql(sql_query)
testResultDataFrame.show(n=5)

Ja seejÀrel salvestame saadud koondatud andmed tabelisse AWS RDS. Andmete salvestamiseks kasutame DataFrame objekti meetodit write:

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()

AWS RDS-iga ĂŒhenduse seadistamisest. Kasutajat ja parooli loodi sammul "AWS PostgreSQL-i juurutamine". Andmebaasi serveri URL-ina tuleks kasutada Endpoint'i, mis kuvatakse Connectivity & security jaotises:

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

Selleks, et Spark ja Kafka korrektselt suhelda, tuleks töö ĂŒles laadida smark-submit kaudu, kasutades artefakti spark-streaming-kafka-0-8_2.11. TĂ€iendavalt rakendame ka artefakti PostgreSQL andmebaasiga suhtlemiseks, edastame need lĂ€bi —packages.

Skripti paindlikkuse huvides, toome sisse ka sÔnumite serveri ja teema nimed, millest soovime andmeid saada.

NĂŒĂŒd on aeg sĂŒsteemi kĂ€ivitada ja selle töökorrasolekut kontrollida:

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

KĂ”ik lĂ€ks korda! Nagu allolevalt pildilt nĂ€ha — rakenduse töö kĂ€igus kuvatakse uusi agregatsioonitulemusi iga 2 sekundi tagant, kuna seadsime paketi intervalli 2 sekundile, kui lĂ”ime StreamingContext objekti:

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

JÀrgmisena teeme lihtsa pÀringu andmebaasi, et kontrollida, kas tabelis on kirjeid transaction_flow:

Apache Kafka ja andmevoogude töötlemine Spark Streaminguga

KokkuvÔte

Antud artiklis kĂ€sitleti nĂ€idet voogedastuse töötlemisest Spark Streaming kasutamisega koos Apache Kafka ja PostgreSQL-ga. Andmemahtude kasvu tĂ”ttu erinevatest allikatest on keeruline ĂŒle hinnata Spark Streaming praktilist vÀÀrtust voogedastusrakenduste ja reaalajas toimivate rakenduste loomisel.

TÀieliku lÀhtekoodi leiate minu repolt GitHub.

Olen valmis arutama seda artiklit, ootan teie kommentaare ning loodan kÔigi murelike lugejate konstruktiivset kriitikat.

Soovin edu!

Ps. Alguses plaaniti kasutada kohalikku PostgreSQL andmebaasi, kuid arvestades minu armastust AWS-i vastu, otsustasin viia andmebaasi pilve. JĂ€rgmises artiklis sellel teemal nĂ€itan, kuidas ellu viia tĂ€ielikult ĂŒlal kirjeldatud sĂŒsteem AWS-is, kasutades AWS Kinesis ja AWS EMR. Hoidke end kursis!

Allikas: habr.com

Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster