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!

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

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

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

Valime PostgreSQL'i ja vajutame nuppu JĂ€rgmine:

Kuna seda nĂ€idet analĂŒĂŒsitakse ainult hariduslikel eesmĂ€rkidel, kasutame tasuta serverit "minimaalses suuruses" (Free Tier):

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:

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:

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

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:

MÀÀrame nime Security group â PostgreSQL, kirjeldame seda, tĂ€psustame, millise VPC-ga see grupp peab olema seotud, ja vajutame nuppu Create:

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:

Naaseme brauseri lehele, kus meil on avatud "Configure advanced settings" ja valime jaotisest VPC security groups â> Choose existing VPC security groups â> PostgreSQL:

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:

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

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

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

