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!

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

Kasutatavad komponendid:
- â 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;
- â 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;
- â ĂŒ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.
- â 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.propertiesLisage faili lÔppu jÀrgmine:
delete.topic.enable = trueEnne 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.propertiesLoome uue teema nimega Transaction:
bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic transactionVeendume, et teema on loodud vajaliku partitsioonide arvu ja replikatsiooniga:
bin/kafka-topics.sh --describe --zookeeper localhost:2181 
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 â . 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:

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 scalaLaadime 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/sparkLisame Spark'i tee bash-faili:
vim ~/.bashrcLisame 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 ~/.bashrcAWS 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:

Valige PostgreSQL ja klÔpsake nuppu JÀrgmine:

Arvestades, et see nÀide on mÔeldud ainult hariduslikel eesmÀrkidel, kasutame tasuta serverit 'minimaalsel tasemel' (Free Tier):

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:

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:

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

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:

Kasutame turvagrupi jaoks nime â PostgreSQL, lisame kirjelduse, valime, mille VPC-ga see grupp peab olema seotud, ja vajutame nuppu Create:

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:

Naaseme brauseri lehele, kus on avatud âConfigure advanced settingsâ ja valime VPC turvagruppide alt â Choose existing VPC security groups â PostgreSQL:

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:

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:
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:

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

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

