Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Tere, Habr! TĂ€na ehitame sĂŒsteemi, mis töötleb Apache Kafka sĂ”numivoolusid Spark Streaming abil ja salvestab töötlemise tulemused AWS RDS pilveandmebaasi.

Kujutame ette, et mĂ”ni krediidiasutus esitab meile ĂŒlesande töötleda sissetulevaid tehinguid «reaalajas» kĂ”igis oma filiaalides. Seda vĂ”ib teha, et kiiresti arvutada avatud valuutapositsioon soovituslikule rahandusele, limiidid vĂ”i finantsilised tulemused tehingute osas jne.

Kuidas seda juhtumit rakendada ilma maagia ja imeliste loitsudeta — loe edasi! Alustame!

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga
(Pildi allikas)

Sissejuhatus

Ilmselgelt pakub suure andmemassi töötlemine reaalajas laialdasi vĂ”imalusi tĂ€napĂ€evastes sĂŒsteemides. Üks populaarsamaid kombinatsioone selleks on tandem Apache Kafka ja Spark Streaming, kus Kafka loob sissetulevate sĂ”numite voogusid ja Spark Streaming töötleb neid pakette mÀÀratud ajavahemiku jooksul.

Rakenduse kindluse suurendamiseks kasutame kontrollpunkte — checkpoints. Selle mehhanismi abil, kui Spark Streaming moodul peab taastama kadunud andmed, peab ta lihtsalt tagasipöörduma viimasest kontrollpunktist ja jĂ€tkama arvutusi sealt.

Arendatava sĂŒsteemi arhitektuur

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Kasutatavad komponendid:

  • Apache Kafka — see on jagatud sĂ”numivahetussĂŒsteem, mis toetab publikatsiooni ja tellimist. Sobib nii iseseisvaks kui ka veebipĂ”hiseks sĂ”numite tarbimiseks. Andmete kadumise vĂ€ltimiseks salvestatakse Kafka sĂ”numid kettale ja replitseeritakse klastris. Kafka sĂŒsteem on ĂŒles ehitatud ZooKeeperi sĂŒnkroonimisteenuse peale;
  • Apache Spark Streaming — Spark komponent voogandmete töötlemiseks. Spark Streaming moodul on ĂŒles ehitatud mikro-pakettide arhitektuurile, kus voogandmed tĂ”lgendatakse kui pidevat vĂ€ikeste andmepakettide jĂ€rjestust. Spark Streaming vĂ”tab vastu andmeid erinevatest allikatest ja ĂŒhendab need vĂ€ikeste pakkideks. Uued pakkide loomiseks kasutatakse regulaarselt ajaintervalli. Iga ajaintervalli alguses luuakse uus pakk, kuhu sisaldatakse kĂ”ik andmed, mis on saabunud selle intervalli jooksul. Intervalli lĂ”pus lĂ”petatakse pakkide suurendamine. Intervalli suurus on mÀÀratud parameetriga, mida nimetatakse partii intervalliks.
  • Apache Spark SQL — ĂŒhendab relatsioonilise töötlemise Spark funktsionaalse programmeerimisega. Struktureerituna mĂ”istetakse andmeid, millel on skeem, see tĂ€hendab ĂŒhtne vĂ€ljade komplekt kĂ”igikirjete jaoks. Spark SQL toetab mitmest struktureeritud andmete allikast sisendi, ja kuna skeemi teavet on olemas, saab see tĂ”husalt eraldada ainult vajalikud rekordite vĂ€ljad ning pakub DataFrame API-sid.
  • AWS RDS — suhteliselt odav pilveline relatsiooniline andmebaas, veebiteenus, mis lihtsustab seadistamist, hooldust ja skaleerimist, mida haldab otse Amazon.

Kafka serveri paigaldamine ja kÀivitamine

Enne Kafka kasutamist tuleb veenduda, et Java on olemas, kuna kasutatakse JVM-i:

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

Loome uue kasutaja Kafka tööks:

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

SeejÀrel laadige alla distributsioon ametlikult Apache Kafka saidilt:

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

Pakkige alla laaditud arhiiv 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 tÀielikult kasutada Apache Kafka kÔiki vÔimalusi. NÀiteks ei saa kustutada teemasid, kategooriaid, gruppe, millele sÔnumeid vÔib avaldada. Selle muutmiseks toimetame konfigureerimisfaili:

vim ~/kafka/config/server.properties

Lisage faili lÔppu jÀrgmine:

delete.topic.enable = true

Enne Kafka serveri kÀivitamist tuleb kÀivitada ZooKeeper server, kasutame abiskripti, mis on kaasas Kafka distributsiooniga:

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

PÀrast seda, 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

Veenduge, et teema oleks loodud vajaliku arvu partitsioonide ja replikatsiooniga:

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

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

