Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

Tere, Habr!

Kas teile meeldib lennata lennukitega? Mulle vĂ€ga meeldib, kuid isoleerimise ajal hakkasin ka analĂŒĂŒsima andmeid ĂŒhe tuntud lennupiletite ressursi — Aviasales — kohta.

TĂ€na kĂ€sitleme Amazon Kinesis'i tööd, ehitame reaalajas analĂŒĂŒsiga voogedastussĂŒsteemi, paigaldame NoSQL andmebaasi Amazon DynamoDB pĂ”hiliseks andmehulgaks ning seadistame SMS-teavituse huvitavate piletite kohta.

KĂ”ik ĂŒksikasjad on allpool! Alustame!

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

Sissejuhatus

Kuna nÀide vajab meil juurdepÀÀsu API Aviasales. JuurdepÀÀs sellele on tasuta ja piiranguteta, peate lihtsalt registreeruma jaama "Arendajatele", et saada oma API token andmete juurde pÀÀsemiseks.

Selle artikli peamine eesmĂ€rk on anda ĂŒlevaade teabe voogedastuse kasutamisest AWS-is, jĂ€ttes kĂ”rvale, et API kaudu saadud andmed ei ole rangelt ajakohased ja need edastatakse vahemĂ€lust, mis moodustatakse Aviasales.ru ja Jetradar.com kasutajate otsingute alusel viimase 48 tunni jooksul.

Kinesis-agendi kaudu saadud lennupiletite andmed, mis on installitud tootmismasinasse, töötlevad automaatselt ja edastavad vajaliku voogu lĂ€bi Kinesis Data Analytics. Töötlemata versioon sellest voogust kirjutatakse otse ladustamisse. DynamoDB-s vĂ€lja töötatud 'toore' andmete ladustamine vĂ”imaldab sĂŒvitsi analĂŒĂŒsida pileteid BI tööriistade, nagu AWS Quick Sight, kaudu.

KĂ€sitleme kahte infrastruktuuri juurutamise varianti:

  • KĂ€sitsi — lĂ€bi AWS Management Console;
  • Terraformi koodist infrastruktuur — laiskade automatiseerijate jaoks;

Arendatava sĂŒsteemi arhitektuur

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Kasutatavad komponendid:

  • Aviasales API — selle API kaudu tagastatavad andmed kasutatakse edasiseks töötamiseks;
  • EC2 tootmisinstants — tavaline virtuaalne masin pilves, kus genereeritakse sisendvoog:
    • Kinesis Agent — see on Java-rakendus, mis installitakse kohalikult masinale ja pakub lihtsat viisi andmete kogumiseks ja saatmiseks Kinesisesse (Kinesis Data Streams vĂ”i Kinesis Firehose). Agent jĂ€lgib pidevalt mÀÀratud kaustades asuvaid failide kogumeid ja saadab uusi andmeid Kinesisesse;
    • API Caller skript — Python-skripti, mis teeb API pĂ€ringuid ja salvestab vastuse kausta, mida jĂ€lgib Kinesis Agent;
  • Kinesis Andmevood — reaalajas andmevoogude teenus laia spektriga skaleerimisvĂ”imetega;
  • Kinesis AnalĂŒĂŒtika — serverivaba teenus, mis lihtsustab reaalajas andmevoogude analĂŒĂŒsi. Amazon Kinesis Data Analytics konfigureerib rakenduste tööks vajalikud ressursid ja skaala automaatselt, et hallata iga sissetuleva andmehulka;
  • AWS Lambda — teenus, mis vĂ”imaldab koodi kĂ€ivitada ilma serverite eraldamise ja seadistamiseta. KĂ”ik arvutusvĂ”imsused skaleeruvad automaatselt iga kutsumise puhul;
  • Amazon DynamoDB — paaride "vĂ”ti-vÀÀrtus" ja dokumentide andmebaas, mis tagab vĂ€hem kui 10 ms latentsuse igas mÔÔtkavas. DynamoDB kasutamisel ei pea jaotama servereid, installima paranduspakette ega haldama neid. DynamoDB skaleerib automaatselt tabeleid, kohandades saadaolevate ressursside hulka ja hoides kĂ”rge jĂ”udluse. SĂŒsteemi haldamise toimingud pole vajalikud;
  • Amazon SNS — tĂ€ielikult hallatav sĂ”numite saatmise teenus ‘vĂ€ljund – tellija’ (Pub/Sub) mudeli kaudu, mille abil saab isoleerida mikroteenuseid, hajutatud sĂŒsteeme ja serverita rakendusi. SNS-i saab kasutada teabe saatmiseks lĂ”ppkasutajatele mobiilsete push-teadete, SMS-sĂ”numite ja e-kirjade kaudu.

Algne ettevalmistus

Andmevoo emuleerimiseks otsustasin kasutada lennupiletite teavet, mida tagastab API Aviasales. Siin on dokumentatsioonis suhteliselt ulatuslik nimekiri erinevatest meetoditest; vĂ”tame neist ĂŒhe — ‘Kuu hinna kalender’, mis tagastab hinnad iga kuu pĂ€eva kohta, rĂŒhmitatuna ĂŒmberistumiste arvu jĂ€rgi. Kui otsinguprogrammi kuud ei edastata, tagastatakse teave kuu kohta, mis jĂ€rgneb praegusele.

Nii et registreerime end, saame oma tokeni.

Allpool on nÀidis pÀring:

http://api.travelpayouts.com/v2/prices/month-matrix?currency=rub&origin=LED&destination=HKT&show_to_affiliates=true&token=TOKEN_API

Ülaltoodud API-andmete saamise meetod koos tokeni edastamisega pĂ€ringus töötab, kuid mulle meeldib rohkem edastada juurdepÀÀs tokeni kaudu pĂ€ises, seega kasutame skriptis api_caller.py just seda meetodit.

NĂ€idis vastus:

{{
   "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
   }]
}

Ülaltoodud API vastuse nĂ€ites on pilet Peterburist Phuketi
 Oh, unistamisest pole mĂ”tet

Kuna ma olen Kazanist ja Phuket on praegu «meie jaoks vaid unistus», otsime pileteid Peterburist Kazani.

Eeldatakse, et teil on juba AWS konto. Tahan kohe rĂ”hutada, et Kinesis ja teavitamine SMS-i kaudu ei kuulu aastase Tasuta kasutustaseme (Free Tier). Kuid isegi sellega arvestades, mĂ”ne dollari investeerides on tĂ€iesti vĂ”imalik ĂŒles ehitada pakutud sĂŒsteem ja sellega mĂ€ngida. Ja loomulikult ei tohi unustada kĂ”ik ressursid kustutada, kui need enam vajalikud ei ole.

Õnneks on DynamoDb ja lambda-funktsioonid meile tinglikult tasuta, kui jÀÀte kuiste tasuta limiitide sisse. NĂ€iteks DynamoDB puhul: 25 GB salvestusruumi, 25 WCU/RUC ja 100 miljonit pĂ€ringut. Ja miljon lambda funktsioonide kutsumist kuus.

SĂŒsteemi kĂ€sitsi juurutamine

Kinesis Data Streams seadistamine

Liigume Kinesis Data Streams teenusesse ja loome kaks uut voogu, igas ĂŒhes shard.

Mis on shard?
Shard on Amazon Kinesis voogude peamine andmeedastuse ĂŒksus. Üks shard vĂ”imaldab sissetulevate andmete edastust kiirusel 1 MB/s ja vĂ€ljuvate andmete edastust kiirusel 2 MB/s. Üks shard toetab kuni 1000 PUT kirjet sekundis. Voogude loomisel tuleb mÀÀrata soovitud shardide arv. NĂ€iteks vĂ”ib luua kaheastmelise voo. See voog tagab sissetulevate andmete edastuse kiirusel 2 MB/s ja vĂ€ljuvate andmete edastuse kiirusel 4 MB/s, toetades kuni 2000 PUT kirjet sekundis.

Mida rohkem shard'e teie voos on, seda suurem on selle lÀbilaskevÔime. PÔhimÔtteliselt skaleeritakse vooge shardide lisamise kaudu. Kuid mida rohkem shard'e teil on, seda kallim see on. Iga shard maksab 1,5 senti tunnis ja lisaks 1,4 senti iga miljoni PUT toimingu kohta (PUT payload units).

Loome uue voo nimega airline_tickets, sellele piisab tÀiesti 1 shardist:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
NĂŒĂŒd loome veel ĂŒhe voo nimega special_stream:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

Produtsendi seadistamine

Andmete töötluse produtsendina on tavaline EC2 instants ĂŒlesande lahendamiseks piisav. See ei pea olema vĂ”imas kallis virtuaalmasin, spottide t2.micro sobib tĂ€iesti.

Oluline mĂ€rkus: nĂ€idisena tuleks kasutada pilti — Amazon Linux AMI 2018.03.0, sest sellega on Kinesis Agendi kiireks kĂ€ivitamiseks vĂ€hem seadistusi.

Liigume EC2 teenusesse, loome uue virtuaalmasina, valime vajaliku AMI tĂŒĂŒbi t2.micro, mis kuulub Free Tieri:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Kuna uuel loodud virtuaalmasinal peab olema Ă”igused Kinesis teenusega suhtlemiseks, tuleb need Ă”igused anda. Parim viis selleks on mÀÀrata IAM Roll. Seega tuleb ekraanil Samuti 3: Konfigureeri instantsi ĂŒksikasjad valida Loo uus IAM Roll:

IAM rolli loomine EC2 jaoks
Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Avanenud aknas valime, et loome uue rolli EC2 jaoks ja lÀheme jaotisse Litsentsid:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
ÕppenĂ€idisena ei pea me kogu ressursside Ă”iguste detailset seadistamist arvesse vĂ”tma, seega valime Amazoni eelkonfigureeritud poliitikad: AmazonKinesisFullAccess ja CloudWatchFullAccess.

Anname sellele rollile mÔne mÔistliku nime, nÀiteks: EC2-KinesisStreams-FullAccess. Tulemusena peaks olema sama, mis on nÀidatud alloleval pildil:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
PÀrast selle uue rolli loomist Àrge unustage seda siduda loodava virtuaalmasina instantsiga:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Sellel ekraanil ei muudeta midagi ja liigutakse jÀrgmistele akendele.

Ketta parameetrid vÔivad jÀÀda vaikeseadeteks, samuti ka sildid (kuigi hea tava on kasutada silte, andes nÀiteks instantsile nime ja mÀÀrates keskkonna).

NĂŒĂŒd oleme vahekaartidel Samm 6: Konfigureeri turgroup, kus on vajalik luua uus vĂ”i nĂ€idata olemasolevat turgroup'i, mis lubab SSH kaudu (port 22) instantsiga ĂŒhendust vĂ”tta. Valige seal Allikas —> Minu IP ja saate instantsi kĂ€ivitada.

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Niipea, kui see lĂ€heb staatuseks running, saate proovida sellele SSH kaudu ĂŒhendust luua.

Kuna Kinesis Agenti kasutamine on vĂ”imaldanud pĂ€rast masinaga edukat ĂŒhendust, peate terminalis sisestama jĂ€rgmised kĂ€sud:

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

Loomine kaust vastuste salvestamiseks API-le:

sudo mkdir /var/log/airline_tickets

Enne agendi kÀivitamist on vajalik selle konfiguratsiooni seadistamine:

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

Faili agent.json sisu peaks olema jÀrgmine:

{
  "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"]
         }
      ]
    }
  ]
}

Nagu nÀha konfiguratsioonifailist, jÀlgib agent kataloogis \/var\/log\/airline_tickets\/ faile, millel on .log laiendus, parsimine ja edastamine voolu airline_tickets.

TaaskÀivitage teenus ja veenduge, et see kÀivitub ja töötab:

sudo service aws-kinesis-agent restart

NĂŒĂŒd laeme alla Python-skripti, mis kĂŒsib andmeid API-lt:

REPO_PATH=https:\/\/raw.githubusercontent.com\/igorgorbenko\/aviasales_kinesis\/master\/producer

wget $REPO_PATH\/api_caller.py -P \/home\/ec2-user\/\nwget $REPO_PATH\/requirements.txt -P \/home\/ec2-user\/\nsudo chmod a+x \/home\/ec2-user\/api_caller.py\nsudo \/usr\/local\/bin\/pip3 install -r \/home\/ec2-user\/requirements.txt

Skript api_caller.py kĂŒsib andmeid Aviasalesilt ja salvestab saadud vastuse katalooge, mida skanneerib Kinesis agent. Selle skripti teostus on piisavalt standardne, seal on klass TicketsApi, mis vĂ”imaldab asĂŒnkroonselt API-d kutsuda. Sellesse klassi edastame pĂ€ise koos mĂ€rgiga ja pĂ€ringu parameetrid:

class TicketsApi:
    """Api caller class."""

    def __init__(self, headers):
        """Init method."""
        self.base_url = BASE_URL
        self.headers = headers

    async def get_data(self, data):
        """Get the data from API query."""
        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('Response status %s: %s',
                            self.base_url, response.status)
                response_json = await response.json()
            except HTTPError as http_err:
                LOGGER.error('Oops! HTTP error occurred: %s', str(http_err))
            except Exception as err:
                LOGGER.error('Oops! An error occurred: %s', str(err))
            return response_json


def prepare_request(api_token):
    """Return the headers and query for the API request."""
    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():
    """Get run the code."""
    if len(sys.argv) != 2:
        print('Usage: 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 has returned %s items', len(response['data']))
        try:
            count_rows = log_maker(response)
            LOGGER.info('%s rows have been saved into %s',
                        count_rows,
                        TARGET_FILE)
        except Exception as e:
            LOGGER.error('Oops! Request result was not saved to file. %s',
                         str(e))
    else:
        LOGGER.error('Oops! API request was unsuccessful %s!', response)

Seadistuste ja agendi töökindluse testimiseks teeme testkÀivituse skripti api_caller.py:

sudo ./api_caller.py TOKEN

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Ja vaatame agendi logides ning Monitoring vahekaardil tulemusi andmevoos airline_tickets:

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

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Nagu nĂ€ha, kĂ”ik töötab ja Kinesis Agent edastab andmeid voosse edukalt. NĂŒĂŒd seadistame consumer'i.

Kinesis Data Analytics seadistamine

Liigume keskse komponendi juurde — loome uue rakenduse Kinesis Data Analytics'isse nimega kinesis_analytics_airlines_app:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Kinesis Data Analytics vĂ”imaldab teha reaalaja andmeanalĂŒĂŒsi Kinesis Streams'is SQL keele abil. See on tĂ€ielikult automaatne teenus (erinevalt Kinesis Streams'ist), mis:

  1. lubab luua uusi vooge (Output Stream) lÀhtudes algandmetest pÀringutest;
  2. pakub voogu vigadest, mis tekkisid rakenduste töö kÀigus (Error Stream);
  3. oskab automaatselt mÀÀrata sisendi andmete skeemi (seda saab vajadusel kĂ€sitsi ĂŒle mÀÀrata).

See ei ole odav teenus — 0.11 USD tunni kohta, seega tuleks seda kasutada ettevaatlikult ja eemaldada pĂ€rast töö lĂ”ppemist.

Ühendame rakenduse andmeallikaga:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Valime voolu, millega kavatseme ĂŒhenduda (airline_tickets):

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
JĂ€rgmisena on vaja lisada uus IAM-rolle, et rakendus saaks voost lugeda ja voogu kirjutada. Selleks ei pea Access permissions plokis midagi muutma:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
NĂŒĂŒd kĂŒsime andmeskeemi avastamist voos, selleks vajutame nuppu „Discover schema“. Tulemusena uuendatakse (loodakse uus) IAM-rollen ja kĂ€ivitatakse andmeskeemi avastamine saadud andmetest voos:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
NĂŒĂŒd tuleb minna SQL-redaktorisse. Selle nupu vajutamisel avaneb aken rakenduse kĂ€ivitamise kĂŒsimusega — valime, mida soovime kĂ€ivitada:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
SQL-redaktori aknasse lisame nii lihtsa pÀringu ja vajutame 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';

Relatsioonilistes andmebaasides töötate tabelitega, kasutades INSERT kĂ€ske kirje lisamiseks ja SELECT kĂ€sku andmete pĂ€rimiseks. Amazon Kinesis Data Analytics'is töötate voogudega (STREAM) ja «pumpadega» (PUMP) — pidevad sisestamispĂ€ringud, mis edastavad andmeid ĂŒhest voost rakenduses teise voogu.

Ülaltoodud SQL-kĂ€sus otsitakse Aerofloti pileteid, mille hind on alla viie tuhande rubla. KĂ”ik kirjed, mis vastavad nendele tingimustele, paigutatakse DESTINATION_SQL_STREAM voogu.

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Sihtkohas valime voolu special_stream ja rippmenĂŒĂŒst rakenduse voonime DESTINATION_SQL_STREAM:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
KÔikide toimingute tulemus peaks olema midagi sarnast allolevale pildile:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

SNS teema loomine ja sellele registreerimine

Liigume Simple Notification Service'i teenusesse ja loome seal uue teema nimega Airlines:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Registreerime sellel teemal, kus anname mobiiltelefoni numbri, kuhu saadetakse SMS-teated:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

Tabeli loomine DynamoDB-s

Töötlemata andmete salvestamiseks voost airline_tickets loome DynamoDB-s sama nimega tabeli. Peamiseks vÔtmeks kasutame record_id:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

Lambda-funktsiooni loomine collector

Loome lambda-funktsiooni nimega Collector, mille ĂŒlesanne on kĂŒsida airline_tickets voogu ja, kui seal leitakse uusi kirjeid, sisestada need kirjed DynamoDB tabelisse. Ilmselgelt peab see lambda, peale vaikimisi Ă”iguste, olema loodud Kinesis andmevoo lugemiseks ja DynamoDB-sse kirjutamiseks.

IAM rolli loomine lambda-funktsioonile collector
Alustame uue IAM rolli loomisega lambdast nimega Lambda-TicketsProcessingRole:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Testimiseks sobivad tÀiesti eelhÀÀlestatud poliitikad AmazonKinesisReadOnlyAccess ja AmazonDynamoDBFullAccess, nagu on nÀidatud alloleval pildil:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

See lambda peab kÀivituma Kinesis'i kÀivitamisel, kui uusi kirjeid lisatakse airline_stream voogu, seega peame lisama uue kÀivitaja:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
JÀÀnud on koodi lisamine ja lambda salvestamine.

"""Voogude analĂŒĂŒsimine ja sisestamine DynamoDB tabelisse."""
import base64
import json
import boto3
from decimal import Decimal

DYNAMO_DB = boto3.resource('dynamodb')
TABLE_NAME = 'airline_tickets'

class TicketsParser:
    """Infoteabe analĂŒĂŒsimine voogudest."""

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

    @staticmethod
    def get_json_data(records):
        """Tagastab deserialiseeritud andmed voogudest."""
        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):
        """Eeltöödelda json andmed."""
        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):
        """Partiis sisestamine tabelisse."""
        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('On lisatud ', len(self.json_data), 'objekti')

def lambda_handler(event, context):
    """AnalĂŒĂŒsi voog ja sisestamine DynamoDB tabelisse."""
    print('Sain sĂŒndmuse:', event)
    parser = TicketsParser(TABLE_NAME, event['Records'])
    parser.run()

Lambda-funktsiooni notifier loomine

Teine Lambda-funktsioon, mis jÀlgib teist voogu (special_stream) ja saadab teate SNS-ile, luuakse sarnaselt. Seega peab see Lambda olema ajaloo lugemise ligipÀÀs Kinesise voogudesse ning saadama teateid antud SNS-teemale, mis edastatakse edasi kÔigile selle teema tellijatele (e-post, SMS jne).

IAM rolli loomine
Esiteks loome IAM rolli Lambda-KinesisAlarm selle Lambda jaoks ja seejÀrel mÀÀrame selle rolli loodavale Lambda'le alarm_notifier:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

See Lambda peab töötama nuhtluse pÔhjal, kui uusi kirjeid lisatakse voogu special_stream, seega tuleb nuhtlust seadistada sarnaselt sellele, kuidas me tegime Lambda Collectoriga.

Selle Lambda seadistamise hĂ”lbustamiseks loome uue keskkonnamuutuja — TOPIC_ARN, kuhu paneme Airlines teema ANR-d (Amazon Resource Names):

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Ja sisestame Lambda koodi, see pole ĂŒldse keeruline:

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='Tere! Olen leidnud huvitavat kraami!',
                           Subject='Lennupiletite alarm')
        print('Alarmiteade on edukalt edastatud')
    except Exception as err:
        print('Eddastamine ebaÔnnestus', str(err))

Tundub, et sĂŒsteemi kĂ€sitsi seadistamine on nĂŒĂŒd lĂ”petatud. JÀÀnud on vaid testida ja veenduda, et oleme kĂ”ik Ă”igesti seadistanud.

Terraformi koodist silti tÔstmine

Vajalik ettevalmistus

Terraform — vĂ€ga mugav avatud lĂ€htekoodiga tööriist infrastruktuuri kĂ€ivitamiseks koodist. Sellel on oma sĂŒntaks, mida on lihtne omandada, ja palju nĂ€iteid, kuidas ja mida kĂ€ivitada. Redaktoris Atom vĂ”i Visual Studio Code on palju mugavaid pluginaid, mis lihtsustavad töö tegemist Terraformiga.

Tarkvara saab alla laadida siit. Terraformi kĂ”igi vĂ”imaluste pĂ”hjalik analĂŒĂŒs ĂŒletab selle artikli piire, seetĂ”ttu piirdume peamiste punktidega.

Kuidas alustada

Kogu projekti kood asub minu repositooriumis. Kloonime enda repositooriumi. Enne kÀivitamist tuleb veenduda, et teil on installitud ja seadistatud AWS CLI, kuna Terraform kontrollib autentimist failis ~/.aws/credentials.

Hea praktikana, enne kogu infrastruktuuri silti tÔstmist, kÀivitage kÀsk plan, et nÀha, mida Terraform meil praegu pilves loomisel on:

terraform.exe plan

KĂŒsimine telefoninumbri sisestamiseks, et saata sellele teatisi. Selles etapis pole selle sisestamine kohustuslik.

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
PĂ€rast programmi tööplaani analĂŒĂŒsimist saame alustada ressursside loomist:

terraform.exe apply

PĂ€rast selle kĂ€su saatmist ilmub taas telefoninumbri sisestamise kĂŒsimus, kirjutame "yes", kui esitatakse kĂŒsimus tegevuste tegeliku tĂ€itmise kohta. See vĂ”imaldab tĂ”sta kogu infrastruktuuri, teha vajalikud EC2 seaded, kĂ€ivitada lambda-funktsioone jne.

Kuna kÔik ressursid on edukalt loodud Terraformi koodi kaudu, peate minema Kinesis Analytics rakenduse detailidesse (kahjuks ei leidnud ma viisi selle tegemiseks otse koodi kaudu).

KĂ€ivitame rakenduse:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
PĂ€rast seda tuleb selgelt mÀÀrata rakenduse voogude nimi, valides rippmenĂŒĂŒst:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
NĂŒĂŒd on kĂ”ik tööks valmis.

Rakenduse testimine

SĂ”ltumata sellest, kuidas te sĂŒsteemi juurutate, kĂ€sitsi vĂ”i Terraformi koodi kaudu, töötab see ĂŒhtemoodi.

Sisestame SSH kaudu EC2 virtuaalmasinasse, kuhu on installitud Kinesis Agent, ja kÀivitame skripti api_caller.py

sudo ./api_caller.py TOKEN

JÀÀb oodata SMS-i teie numbrile:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
SMS — sĂ”num tuleb telefonile praktiliselt 1 minuti jooksul:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus
JÀÀn ootama, et nĂ€ha, kas DynamoDB andmebaasis on salvestatud kirjed edasiseks, pĂ”hjalikumaks analĂŒĂŒsiks. Tabel airline_tickets sisaldab enam-vĂ€hem selliseid andmeid:

Aviasales API integreerimine Amazon Kinesisiga ja serverless lihtsus

KokkuvÔte

Töötamise kĂ€igus on loodud veebipĂ”hine andmete töötlemise sĂŒsteem Amazon Kinesis alusel. KĂ€sitleti Kinesis Agendi kasutamise vĂ”imalusi koos Kinesis Data Streams'i ja reaalajas analĂŒĂŒtika Kinesis Analytics'iga SQL kĂ€skude abil, samuti Amazon Kinesis'i integreerimist teiste AWS teenustega.

Üksikasjalikult kirjeldatud sĂŒsteemi oleme juurutatud kahte moodi: suhteliselt pika kĂ€sitsi ja kiire koodi Terraform abil.

Kogu projekti lÀhtekood on saadaval minu GitHubi hoidlas, soovitan sellega tutvuda.

Olen valmis artiklit arutama, ootan teie kommentaare. Lootes konstruktiivset kriitikat.

Soovin edu!

Allikas: habr.com

Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster