Përshëndetje, Habr!
A ju pĂ«lqen tĂ« fluturoni me avionĂ«? UnĂ« i dua, por gjatĂ« izolimit mĂ«sova gjithashtu tĂ« analizoj tĂ« dhĂ«nat pĂ«r biletat e avionĂ«ve nga njĂ« burim tĂ« njohur â Aviasales.
Sot do të shqyrtojmë funksionimin e Amazon Kinesis, do të ndërtojmë një sistem streaming me analitikë në kohë reale, do të vendosim një bazë të dhënash NoSQL Amazon DynamoDB si magazinën kryesore të të dhënave dhe do të paramendojmë njoftimin përmes SMS për biletat e interesit.
Të gjitha detajet më poshtë! Le të fillojmë!

Hyrje
Për shembull, na nevojitet qasje në . Qasja në të ofrohet falas dhe pa kufizime, thjesht duhet të regjistroheni në seksionin "Zhvilluesit" për të marrë token e API për qasje në të dhënat.
Qëllimi kryesor i këtij artikulli është të japë një kuptim të përgjithshëm të përdorimit të transmetimit të informacionit në AWS, ne e lëmë mënjanë faktin se të dhënat e kthyera nga API i përdorur nuk janë plotësisht të azhurnuara dhe transmetohen nga një memorie që formohet mbi bazën e kërkimeve të përdoruesve në sitet Aviasales.ru dhe Jetradar.com gjatë 48 orëve të fundit.
Të dhënat për biletat e avionit që merren përmes API-së Kinesis-agent, e instaluar në makinën prodhuese, do të parse dhe dërgojë automatikisht në rrjedhën e duhur përmes Kinesis Data Analytics. Versioni i papërpunuar i kësaj rrjedhe do të shkruhet direkt në depo. Depoja e zhvilluar në DynamoDB për të dhënat 'të papërpunuara' do të lejojë analizë më të thellë të biletave përmes mjeteve BI, për shembull, AWS Quick Sight.
Ne do të shqyrtojmë dy opsione për shpërndarjen e gjithë infrastrukturës:
- Dorezimi manual - përmes AWS Management Console;
- Infrastruktura nga kodi Terraform - për automatizuesit e lenë;
Arkitektura e sistemit në zhvillim

Komponentët e përdorur:
- - të dhënat që kthehen nga ky API do të përdoren për të gjitha aktivitetet e mëtejshme;
- - një makinë virtuale e zakonshme në cloud, në të cilën do të gjenerohet rrjedha hyrëse e të dhënave:
- - kjo është një aplikacion Java, i instaluar lokal në makinë, që ofron një mënyrë të thjeshtë për të mbledhur dhe dërguar të dhëna në Kinesis (Kinesis Data Streams ose Kinesis Firehose). Agjenti ndjek vazhdimisht një grup skedarësh në drejtoritë e caktuara dhe dërgon të dhëna të reja në Kinesis;
- - një skript Python, që bën kërkesa në API dhe grumbullon përgjigjen në një dosje që monitoron Kinesis Agent;
- â njĂ« shĂ«rbim i transmetimit tĂ« tĂ« dhĂ«nave nĂ« kohĂ« reale me mundĂ«si tĂ« gjera shkallĂ«zimi;
- â njĂ« shĂ«rbim pa server qĂ« lehtĂ«son analizĂ«n e tĂ« dhĂ«nave tĂ« transmetuara nĂ« kohĂ« reale. Amazon Kinesis Data Analytics rregullon burimet pĂ«r funksionimin e aplikacioneve dhe automatikisht shkallĂ«zohet pĂ«r tĂ« pĂ«rballuar çdo vĂ«llim tĂ« tĂ« dhĂ«nave qĂ« hyn;
- â njĂ« shĂ«rbim qĂ« lejon ekzekutimin e kodit pa rezervim dhe konfigurim serveresh. TĂ« gjitha burimet kompjuterike shkallĂ«zohen automatikisht me çdo thirrje;
- â njĂ« bazĂ« tĂ« dhĂ«nash 'çelĂ«s-vlerĂ«' dhe dokumentesh, e cila siguron njĂ« vonesĂ« mĂ« tĂ« vogĂ«l se 10 milisekonda nĂ« çdo shkallĂ«. Me pĂ«rdorimin e DynamoDB, nuk nevojitet shpĂ«rndarja e serverĂ«ve tĂ« ndonjĂ« lloji, as instalimi i patch-eve apo menaxhimi i tyre. DynamoDB automatikisht shkallĂ«zon tabelat duke rregulluar vĂ«llimin e burimeve tĂ« disponueshme dhe duke ruajtur njĂ« performancĂ« tĂ« lartĂ«. AsnjĂ« veprim pĂ«r administrimin e sistemit nuk Ă«shtĂ« i nevojshĂ«m;
- â njĂ« shĂ«rbim i plotĂ« i menaxhuar pĂ«r dĂ«rgimin e mesazheve sipas modelit "botues - abonues" (Pub/Sub), i cili mundĂ«son izolimin e mikroshĂ«rbimeve, sistemeve tĂ« shpĂ«rndara dhe aplikacioneve pa server. SNS mund tĂ« pĂ«rdoret pĂ«r tĂ« dĂ«rguar informacion te pĂ«rdoruesit e fundit nĂ«pĂ«rmjet njoftimeve mobile push, mesazheve SMS dhe email-eve.
Përgatitja fillestare
PĂ«r tĂ« emuluar njĂ« rrjedhĂ« tĂ« dhĂ«nash, kam vendosur tĂ« pĂ«rdor informacionin mbi biletat e aeroplanit, qĂ« kthehet nga API Aviasales. NĂ« njĂ« listĂ« tĂ« gjerĂ« me metoda tĂ« ndryshme, do tĂ« marrim njĂ« prej tyre â "Kalendari i çmimeve pĂ«r muajin", i cili kthen çmimet pĂ«r çdo ditĂ« tĂ« muajit, tĂ« grumbulluara sipas numrit tĂ« ndĂ«rrimeve. NĂ«se nuk dĂ«rgoni nĂ« kĂ«rkesĂ« muajin e kĂ«rkimit, do tĂ« kthehet informacioni pĂ«r muajin qĂ« vjen pas atij aktual.
Pra, regjistrohemi, marrim token tonë.
Shembulli i kërkesës më poshtë:
http://api.travelpayouts.com/v2/prices/month-matrix?currency=rub&origin=LED&destination=HKT&show_to_affiliates=true&token=TOKEN_APIMetoda e përshkruar më sipër për marrjen e të dhënave nga API me caktimin e tokenit në kërkesë do të funksionojë, por më pëlqen më shumë të kaloj tokenin e aksesit përmes titullit, prandaj në skriptin api_caller.py do të përdorim pikërisht këtë metodë.
Shembulli i përgjigjes:
{{
"success":true,
"data":[{
"show_to_affiliates":true,
"trip_class":0,
"origin":"LED",
"destination":"HKT",
"depart_date":"2015-10-01",
"return_date":"",
"number_of_changes":1,
"value":29127,
"found_at":"2015-09-24T00:06:12+04:00",
"distance":8015,
"actual":true
}]
}
Në shembullin e përgjigjes API më sipër tregohet një biletë nga Shën Petersburg në Phuket... Eh, çfarë ëndrrash...
Duke qenë se jam nga Kazani dhe Phuket tani "na ëndërron vetëm", le të kërkojmë bileta nga Shën Petersburg në Kazan.
Kuptohet se ju tashmë keni një llogari në AWS. Dëshiroj të tërheq vëmendjen tuaj specifikisht, që Kinesis dhe dërgimi i njoftimeve përmes SMS nuk përfshihen në . Por edhe pse kjo është e vërtetë, duke menduar disa dollarë, është krejtësisht e mundur të ndërtohet sistemi i propozuar dhe të luhet me të. Sigurisht, mos harroni të fshini të gjitha burimet pasi të mos ju nevojiten më.
Me fat, DynamoDb dhe funksionet lambda do të jenë relativisht falas për ne, nëse qëndrojmë brenda limiteve mujore falas. Për shembull, për DynamoDB: 25 GB ruajtje, 25 WCU/RCU dhe 100 milion kërkesa. Dhe një milion thirrje funksionesh lambda në muaj.
Depoj i dorës së sistemit
Konfigurimi i Kinesis Data Streams
Të kalojmë te shërbimi Kinesis Data Streams dhe të krijojmë dy kanale të reja me një shard në secilin.
ĂfarĂ« Ă«shtĂ« njĂ« shard?
Një shard është një njësi kryesore e transferimit të të dhënave në Amazon Kinesis. Një segment ofron transferim të të dhënave hyrëse me shpejtësi 1 MB/s dhe transferim të të dhënave dalëse me shpejtësi 2 MB/s. Një segment mbështet deri në 1000 regjistrime PUT në sekondë. Kur krijoni një kanal të dhënash, duhet të specifikoni numrin e nevojshëm të segmenteve. Për shembull, mund të krijoni një kanal të dhënash me dy segmente. Ky kanal të dhënash do të ofrojë transferim të të dhënave hyrëse me shpejtësi 2 MB/s dhe transferim të të dhënave dalëse me shpejtësi 4 MB/s, duke mbështetur deri në 2000 regjistrime PUT në sekondë.
Sa mĂ« shumĂ« sharde tĂ« keni nĂ« kanalin tuaj, aq mĂ« e lartĂ« Ă«shtĂ« kapaciteti i tij. Pra, kanalet zmadhohen duke shtuar sharde. Por sa mĂ« shumĂ« sharde tĂ« keni, aq mĂ« e lartĂ« Ă«shtĂ« kostoja. Ădo shard kushton 1,5 cent nĂ« orĂ« dhe pĂ«r mĂ« tepĂ«r 1.4 cent pĂ«r çdo milion operacione shtimi nĂ« kanal (PUT payload units).
Le të krijojmë një kanal të ri me emrin airline_tickets, i cili do të ketë mjaft 1 shard:

Tani le të krijojmë një tjetër kanal me emrin special_stream:

Konfigurimi i prodhuesit
Si producent të dhënash për të analizuar detyrën, mjafton të përdorni një instancë EC2 të zakonshme. Nuk duhet të jetë një makinë virtuale e fuqishme dhe e kushtueshme, një t2.micro spot është plotësisht e mjaftueshme.
NjĂ« shĂ«nim i rĂ«ndĂ«sishĂ«m: pĂ«r shembullin, duhet tĂ« pĂ«rdorni imazhin â Amazon Linux AMI 2018.03.0, me tĂ« cilin ka mĂ« pak konfigurime pĂ«r tĂ« nisur shpejt Kinesis Agent.
Kalohet te shërbimi EC2, krijohet një makinë virtuale e re, zgjidhet AMI i duhur me tipin t2.micro, i cili përfshihet në Free Tier:

Që makina virtuale e krijuar së fundmi të mund të bashkëpunojë me shërbimin Kinesis, është e nevojshme t'i jepet të drejta. Mënyra më e mirë për ta bërë këtë është të caktoni një IAM Role. Prandaj, në ekranin Step 3: Configure Instance Details duhet të zgjidhet Krijo IAM të re:
Krijimi i rolit IAM për EC2

Në dritaren e hapur, zgjidhni se roli i ri krijohet për EC2 dhe kaloni në seksionin Permissions:

Në shembullin e mësimit, mund të mos merremi me të gjitha thellësitë e konfigurimit të detajuar të të drejtave për burimet, kështu që do të zgjedhim politikat e parakonfiguruara nga Amazon: AmazonKinesisFullAccess dhe CloudWatchFullAccess.
Do t'i japim një emër me kuptim këtij roli, për shembull: EC2-KinesisStreams-FullAccess. Si rezultat, duhet të dalë e njëjta gjë siç tregohet në imazhin më poshtë:

Pas krijimit të këtij roli të ri, mos harro të lidhe atë me instancën e makinës virtuale që po krijon:

Nuk kemi asgjë tjetër për të ndryshuar në këtë ekran dhe kalojmë në dritaret e tjera.
Parametrat e hard diskut mund t'i lëmë siç janë, po ashtu edhe etiketat (megjithatë, një praktikë e mirë është të përdorim etiketa, të paktën duke i dhënë emër instancës dhe duke treguar ambientin).
Tani jemi nĂ« skedĂ«n Hapi 6: Konfiguro Grupi i SigurisĂ«, ku duhet tĂ« krijojmĂ« njĂ« tĂ« ri ose tĂ« specifikojmĂ« atĂ« qĂ« keni pĂ«r grupin e sigurisĂ«, qĂ« lejon lidhjen pĂ«rmes ssh (porti 22) nĂ« instancĂ«. Zgjidhni aty Burimin â> IP-ja ime dhe mund tĂ« nisin instancĂ«n.

Sapo të kalojë në statusin running, mund të provoni të lidheni me të përmes ssh.
Për të mundësuar funksionimin me Kinesis Agent, pas lidhjes së suksesshme me makinën, duhet të shkruani komandat e mëposhtme në terminal:
sudo yum -y update
sudo yum install -y python36 python36-pip
sudo /usr/bin/pip-3.6 install --upgrade pip
sudo yum install -y aws-kinesis-agent
Le të krijojmë një dosje për ruajtjen e përgjigjeve të API:
sudo mkdir /var/log/airline_ticketsPara se të nisim agjentin, duhet të konfigurojmë skedarin e tij:
sudo vim /etc/aws-kinesis/agent.jsonPërmbajtja e skedarit agent.json duhet të ketë formën e mëposhtme:
{
"cloudwatch.emitMetrics": true,
"kinesis.endpoint": "",
"firehose.endpoint": "",
"flows": [
{
"filePattern": "/var/log/airline_tickets/*log",
"kinesisStream": "airline_tickets",
"partitionKeyOption": "RANDOM",
"dataProcessingOptions": [
{
"optionName": "CSVTOJSON",
"customFieldNames": ["cost","trip_class","show_to_affiliates",
"return_date","origin","number_of_changes","gate","found_at",
"duration","distance","destination","depart_date","actual","record_id"]
}
]
}
]
}
Si duket në skedarin e konfigurimit, agjenti do të monitorojë në direktorinë /var/log/airline_tickets/ skedarët me zgjedhjen .log, do t'i analizojë ato dhe do t'i dërgojë në rrjedhën airline_tickets.
Rinisim shërbimin dhe sigurohemi që ai është aktive dhe po punon:
sudo service aws-kinesis-agent restartTani, do të shkarkojmë skriptin Python, i cili do të kërkojë të dhëna nga API:
REPO_PATH=https://raw.githubusercontent.com/igorgorbenko/aviasales_kinesis/master/producer
wget $REPO_PATH/api_caller.py -P /home/ec2-user/
wget $REPO_PATH/requirements.txt -P /home/ec2-user/
sudo chmod a+x /home/ec2-user/api_caller.py
sudo /usr/local/bin/pip3 install -r /home/ec2-user/requirements.txt
Skripti api_caller.py kërkon të dhëna nga Aviasales dhe ruan përgjigjen e marrë në direktorinë që skanon agjenti Kinesis. Realizimi i këtij skripti është mjaft standard, ka klasën TicketsApi, e cila lejon të qasemi asinkronisht në API. Në këtë klasë ne jemi duke dërguar një titull me tokenin dhe parametrat e kërkesës:
class TicketsApi:
"""Klasa e thirrjes së API-së."""
def __init__(self, headers):
"""Metoda e inicializimit."""
self.base_url = BASE_URL
self.headers = headers
async def get_data(self, data):
"""Merr të dhënat nga kërkesa e API-së."""
response_json = {}
async with ClientSession(headers=self.headers) as session:
try:
response = await session.get(self.base_url, data=data)
response.raise_for_status()
LOGGER.info('Statusi i përgjigjes %s: %s',
self.base_url, response.status)
response_json = await response.json()
except HTTPError as http_err:
LOGGER.error('Ups! Ndodhi një gabim HTTP: %s', str(http_err))
except Exception as err:
LOGGER.error('Ups! Ndodhi një gabim: %s', str(err))
return response_json
def prepare_request(api_token):
"""Kthehet headers dhe kërkesa për thirrjen API."""
headers = {'X-Access-Token': api_token,
'Accept-Encoding': 'gzip'}
data = FormData()
data.add_field('currency', CURRENCY)
data.add_field('origin', ORIGIN)
data.add_field('destination', DESTINATION)
data.add_field('show_to_affiliates', SHOW_TO_AFFILIATES)
data.add_field('trip_duration', TRIP_DURATION)
return headers, data
async def main():
"""Merr kodin për të ekzekutuar."""
if len(sys.argv) != 2:
print('Përdorimi: api_caller.py ')
sys.exit(1)
return
api_token = sys.argv[1]
headers, data = prepare_request(api_token)
api = TicketsApi(headers)
response = await api.get_data(data)
if response.get('success', None):
LOGGER.info('API ka kthyer %s objekte', len(response['data']))
try:
count_rows = log_maker(response)
LOGGER.info('%s rreshta janë ruajtur në %s',
count_rows,
TARGET_FILE)
except Exception as e:
LOGGER.error('Ups! Rezultati i kërkesës nuk u ruajt në skedarin. %s',
str(e))
else:
LOGGER.error('Ups! Kërkesa për API dështoi %s!', response)
Për të testuar saktësinë e konfigurimeve dhe funksionalitetin e agjentit, do të bëjmë një provë lançimi me skriptin api_caller.py:
sudo ./api_caller.py TOKEN 
Dhe shikojmë rezultatin e funksionimit në të dhënat e Agjentit dhe në skedën Monitoring në rrjedhën e të dhënave airline_tickets:
tail -f /var/log/aws-kinesis-agent/aws-kinesis-agent.log 

Siç duket, gjithçka funksionon dhe Agjenti Kinesis dërgon me sukses të dhëna në rrjedhë. Tani do të konfigurojmë konsumatorin.
Konfigurimi i Kinesis Data Analytics
TĂ« kalojmĂ« nĂ« komponentin kryesor tĂ« gjithĂ« sistemit â do tĂ« krijojmĂ« njĂ« aplikacion tĂ« ri nĂ« Kinesis Data Analytics me emrin kinesis_analytics_airlines_app:

Kinesis Data Analytics lejon ekzekutimin e analizës së të dhënave në kohë reale nga Kinesis Streams duke përdorur gjuhën SQL. Ky është një shërbim plotësisht automatik (ndryshe nga Kinesis Streams), i cili:
- lejon krijimin e rrjedhave të reja (Output Stream) në bazë të pyetjeve për të dhënat burimore;
- siguron një rrjedhë me gabime që kanë ndodhur gjatë funksionimit të aplikacioneve (Error Stream);
- ka aftësinë për të identifikuar automatikisht skemën e të dhënave hyrëse (mund të rregullohet manualisht nëse është e nevojshme).
Ky shĂ«rbim nuk Ă«shtĂ« i lirĂ« â 0.11 USD nĂ« orĂ« punĂ«, prandaj duhet tĂ« pĂ«rdoret me kujdes dhe tĂ« fshihet kur tĂ« pĂ«rfundojĂ« punĂ«n.
Do të lidhi aplikacionin me burimin e të dhënave:

Zgjidhni rrjedhën në të cilën do të lidhni (airline_tickets):

Për më pas, është e nevojshme të bashkangjitni një Rol të ri IAM në mënyrë që aplikacioni të mund të lexojë dhe shkruajë në rrjedhë. Mjafton të mos ndryshoni asgjë në bllokun e Lejeve të Qasjes:

Tani do të kërkojmë zb discover schema të të dhënave në rrjedhë, për këtë, klikoni në butonin 'Discover schema'. Si rezultat, do të përditësohet (do të krijohet një e re) Roli IAM dhe do të nisë zb discovery schema nga të dhënat që tashmë kanë mbërritur në rrjedhë:

Tani Ă«shtĂ« e nevojshme tĂ« kaloni nĂ« redaktorin SQL. Kur tĂ« klikoni nĂ« kĂ«tĂ« buton, do tĂ« shfaqet njĂ« dritare me pyetje pĂ«r tĂ« nisur aplikacionin â zgjidhni atĂ« qĂ« dĂ«shironi tĂ« nisni:

Në dritaren e redaktorit SQL, vendosni këtë kërkesë të thjeshtë dhe klikoni Save and Run SQL:
CREATE OR REPLACE STREAM "DESTINATION_SQL_STREAM" ("cost" DOUBLE, "gate" VARCHAR(16));
CREATE OR REPLACE PUMP "STREAM_PUMP" AS INSERT INTO "DESTINATION_SQL_STREAM"
SELECT STREAM "cost", "gate"
FROM "SOURCE_SQL_STREAM_001"
WHERE "cost" < 5000
and "gate" = 'Aeroflot';
NĂ« databazat relacional, punoni me tabela duke pĂ«rdorur operatorĂ«t INSERT pĂ«r tĂ« shtuar regjistrime dhe operatorin SELECT pĂ«r tĂ« kĂ«rkuar tĂ« dhĂ«na. NĂ« Amazon Kinesis Data Analytics, punoni me rrjedha (STREAM) dhe "pumpa" (PUMP) â kĂ«rkesa tĂ« vazhdueshme tĂ« futjes qĂ« fusin tĂ« dhĂ«na nga njĂ« rrjedhĂ« nĂ« aplikacion nĂ« njĂ« tjetĂ«r rrjedhĂ«.
Në kërkesën SQL të paraqitur më sipër, ndodhet një kërkim për biletat e Aeroflotit me çmim nën pesë mijë rubla. Të gjitha regjistrimet që bien nën këto kushte do të vendosen në rrjedhën DESTINATION_SQL_STREAM.

Në bllokun Destination, zgjedhim rrjedhën special_stream, dhe në listën e rënëshme emrin e rrjedhës brenda aplikacionit DESTINATION_SQL_STREAM:

Si rezultat i të gjitha manipulimeve, duhet të kemi diçka të ngjashme me imazhin më poshtë:

Krijimi dhe abonimi në temën SNS
Kalo në shërbimin Simple Notification Service dhe krijo një temë të re me emrin Airlines:

Krijo një abonim për këtë temë, ku specifikon numrin e telefonit mobil, në të cilin do të vijnë njoftimet SMS:

Krijimi i tabelës në DynamoDB
Për ruajtjen e të dhënave të papërpunuara nga rrjedha airline_tickets, do të krijojmë një tabelë në DynamoDB me të njëjtin emër. Si çelës primar do të përdorim record_id:

Krijimi i funksionit lambda collector
Do tĂ« krijojmĂ« njĂ« funksion lambda me emrin Collector, i cili do tĂ« ketĂ« detyrĂ« tĂ« pyetĂ«sisĂ« pĂ«r rrjedhĂ«n airline_tickets dhe, nĂ« rast se gjenden aty shĂ«nime tĂ« reja, t'i fusĂ« kĂ«to shĂ«nime nĂ« tabelĂ«n DynamoDB. ĂshtĂ« e qartĂ« se pĂ«rveç lejeve tĂ« paracaktuara, kjo lambda duhet tĂ« ketĂ« qasje pĂ«r tĂ« lexuar rrjedhĂ«n e tĂ« dhĂ«nave Kinesis dhe tĂ« shkruajĂ« nĂ« DynamoDB.
Krijimi i rolit IAM për funksionin lambda collector
Së pari, le të krijojmë një rol të ri IAM për lambda me emrin Lambda-TicketsProcessingRole:

Për një shembull testimi, politikat e paracaktuara AmazonKinesisReadOnlyAccess dhe AmazonDynamoDBFullAccess do të ishin të mjaftueshme, siç tregohet në imazhin më poshtë:


Kjo lambda duhet të startohet nga një alarm i Kinesis kur shënohen shënime të reja në rrjedhën airline_stream, prandaj duhet të shtojmë një alarm të ri:


E mbetur është të fusim kodin dhe të ruajmë lambdën.
"""Parsing the stream and inserting into the DynamoDB table."""
import base64
import json
import boto3
from decimal import Decimal
DYNAMO_DB = boto3.resource('dynamodb')
TABLE_NAME = 'airline_tickets'
class TicketsParser:
"""Parsing info from the Stream."""
def __init__(self, table_name, records):
"""Init method."""
self.table = DYNAMO_DB.Table(table_name)
self.json_data = TicketsParser.get_json_data(records)
@staticmethod
def get_json_data(records):
"""Return deserialized data from the stream."""
decoded_record_data = ([base64.b64decode(record['kinesis']['data'])
for record in records])
json_data = ([json.loads(decoded_record)
for decoded_record in decoded_record_data])
return json_data
@staticmethod
def get_item_from_json(json_item):
"""Pre-process the json data."""
new_item = {
'record_id': json_item.get('record_id'),
'cost': Decimal(json_item.get('cost')),
'trip_class': json_item.get('trip_class'),
'show_to_affiliates': json_item.get('show_to_affiliates'),
'origin': json_item.get('origin'),
'number_of_changes': int(json_item.get('number_of_changes')),
'gate': json_item.get('gate'),
'found_at': json_item.get('found_at'),
'duration': int(json_item.get('duration')),
'distance': int(json_item.get('distance')),
'destination': json_item.get('destination'),
'depart_date': json_item.get('depart_date'),
'actual': json_item.get('actual')
}
return new_item
def run(self):
"""Batch insert into the table."""
with self.table.batch_writer() as batch_writer:
for item in self.json_data:
dynamodb_item = TicketsParser.get_item_from_json(item)
batch_writer.put_item(dynamodb_item)
print('Has been added ', len(self.json_data), 'items')
def lambda_handler(event, context):
"""Parse the stream and insert into the DynamoDB table."""
print('Got event:', event)
parser = TicketsParser(TABLE_NAME, event['Records'])
parser.run()
Krijimi i funksionit lambda notifier
Funksioni i dytë lambda, i cili do të monitorojë rrjedhën e dytë (special_stream) dhe do të dërgojë njoftime në SNS, krijohet në mënyrë të ngjashme. Pra, ky lambda duhet të ketë akses për të lexuar nga Kinesis dhe të dërgojë mesazhe në temën e caktuar SNS, e cila më pas do të dërgohet nga shërbimi SNS për të gjithë abonentët e kësaj teme (email, SMS, etj).
Krijimi i rolit IAM
Fillimisht krijojmë rolin IAM Lambda-KinesisAlarm për këtë lambda dhe pastaj e caktojmë këtë rol për lambda-n që krijohet alarm_notifier:


Ky lambda duhet të funksionojë në bazë të aktivizimit për hyrjen e shënimeve të reja në rrjedhën special_stream, prandaj është e nevojshme të konfigurohet aktivizimi në mënyrë të ngjashme me atë që bëmë për lambda-n Collector.
PĂ«r lehtĂ«simin e konfigurimit tĂ« kĂ«tij lambda, do tĂ« prezantojmĂ« njĂ« variabĂ«l tĂ« re ambientale â TOPIC_ARN, ku do tĂ« vendosim ANR (Emrat e Burimeve Amazon) tĂ« temĂ«s Airlines:

Dhe vendosim kodin e lambda-s, ai është krejtësisht i thjeshtë:
import boto3
import base64
import os
SNS_CLIENT = boto3.client('sns')
TOPIC_ARN = os.environ['TOPIC_ARN']
def lambda_handler(event, context):
try:
SNS_CLIENT.publish(TopicArn=TOPIC_ARN,
Message='Përshëndetje! Kam gjetur një gjë interesante!',
Subject='Alarm për biletat e aviacionit')
print('Mesazhi i alarmit është dërguar me sukses')
except Exception as err:
print('Dështim në dërgesë', str(err))
Duket se konfigurimi manual i sistemit ka përfunduar. Tani mbetet vetëm të testojmë dhe të sigurohemi që kemi konfiguruar gjithçka siç duhet.
Implementimi nga kodi Terraform
Përgatitja e nevojshme
â Ă«shtĂ« njĂ« mjet shumĂ« i pĂ«rshtatshĂ«m open-source pĂ«r ndĂ«rtimin e infrastrukturĂ«s nga kodi. Ka njĂ« sintaksĂ« tĂ« vetĂ«n, e cila Ă«shtĂ« e lehtĂ« pĂ«r t'u mĂ«suar dhe shumĂ« shembuj se si dhe çfarĂ« tĂ« implementoni. NĂ« redaktorin Atom ose Visual Studio Code ka shumĂ« pluginĂ« tĂ« dobishĂ«m qĂ« lehtĂ«sojnĂ« punĂ«n me Terraform.
Distribucioni mund të shkarkohet . Shpjegimi i plotë i të gjitha mundësive të Terraform kalon përmasat e këtij artikulli, prandaj do të ndalemi te aspektet kryesore.
Si të nisni
Kodi i plotë i projektit ndodhet . Klononi depozitat për vete. Para se të nisin, sigurohuni që keni të instaluar dhe të konfiguruar AWS CLI, pasi Terraform do të kërkojë kredencialet në skedarin ~/ .aws/credentials.
Një praktikë e mirë është të ekzekutoni komandën plan para se të implementoni tërë infrastrukturën, për të parë se çfarë do të krijojë Terraform aktualisht në cloud:
terraform.exe planDo të proponohet të jepni numrin e telefonit për të marrë njoftime. Në këtë fazë, dhënia e tij nuk është e detyrueshme.

Pas analizimit të planit të punës së programit, mund të nisim krijimin e burimeve:
terraform.exe applyPas dërgimit të kësaj komande, do të shfaqet përsëri një kërkesë për të futur numrin e telefonit, shkruani "yes" kur të shfaqet pyetja për realizimin real të veprimeve. Kjo do të mundësojë ngritjen e gjithë infrastrukturës, konfiguroj të gjitha cilësimet e nevojshme EC2, të deploy funksionet lambda etj.
Pas krijimit të suksesshëm të të gjitha burimeve përmes kodit Terraform, duhet të hyjmë në detajet e aplikacionit Kinesis Analytics (fatkeqësisht, nuk e gjetëm se si ta bëjmë këtë menjëherë nga kodi).
Nisim aplikacionin:

Pas kësaj, duhet të caktoni me qartë emrin e stream-it brenda aplikacionit, duke zgjedhur nga lista e rënies:


Tani gjithçka është gati për punë.
Testimi i funksionit të aplikacionit
Pavarësisht se si e keni deployuar sistemin, me dorë apo përmes kodit Terraform, ai do të funksionojë njësoj.
Hyni përmes SSH në makinë virtuale EC2, ku është instaluar Kinesis Agent, dhe nisni skriptin api_caller.py
sudo ./api_caller.py TOKENTani duhen pritur SMS në numrin tuaj:

SMS â mesazhi arrin nĂ« telefon pothuajse brenda 1 minute:

Mbështetje për të parë nëse regjistrimet e bazës së të dhënave DynamoDB janë ruajtur për analizë më të hollësishme. Tabela airline_tickets përmban të dhëna të ngjashme si këto:

Përfundimi
Gjatë punës së kryer, u ndërtua një sistem online për përpunimin e të dhënave mbi bazën e Amazon Kinesis. U shqyrtuan opsionet për përdorimin e Kinesis Agent në lidhje me Kinesis Data Streams dhe analitikën në kohë reale Kinesis Analytics duke përdorur komandat SQL, si dhe ndërveprimin e Amazon Kinesis me shërbime të tjera të AWS.
Sistemin e përmendur më sipër e kemi vendosur në dy mënyra: një manuale e gjatë dhe një të shpejtë nga kodi Terraform.
I gjithë kodi burimor i projektit është i disponueshëm , ju ftoj ta shqyrtoni.
Jam i gatshëm të diskutoj artikullin, po pres komentet tuaja. Shpresoj për kritika konstruktive.
Ju uroj sukses!
Burimi: habr.com