JĂ€tame vahele testimise hetked produtsendi ja tarbija jaoks uue loodud teema puhul. Üksikasjalikumalt selle kohta, kuidas saatmist ja sĂ”numite vastuvĂ”ttu testida, on kirjas ametlikus dokumentatsioonis — Saada mĂ”ned sĂ”numid. Aga liigume edasi Pythonis produtsendi kirjutamise juurde, kasutades KafkaProducer API.

Produtsendi kirjutamine

Produtsent genereerib juhuslikke andmeid — 100 sĂ”numit iga sekundi jooksul. Juhuslike andmete all mĂ”istame sĂ”nastikku, mis koosneb kolmest valdkonnast:

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

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

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 sÔnumi serverisse, soovitud teemas, 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('--> SÔnum on saadetud teemale: 
            {}, partitsioon: {}, positsioon: {}' 
            .format(record_metadata.topic,
                record_metadata.partition,
                record_metadata.offset ))   
                             
except Exception as e:
    print('--> Tundub, et tekkis viga: {}'.format(e))

finally:
    producer.flush()

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

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

See tĂ€hendab, et kĂ”ik töötab nii nagu soovisime — produtsent genereerib ja saadab sĂ”numeid meie soovitud teemale.
JÀrgmine samm on Spark'i installeerimine ja selle sÔnumivooge töötlemine.

Apache Spark'i paigaldamine

Apache Spark on universaalne ja kÔrge jÔudlusega klastritehnoloogia platvorm.

JĂ”udluse poolest ĂŒletab Spark populaarseid MapReduce mudeli rakendusi, pakkudes samal ajal laiemat toetust erinevat tĂŒĂŒpi arvutustele, sealhulgas interaktiivsetele pĂ€ringutele ja voogedastusele. Kiirus mĂ€ngib suurt rolli suurte andmemahtude töötlemisel, kuna just kiirus vĂ”imaldab töötada interaktiivselt, vĂ€ltides minutite vĂ”i tundide ootamist. Üks Spark'i suurimaid eeliseid, mis tagab sellise kĂ”rge kiirus, on vĂ”ime teostada arvutusi mĂ€lus.

KÀesolev raamistik on kirjutatud Scala keeles, seetÔttu tuleb see esmalt paigaldada:

sudo apt-get install scala

Laadime ametlikult Spark'i distributsiooni alla:

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

Pakkime arhiivi lahti:

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

Lisame tee Spark'i bash-faili:

vim ~/.bashrc

Lisa toimetajaga jÀrgmised read:

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

KÀivita allolev kÀsk pÀrast muudatuste tegemist bashrc-s:

source ~/.bashrc

AWS PostgreSQL'i rakendamine

Peame veel rakendama andmebaasi, kuhu laadime voogudest töödeldud teabe. Selleks kasutame AWS RDS teenust.

Logi sisse AWS konsooli —> AWS RDS —> Andmebaasid —> Loo andmebaas:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Valime PostgreSQL'i ja vajutame nuppu JĂ€rgmine:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Kuna seda nĂ€idet analĂŒĂŒsitakse ainult hariduslikel eesmĂ€rkidel, kasutame tasuta serverit "minimaalses suuruses" (Free Tier):
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

SeejĂ€rel mĂ€rkime Free Tier ploki, ja pĂ€rast seda pakutakse meile automaatselt t2.micro instantsi — kuigi see on nĂ”rk, on see tasuta ja sobib meie ĂŒlesande jaoks hĂ€sti:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

JÀrgmised on vÀga olulised asjad: andmebaasi instantsi nimi, administraator-kasutaja nimi ja tema parool. Nimetame instantsi: myHabrTest, administraator-kasutaja: habr, parool: habr12345 ja vajutame nuppu JÀrgmine:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

JÀrgmisel lehel on seaded, mis vastutavad meie andmebaasi serveri vÀliskÀideldavuse (Avalik ligipÀÀs) ja portide kÀttesaadavuse eest:

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Loome VPC turvagruppide jaoks uue konfiguratsiooni, mis vÔimaldab vÀljastpoolt meie andmebaasi serverisse pÀÀseda kaudu pordi 5432 (PostgreSQL).
Liigume teises brauseriaknas AWS konsooli VPC armatuurlauale —> Turvagruppide haldamine —> Loo turvagrupp:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

MÀÀrame nime Security group – PostgreSQL, kirjeldame seda, tĂ€psustame, millise VPC-ga see grupp peab olema seotud, ja vajutame nuppu Create:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

TĂ€idame vĂ€rskelt loodud grupile Inbound rules sadamale 5432, nagu on nĂ€idatud alloleval pildil. Sadamat ei ole vaja kĂ€kitehnikat mÀÀrata, vaid valida PostgreSQL rippmenĂŒĂŒst Type.

Tegelikult tÀhendab vÀÀrtus ::/0 juurdepÀÀsu sissetulevale liiklusele serverile kogu maailmast, mis ei ole tÀiesti Ôige, kuid nÀite mÔistmiseks lubame endale sellist lÀhenemist:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Naaseme brauseri lehele, kus meil on avatud "Configure advanced settings" ja valime jaotisest VPC security groups –> Choose existing VPC security groups –> PostgreSQL:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

SeejĂ€rel, jaotises Database options –> Database name –> mÀÀrame nime – habrDB.

ÜlejÀÀnud parameetrid, vĂ€lja arvatud vĂ”ib-olla varukoopiate keelamine (backup retention period – 0 days), jĂ€lgimine ja Performance Insights, saame jĂ€tta vaikevÀÀrtustele. Vajutame nuppu Create database:
Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Voogude töötleja

LÔppfaasiks on Spark-i töö arendamine, mis igal kahel sekundil töötleb uusi andmeid, mis on arrivedud Kafka kaudu ja salvestab tulemuse andmebaasi.

Nagu eespool mainitud, on kontrollpunktid (checkpoints) pÔhimehhanism SparkStreamingus, mis peab olema seadistatud talitluse katkemise vÀltimiseks. Kasutame kontrollpunkte ja juhul, kui protseduur kukub kokku, peab Spark Streaming moodul kaotatud andmete taastamiseks lihtsalt naasma viimasele kontrollpunktile ja jÀtkama kalkulatsioone sellelt.

Kontrollpunkti saab lubada, mÀÀrates katalooge usaldusvÀÀrses, vĂ€ljatĂ”rjutavas failisĂŒsteemis (nt HDFS, S3 jne), kus kontrollpunkti teave salvestatakse. Seda tehakse nĂ€iteks:

streamingContext.checkpoint(checkpointDirectory)

Meie nÀites kasutame jÀrgmist lÀhenemist, nimelt kui checkpointDirectory eksisteerib, siis konteksti taastatakse kontrollpunkti andmetest. Kui kataloog ei eksisteeri (st see kÀivitatakse esmakordselt), siis kutsutakse vÀlja funktsioon functionToCreateContext uue konteksti loomiseks ja DStreamide seadistamiseks:

from pyspark.streaming import StreamingContext

context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)

Loome objekti DirectStream, et ĂŒhendada teema «transaction» abil meetodi createDirectStream kaudu, kasutades KafkaUtils-i:

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

Parsime sisenevaid 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-i, teeme lihtsa rĂŒhmitamise ja kuvame tulemuse konsoolile:

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

Saame pÀringu teksti ja kÀivitame selle Spark SQL-is:

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

SeejÀrel salvestame saadud agrigeeritud andmed tabelisse AWS RDS. Tulemuste salvestamiseks andmebaasitabelisse kasutame DataFrame'i 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()

MĂ”ned sĂ”nad AWS RDS-iga ĂŒhenduse seadistamisest. Kasutaja ja parooli lĂ”ime me "AWS PostgreSQL-i juurutamise" etapis. Andmebaasi serveri url-ina tuleks kasutada Endpoint'i, mis kuvatakse Connectivity & security jaotises:

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

Selleks, et Spark ja Kafka Ă”igesti koos töötaksid, tuleks töö kĂ€ivitada spark-submit kaudu artefakti abil spark-streaming-kafka-0-8_2.11. Samuti rakendame PostgreSQL-i andmebaasi suhtlemiseks vajalikku artefakti, mille edastame lĂ€bi —packages.

Skripti paindlikkuse tagamiseks toome ka sÔnumiserveri nime ja teema, kust soovime andmeid saada, sisse sisendparameetriteks.

Nii et on aeg sĂŒsteem kĂ€ivitada ja selle töökindlust 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 Ă”nnestus! Nagu allolevalt pildilt nĂ€ha — rakenduse töötamise ajal kuvatakse uusi agrigeeritud tulemusi iga 2 sekundi jĂ€rel, kuna seadsime paketide intervalli 2 sekundiks, kui loome StreamingContext objekti:

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

SeejÀrel teeme lihtsa pÀringu andmebaasi, et kontrollida tabelis salvestuste olemasolu transaction_flow:

Apache Kafka ja andmete voogude töötlemine Spark Streaming-iga

KokkuvÔte

Selles artiklis kĂ€sitletakse nĂ€idet andmevoo töötlemisest, kasutades Spark Streaming koos Apache Kafka ja PostgreSQL-iga. Andmemahtude suurenedes erinevatest allikatest on raske ĂŒle hinnata Spark Streaming'i praktilist vÀÀrtust voograkenduste ja reaalajas töötavate rakenduste loomisel.

Kogu algkoodi leiate minu repolt aadressilt GitHub.

Olen hea meelega valmis arutama seda artiklit, ootan teie kommentaare ja loodan konstruktiivset kriitikat kÔigilt, kes on huvitatud.

Soovin edu!

Ps. Alguses oli plaanis kasutada kohalikke PostgreSQL andmebaase, kuid arvestades minu armastust AWS-i vastu, otsustasin viia andmebaasi pilve. JĂ€rgmises artiklis selle teema kohta nĂ€itan, kuidas rakendada tĂ€ielikult ĂŒlaltoodud sĂŒsteemi AWS-is, kasutades AWS Kinesis ja AWS EMR. JĂ€lgige uudiseid!

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster