Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Përshëndetje, Habr!

A ju pëlqen të flisni me avionë? Unë e adhuroj, por gjatë vetëizolimit fillova edhe 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 analizë në kohë reale, do të vendosim bazën e të dhënave NoSQL Amazon DynamoDB si ruajtjen kryesore të të dhënave dhe do të konfigurojmë njoftime përmes SMS për biletat interesante.

Të gjitha detajet më poshtë! Le të fillojmë!

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Hyrje

Për shembullin tonë do të kemi nevojë për qasje në API-në e Aviasales. Qasja në të ofrohet falas dhe pa kufizime, është e nevojshme vetëm të regjistroheni në seksionin "Për Programuesit" për të marrë token tuaj API për qasje në të dhëna.

Qëllimi kryesor i këtij artikulli është të japë një kuptim të përgjithshëm për përdorimin e transmetimit të informacionit në AWS, ne nuk do të flasim për të dhënat e kthyer nga API i përdorur, të cilat nuk janë domosdoshmërisht të sakta dhe dërgohen nga cache që formohet në bazë të kërkesave të përdoruesve në sitet Aviasales.ru dhe Jetradar.com gjatë 48 orëve të fundit.

Të dhënat e marra përmes API-së për biletat e avionëve do të parsehen automatikisht nga Kinesis-agent, i instaluar në makinën prodhuese, dhe do të dërgohen në rrjedhën përkatëse përmes Kinesis Data Analytics. Versioni i papërpunuar i kësaj rrjedhe do të shkruhet drejtpërdrejt në ruajtje. Ruajtja e ‘papërpunuar’ e vendosur në DynamoDB do të lejojë një analizë më të thellë të biletave përmes mjeteve BI, për shembuj, AWS Quick Sight.

Ne do të shqyrtojmë dy mundësi për implementimin e gjithë infrastrukturës:

  • Dora — përmes AWS Management Console;
  • Infrastruktura nga kodi Terraform — për automatizuesit e lenë;

Arkitektura e sistemit në zhvillim

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Komponentët e përdorur:

  • API Aviasales — të dhënat që kthehen nga ky API do të përdoren për të gjithë punën e mëtejshme;
  • EC2 Prodhuesi i Instancave — një makinë virtuale normale në cloud, mbi të cilën do të krijohet rrjedha hyrëse e të dhënave:
    • Kinesis Agent — ky është një aplikacion Java, i instaluar lokal në makinë, i cili 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 monitoron vazhdimisht një grup skedarësh në direktorët e caktuar dhe dërgon të dhëna të reja në Kinesis;
    • Skripti i Thirrësit të API-së — një skript Python, i cili bën kërkesa në API dhe ruan përgjigjen në një dosje, të cilën e monitoron Kinesis Agent;
  • Kinesis Data Streams — shërbim transmetimi të dhënash në kohë reale me mundësi të gjera për shkallëzim;
  • Kinesis Analytics — shërbim pa server, i cili thjeshton analizën e të dhënave në kohë reale. Amazon Kinesis Data Analytics konfigurón burimet për të punuar me aplikacionet dhe automatikisht shkallëzohet për të përpunuar çdo sasi të dhënash në hyrje;
  • AWS Lambda — shërbim që lejon ekzekutimin e kodit pa rezervuar dhe konfiguruar serverë. Të gjitha kapacitetet llogaritëse shkallëzohen automatikisht për çdo thirje;
  • Amazon DynamoDB — databazë e çiftit "çelës-vlerë" dhe dokumenteve, e cila siguron një vonesë më të vogël se 10 milisekonda, pavarësisht nga shkalla. Me përdorimin e DynamoDB nuk është e nevojshme të shpërndahen serverë, të instalohen patches ose të menaxhohen ato. DynamoDB automatikisht shkallëzon tabelat, duke rregulluar sasinë e burimeve të disponueshme dhe ruajtur performancën e lartë. Nuk kërkohen veprime për administrimin e sistemit;
  • Amazon SNS — shërbim i plotë menaxhimi për dërgimin e mesazheve sipas modelit "botues - abonent" (Pub/Sub), me ndihmën e të cilit mund të izoloheshin mikro-shërbimet, sistemet e shpërndara dhe aplikacionet pa server. SNS mund të përdoret për të dërguar informacion te përdoruesit përfundimtarë përmes njoftimeve mobile, mesazheve SMS dhe email-eve.

Përgatitja fillestare

Për të simuluar një rrjedhë të dhënash, kam vendosur të përdor informacionin mbi biletat e avionit, të kthehet nga API Aviasales. Në dokumentacionin një listë mjaft të gjerë metodash të ndryshme, le të marrim një nga ato — "Kalendari i çmimeve për një muaj", i cili kthen çmimet për çdo ditë të muajit, të grupezuara sipas numrit të ndërrimeve. Nëse nuk dërgohet muaji i kërkimit në kërkesë, do të kthehet informacion për muajin që vjen pas aktualit.

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_API

Mënyra e përshkruar për të marrë të dhëna nga API duke përfshirë tokenin në kërkesë do të funksionojë, por më pëlqen më shumë të dërgoj tokenin e aksesit përmes header-it, kështu që në skenarin api_caller.py do ta përdorim këtë mënyrë.

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 së API-të më lart tregohet një biletë nga Shën Petersburgu për në Phuket... Ah, çfarë ëndrrash...
Duke qenë se unë jam nga Kazani dhe Phuket tani është një ëndërr për ne, le të kërkojmë bileta nga Shën Petersburgu në Kazan.

Supozimi është se ju tashmë keni një llogari në AWS. Dëshiroj të theksoj se Kinesis dhe dërgimi i njoftimeve përmes SMS nuk përfshihen në ofertën vjetore Free Tier (përdorim falas). Por edhe përkundër kësaj, duke marrë parasysh disa dollarë, është plotësisht e mundur të ndërtosh sistemin e propozuar dhe të eksperimentosh me të. Dhe, natyrisht, nuk duhet harrosh të fshish të gjitha burimet pasi të mos jenë më të nevojshme.

Fatmirësisht, DynamoDb dhe funksionet lambda do të jenë për ne në mënyrë të kushtëzuar falas nëse mbetet brenda kufijve mujorë falas. Për shembull, për DynamoDB: 25 GB hapësirë ruajtjeje, 25 WCU/RCU dhe 100 milion kërkesash. Dhe një milion thirrje funksionesh lambda në muaj.

Depo me dorë të sistemit

Konfigurimi i Kinesis Data Streams

Shkoni te shërbimi Kinesis Data Streams dhe krijoni dy rrjedha të reja me një shard për secilën.

Çfarë është një shard?
Shard është njësi kryesore për transmetimin e të dhënave në rrjedhën Amazon Kinesis. Një segment siguron transmetimin e të dhënave hyrëse me shpejtësi 1 MB/s dhe transmetimin e 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ë rrjedhë të dhënash, duhet të specifikoni numrin e dëshiruar të segmenteve. Për shembull, mund të krijoni një rrjedhë të dhënash me dy segmente. Kjo rrjedhë e dhënash do të sigurojë transmetimin e të dhënave hyrëse me shpejtësi 2 MB/s dhe transmetimin e 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 në rrjedhën tuaj - aq më shumë kapacitet ka. Në thelb, kështu masivizohen rrjedhat - duke shtuar sharde. Por sa më shumë sharde, aq më e lartë është çmimi. Çdo shard kushton 1.5 cent në orë dhe përveç kësaj 1.4 cent për çdo milion operacione që shtohen në rrjedhë (njesi ngarkimi PUT).

Le të krijojmë një rrjedhë të re me emrin airline_tickets, një shard është mjaft i mjaftueshëm:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Tani le të krijojmë një tjetër rrjedhë me emrin special_stream:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Konfigurimi i prodhuesit

Si prodhues të dhënash për të trajtuar problemin, mjafton të përdorim një instancë të zakonshme EC2. Nuk duhet të jetë një makinë virtuale e fuqishme dhe e shtrenjtë, një t2.micro spote do të ishte më thanë e mjaftueshme.

Shënim i rëndësishëm: për shembull duhet të përdorni image — Amazon Linux AMI 2018.03.0, me të cilin ka më pak konfigurime për një lansim të shpejtë të Kinesis Agent.

Shkkojmë te shërbimi EC2, krijojmë një makinë virtuale të re, zgjedhim AMI-në e nevojshme me llojin t2.micro, i cili përfshihet në Free Tier:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Për të siguruar që makina virtuale e sapokrijuar të mund të bashkëpunojë me shërbimin Kinesis, është e nevojshme t'i japim asaj të drejtat përkatëse. Mënyra më e mirë për ta bërë këtë është të caktuara një IAM Role. Prandaj, në ekranin Step 3: Configure Instance Details duhet të zgjidhet Krijoni rol të ri IAM:

Krijimi i rolit IAM për EC2
Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Në dritaren që hapet, zgjedhim që rolin e ri ta krijojmë për EC2 dhe kalojmë te seksioni Permissions:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Në shembullin mësimor mund të mos merremi me të gjitha nuancat e konfigurimit granular të të drejtave mbi burimet, prandaj do të zgjedhim politikat e paracaktuara nga Amazon: AmazonKinesisFullAccess dhe CloudWatchFullAccess.

Do t'i japim një emër kuptimplotë këtij roli, për shembull: EC2-KinesisStreams-FullAccess. Si rezultat, duhet të kemi të njëjtin rezultat siç është treguar në figurën më poshtë:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Pas krijimit të këtij roli të ri, mos harroni ta lidhni atë me instancën e makinë virtuale që po krijoni:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Nuk ndryshojmë asgjë tjetër në këtë ekran dhe kalojmë në ekranet e tjera.

Opsionet e diskut të fortë mund t'i lëmë siç janë nga fabrika, etiketat gjithashtu (në një praktikë të mirë është të përdoren etiketat, të paktën për të dhënë emër instancës dhe të tregoni ambientin).

Tani jemi në skedën Step 6: Configure Security Group, ku duhet të krijojmë një grup të ri të sigurisë ose të specifikojmë një grup të sigurisë ekzistues që lejon lidhjen përmes ssh (porta 22) në instancë. Zgjidhni aty Source —> My IP dhe mund të filloni instancën.

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Sapo të kalojë në statusin running, mund të provoni të lidheni me të përmes ssh.

Për të pasur mundësinë e punës me Kinesis Agent, pas lidhjes së suksesshme me makinën, është e nevojshme të jepni 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

Të krijojmë një dosje për ruajtjen e përgjigjeve nga API:

sudo mkdir /var/log/airline_tickets

Para se të nisni agjentin, është e nevojshme të konfiguroni konfigurimin e tij:

sudo vim /etc/aws-kinesis/agent.json

Përmbajtja e skedarit agent.json duhet të ketë këtë format:

{
  "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ç poeohet nga skedari i konfigurimit, agjenti do të monitorojë në direktorinë \/var\/log\/airline_tickets\/ skedarët me zgjerim .log, do t'i analizojë dhe do t'i dërgojë në kanalin airline_tickets.

Rindizim shërbimin dhe sigurohemi që ai ka filluar dhe po funksionon:

sudo service aws-kinesis-agent restart

Tani do të shkarkojmë skriptin Python që do të kërkojë të dhënat 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. Zbatimi i këtij skripti është mjaft standard, ka një klasë TicketsApi që lejon të bëhen thirrje asinkrone në API. Në këtë klasë kalojmë kryemin e tokenit dhe parametrat e kërkesës:

class TicketsApi:
    """Klasë e thirrjes së API."""

    def __init__(self, headers):
        """Metoda e inicializimit."""
        self.base_url = BASE_URL
        self.headers = headers

    async def get_data(self, data):
        """Merrni të dhënat nga kërkesa API."""
        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! Ka ndodhur një gabim HTTP: %s', str(http_err))
            except Exception as err:
                LOGGER.error('Ups! Ka ndodhur një gabim: %s', str(err))
            return response_json


def prepare_request(api_token):
    """Kthejini headers dhe kërkesën 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():
    """Kthejeni kodin në ekzekutim."""
    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 artikuj', 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 API ishte e pasuksesshme %s!', response)

Për të testuar saktësinë e konfigurimeve dhe funksionimin e agjentit, do të bëjmë një ekzekutim testues të skenarit api_caller.py:

sudo ./api_caller.py TOKEN

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Dhe shikojmë rezultatin e punës në logjet e Agjentit dhe në seksionin Monitoring në rrjedhën e dhënave airline_tickets:

tail -f /var/log/aws-kinesis-agent/aws-kinesis-agent.log

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Siç duket, gjithçka funksionon dhe Kinesis Agent po dërgon me sukses të dhënat në rrjedhë. Tani do të konfiguroni konsumatorin.

Konfigurimi i Kinesis Data Analytics

Të kalojmë në komponentin qendror të gjithë sistemit - do të krijojmë një aplikacion të ri në Kinesis Data Analytics me emrin kinesis_analytics_airlines_app:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Kinesis Data Analytics lejon kryerjen e analitikës së të dhënave në kohë reale nga Kinesis Streams duke përdorur gjuhën SQL. Ky është një shërbim plotësisht automatikisht i shkallëzuar (në krahasim me Kinesis Streams), i cili:

  1. lejon krijimin e rrjedhave të reja (Output Stream) mbi bazën e kërkesave për të dhënat origjinale;
  2. jep një rrjedhë për gabimet që ndodhin gjatë punës së aplikacioneve (Error Stream);
  3. është në gjendje të përcaktojë automatikisht skemën e të dhënave hyrëse (e cila mund të tejkalohet manualisht në rast nevoje).

Ky është një shërbim relativisht i shtrenjtë — 0.11 USD për orë, prandaj duhet ta përdorni me kujdes dhe ta fshini pas përfundimit të punës.

Të lidhim aplikacionin me burimin e të dhënave:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Zgjidhni shtratin, me të cilin po përpiqemi të lidhim (airline_tickets):

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Më pas, është e nevojshme të bashkangjitni një Rol IAM të ri në mënyrë që aplikacioni të mund të lexojë nga shtrati dhe të shkruajë në shtrat. Për këtë, mjafton të mos ndryshoni asgjë në bllokun e Lejeve të Qasjes:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Tani do të kërkojmë zbDiscoverimin e strukturës së të dhënave në shtrat, për këtë klikoni në butonin «Discover schema». Si rezultat, do të azhurnohet (do të krijohet një rol i ri) IAM dhe do të fillojë zbDiscoverimi i strukturës nga të dhënat që tashmë kanë mb arriv këtë shtrat:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Tani është e nevojshme të kaloni në redaktorin SQL. Duke klikuar në këtë buton, do të dalë një dritare me pyetje për fillimin e aplikacionit — zgjidhni atë që dëshironi të nisni:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Në dritaren e redaktorit SQL, ngjitni 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ë bazat e të dhënave 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 shtrata (STREAM) dhe "pumps" (PUMP) — kërkesa të vazhdueshme të futjes, të cilat futin të dhëna nga një shtrat në aplikacion në një tjetër shtrat.

Në kërkesën SQL të paraqitur më sipër, bëhet kërkimi i biletave të Aeroflotit për çmime më të ulta se pesë mijë rubla. Të gjitha regjistrimet që përmbushin këto kushte do të vendosen në shtratin DESTINATION_SQL_STREAM.

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Në bllokun Destination, zgjidhni shtratin special_stream, ndërsa në listën e rënëse In-application stream name DESTINATION_SQL_STREAM:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Si rezultat i të gjitha manovrave duhet të kemi diçka që i ngjan figurës më poshtë:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Krijimi dhe abonimi në temën SNS

Kaloni në shërbimin Simple Notification Service dhe krijoni një temë të re me emrin Airlines:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Regjistroni abonimin në këtë temë, ku specifikoni numrin e telefonit celular, në të cilin do të marrin SMS njoftime:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Krijimi i një tabele në DynamoDB

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

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Krijimi i funksionit lambda collector

Do të krijojmë një funksion lambda të quajtur Collector, i cili do të ketë detyrë të monitorojë rrjedhën airline_tickets dhe, në rast të gjetjes së shënimeve të reja, t'i vendosë ato në tabelën DynamoDB. Është e qartë se përveç të drejtave të parazgjedhura, kjo lambda duhet të ketë qasje në leximin e rrjedhës së të dhënave Kinesis dhe shkrimin në DynamoDB.

Krijimi i rolit IAM për funksionin lambda collector
Për fillim, le të krijojmë një rol të ri IAM për lambdën e quajtur Lambda-TicketsProcessingRole:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Për një shembull testimi, politikat e paravendosura AmazonKinesisReadOnlyAccess dhe AmazonDynamoDBFullAccess janë të mjaftueshme, siç tregohet në përgjigjen më poshtë:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Ky lambda duhet të kick-start nëpërmjet një trigger nga Kinesis kur të hyjnë shënime të reja në rrjedhën airline_stream, prandaj duhet të shtonim një trigger të ri:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Duhet vetëm të vendosni kodin dhe të ruani lambdën.

"""Analyzing 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:
    """Analizimi i informacionit nga Rryma."""

    def __init__(self, table_name, records):
        """Metoda iniziale."""
        self.table = DYNAMO_DB.Table(table_name)
        self.json_data = TicketsParser.get_json_data(records)

    @staticmethod
    def get_json_data(records):
        """Kthe të dhënat e deserialize-ura nga rryma."""
        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):
        """Parapërpunoni të dhënat json."""
        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('Ka qenë e shtuar ', len(self.json_data), 'artikuj')

def lambda_handler(event, context):
    """Analizoni rrjedhën dhe vendosni në tabelën DynamoDB."""
    print('Mori ngjarje:', 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. Prandaj, kjo lambda duhet të ketë qasje për të lexuar nga Kinesis dhe për të dërguar mesazhe në temën e caktuar SNS, e cila më pas do të dërgohet nga shërbimi SNS të gjithë abonentëve të kësaj teme (email, SMS, etj).

Krijimi i rolit IAM
Së pari krijojmë rolin IAM Lambda-KinesisAlarm për këtë lambda, dhe më pas e caktojmë këtë rol për lambdën që po krijojmë alarm_notifier:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Kjo lambda duhet të funksionojë me një trigger për të kapur regjistrime të reja në rrjedhën special_stream, prandaj është e nevojshme të konfiguroni trigger-in në mënyrë të ngjashme me atë që bëmë për lambdën Collector.

Për lehtësimin e konfigurimit të kësaj lambde, do të prezantojmë një variabël të re ambienti — TOPIC_ARN, në të cilën do të vendosim ARN-në (Amazon Resource Names) të temës Airlines:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Dhe vendosim kodin e lambdës, ai nuk është shumë i komplikuar:

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 diçka interesante!',
                           Subject='Alarm për biletat e fluturimeve')
        print('Mesazhi i alarmit u dorëzua me sukses')
    except Exception as err:
        print('Dështimi i dorëzimit', str(err))

Duket se këtu përfundon konfigurimi manual i sistemit. Ngeli vetëm të testojmë dhe të sigurohemi që kemi konfiguruar gjithçka siç duhet.

Depoy nga kodi Terraform

Përgatitja e nevojshme

Terraform — është një mjet shumë i përshtatshëm open-source për shpërndarjen e infrastrukturës nga kodi. Ai ka sintaksën e tij, e cila është e lehtë për t'u zotëruar dhe një numër të madh shembujsh, si dhe çfarë duhet të shpërndarohet. Në redaktorin Atom ose Visual Studio Code ka shumë plug-in-e të përshtatshme që ndihmojnë në punën me Terraform.

Mund të shkarkoni distribucionin këtu. Një analizë e detajuar e të gjitha mundësive të Terraform kalon përmasat e këtij artikulli, prandaj do të fokusohemi në momentet kryesore.

Si të nisni

Kodi i plotë i projektit ndodhet në repo-n time. Kloni repo-n tek vetja. Para fillimit, sigurohuni që keni instaluar dhe konfiguroni AWS CLI, pasi Terraform do të kërkojë kredencialet në skedarin ~/.aws/credentials.

Një praktikë e mirë para shpërndarjes së gjithë infrastrukturës është të ekzekutoni komandën plan, për të parë se çfarë do të krijojë tani Terraform në re:

terraform.exe plan

Do të propozoni të ingresoni numrin e telefonit për të dërguar njoftimet. Në këtë fazë, futja e tij nuk është e domosdoshme.

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Duke analizuar planin e punës së programit, mund të fillojmë krijimin e burimeve:

terraform.exe apply

Pas 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ë pyetet për realizimin e veprimeve. Kjo do të lejojë të ngrihen të gjithë infrastrukturën, të kryhet e gjithë konfigurimi i nevojshëm EC2, të dislokohen funksionet Lambda, etj.

Pasi të gjitha burimet të kenë qenë me sukses të krijuara përmes kodit Terraform, duhet të hyrni në detajet e aplikacionit Kinesis Analytics (fatkeqësisht, nuk kam gjetur si ta bëj këtë direkt nga kodi).

Nisëm aplikacionin:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Pas kësaj, është e nevojshme të caktohet qartë emri i stream-it brenda aplikacionit, duke zgjedhur nga lista e rënëse:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Tani gjithçka është gati për punë.

Testimi i funksionimit të aplikacionit

Pavarësisht se si e keni dislokuar sistemin, manualisht ose përmes kodit Terraform, do të funksionojë njëjtë.

Hyni me SSH në makinën virtuale EC2, ku është instaluar Kinesis Agent dhe nisni skriptin api_caller.py

sudo ./api_caller.py TOKEN

Tani duhet të presim SMS në numrin tuaj:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
SMS—mesazhi arrin në telefon pothuajse brenda 1 minute:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless
Tani mbetet të shikojmë nëse regjistrimet janë ruajtur në bazën e të dhënave DynamoDB për një analizë më të detajuar për më vonë. Tabela airline_tickets përmban rreth të dhënave si këto:

Integrimi i Aviasales API me Amazon Kinesis dhe thjeshtësia serverless

Përfundim

Gjatë punës së kryer, është ndërtuar një sistem i përpunimit të të dhënave në kohë reale mbi bazën e Amazon Kinesis. Janë shqyrtuar mundësitë e përdorimit të Kinesis Agent në lidhje me Kinesis Data Streams dhe analizat në kohë reale të Kinesis Analytics duke përdorur komanda SQL, si dhe ndërveprimi i Amazon Kinesis me shërbime të tjera AWS.

Sistemi e përshkruar më sipër e kemi dislokuar në dy mënyra: në mënyrë manuale të ngadalshme dhe shpejt nga kodi Terraform.

I gjithë kodi burimor i projektit është i disponueshëm në repositorin tim në GitHub, sugjeroj ta shikoni atë.

Me kënaqësi jam i gatshëm të diskutoj për artikullin, pres komentet tuaja. Shpresoj për kritikë konstruktive.

Ju uroj suksese!

Burimi: habr.com

Blini hostim të besueshëm për faqe interneti me mbrojtje DDoS, serverë VPS VDS 🔥 Blini hostim të besueshëm për faqe interneti me mbrojtje DDoS, serverë VPS VDS - ProHoster