Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

Tere, Habr!

Kas teile meeldib lennata? Mulle meeldib, kuid isoleerimise ajal olen hakanud ka analĂŒĂŒsima andmeid tuntud lennupiletite portaalist — Aviasales.

TĂ€na vaatame Amazon Kinesis’i tööpĂ”himĂ”tteid, loome voogedastuse sĂŒsteemi reaalajas analĂŒĂŒsiga, seadistame NoSQL andmebaasi Amazon DynamoDB peamiseks andmehoidla ning seadistame SMS-teavituse huvitavate piletite kohta.

KÔik detailid allpool! LÀhme!

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

Sissejuhatus

NĂ€iteks vajame ligipÀÀsu Aviasales API-le. LigipÀÀs sellele on tasuta ja piiranguteta, tuleb lihtsalt registreeruda ja saada oma API token andmete saamiseks jaotises „Arendajatele“.

Selle artikli peamine eesmĂ€rk on anda ĂŒldine ĂŒlevaade andmete voogedastuse kasutamisest AWS-is. JĂ€tame kĂ”rvale, et API kaudu tagasi saadud andmed ei pruugi olla praegused ja need edastatakse vahemĂ€lust, mis moodustatakse Aviasales.ru ja Jetradar.com kasutajate otsingute pĂ”hjal viimase 48 tunni jooksul.

API kaudu saadud lennupiletite andmed (Kinesis-agent, mis on paigaldatud tootmismasinale) analĂŒĂŒsib automaatselt ja edastab vajalikku voogu Kinesis Data Analytics kaudu. Tootmisvoog kirjutatakse otse andmehoidlasse. DynamoDB-s rakendatud „tooĆŸit“ andmehoidla vĂ”imaldab lĂ€bi viia sĂŒgavamat analĂŒĂŒsi piletite kohta BI tööriistade kaudu, nĂ€iteks AWS Quick Sight.

KĂ€sitleme kahte infrastruktuuri juurutamise varianti:

  • KĂ€sitsi — lĂ€bi AWS Management Console;
  • Koodi Terraformi alusel — lenijatele automatiseerijatele;

Arendatava sĂŒsteemi arhitektuur

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Kasutatavad komponendid:

  • Aviasales API — selle API kaudu tagastatud andmeid kasutatakse kogu edasise töötlemise jaoks;
  • EC2 tootmisinstants — tavaline virtuaalne masin pilves, kus genereeritakse sisendvoog;
    • Kinesis Agent — Java-rakendus, mis paigaldatakse kohalikult masinale ja mis pakub lihtsat viisi andmete kogumiseks ja edastamiseks Kinesis'esse (Kinesis Data Streams vĂ”i Kinesis Firehose). Agent jĂ€lgib pidevalt mÀÀratud kaustade failide kogumeid ja edastab uusi andmeid Kinesis'esse;
    • API Caller skript — Python-skript, mis teeb API pĂ€ringuid ja salvestab vastuse kausta, mida jĂ€lgib Kinesis Agent;
  • Kinesis Data Streams — reaalajas andmete voogude edastamise teenus, millel on ulatuslikud skaleerimisvĂ”imalused;
  • Kinesis Analytics — serverita teenus, mis lihtsustab reaalajas voogandmete analĂŒĂŒsi. Amazon Kinesis Data Analytics konfigureerib rakenduste tööks vajalikud ressursid ja skaleerib automaatselt, et hallata mis tahes andmevooge;
  • AWS Lambda — teenus, mis vĂ”imaldab kĂ€ivitada koodi ilma serverite reserveerimise ja seadistamiseta. KĂ”ik arvutusvĂ”imsused skaleeruvad automaatselt iga vĂ€ljakutse jaoks;
  • Amazon DynamoDB — paaride "vĂ”ti-vÀÀrtus" ja dokumentide andmebaas, mis tagab alla 10 millisekundi latentsuse igas mÔÔtkavas. DynamoDB kasutamine ei nĂ”ua serverite jaotamist, nende vĂ€rskendamist ega haldamist. DynamoDB skaleerib tabeleid automaatselt, kohandades saadavalolevate ressursside mahtu ja sĂ€ilitades kĂ”rge jĂ”udluse. SĂŒsteemi haldamise tegevusi ei ole vaja;
  • Amazon SNS — tĂ€ielikult hallatav sĂ”numite saatmise teenus, mis pĂ”hineb "vĂ€ljaandja - tellija" (Pub/Sub) mudelil, mille kaudu saab isoleerida mikroteenuseid, jaotatud sĂŒsteeme ja serverita rakendusi. SNS-i saab kasutada teabe saatmiseks lĂ”ppkasutajatele mobiilsete push-teate, SMS-i ja e-kirjade kaudu.

Algne ettevalmistus

Andmevoo simuleerimiseks otsustasin kasutada lennupiletite teavet, mille tagastab Aviasalesi API. dokumentatsioon suur hulk erinevaid meetodeid, vĂ”tame ĂŒhe neist — "Hinnakalender kuuks", mis tagastab igapĂ€evased hinnad, grupeerituna ĂŒmberistumiste arvu. Kui otsingus kuud ei edastata, tagastatakse jĂ€rgmise kuu teave.

Nii et registreerime end, saame oma tokeni.

Allpool on nÀide pÀringust:

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

Ülaltoodud meetod API-st andmete saamiseks, kasutades pĂ€ringus tokenit, töötab, kuid ma eelistan edastada juurdepÀÀsutokeni lĂ€bi pĂ€ise, seega kasutame skripti api_caller.py jaoks just seda meetodit.

NĂ€ide vastusest:

{{
   "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 Phuketti... Ah, milleks unistada...
Kuna ma olen Kazanist ja Phukett on meil praegu "ainult unistus", otsime pileteid Peterburist Kazani.

Eeldatakse, et teil on juba AWSi konto. Soovin kohe eraldi tĂ€helepanu juhtida, et Kinesis ja SMS-ide saatmine ei kuulu aastase Free Tier (tasuta kasutamine). Kuid isegi sel juhul, pidades meeles paar dollarit, on tĂ€iesti vĂ”imalik luua soovitatud sĂŒsteem ja sellega mĂ€ngida. Ja muidugi, Ă€rge unustage eemaldada kĂ”iki ressursse pĂ€rast nende vajaduse lĂ”ppemist.

Õnneks on DynamoDb ja lambda-funktsioonid meie jaoks tinglikult tasuta, kui jÀÀme kuueelarve tasuta piiridesse. NĂ€iteks DynamoDB puhul: 25 GB salvestusruumi, 25 WCU/RCU ja 100 miljonit pĂ€ringut. Ja miljon lambda funktsiooni kutsumist kuus.

SĂŒsteemi kĂ€sitsi juurutamine

Kinesis Data Streamsi seadistamine

Liigume Kinesis Data Streamsi teenusesse ja loome kaks uut voogu, igaĂŒhel ĂŒks shard.

Mis on shard?
Shard on Amazon Kinesis voogude andmete edastamise pĂ”hielement. Üks shard tagab sisendandmete edastamise kiirusel 1 MB/s ja vĂ€ljundandmete edastamise kiirusel 2 MB/s. Üks shard toetab kuni 1000 PUT faili sekundis. Andmevoo loomisel tuleb mÀÀrata vajalik shardide arv. NĂ€iteks saab luua andmevoo kahe shardiga. See andmevoog tagab sisendandmete edastamise kiirusel 2 MB/s ja vĂ€ljundandmete edastamise kiirusel 4 MB/s, toetades kuni 2000 PUT faili sekundis.

Mida rohkem shard'e teie voos on, seda suurem on selle lĂ€bilaskevĂ”ime. Üldiselt skaleeritakse vooge shardide lisamise kaudu. Aga mida rohkem shard'e teil on, seda kallim see on. Iga shard maksab 1,5 senti tunnis ja lisaks 1,4 senti iga miljoni vookuulutuse (PUT payload units) kohta.

Loome uue voolu nimega airline_tickets, talle piisab tĂ€iesti ĂŒhest shardist:

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

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

Tootja seadistamine

Andmete tootjana piisab ĂŒlesande lahendamiseks tavalise EC2 instantsi kasutamisest. See ei pea olema vĂ”imas kallis virtuaalmasin, piisab tĂ€iesti spottest t2.micro-st.

Oluline mĂ€rkus: nĂ€idisena tuleks kasutada image'i — Amazon Linux AMI 2018.03.0, millel on vĂ€hem seadistusi Kinesis Agendi kiireks kĂ€ivitamiseks.

Liigume EC2 teenusse, loome uue virtuaalmasina, valime sobiva AMI tĂŒĂŒbi t2.micro, mis kuulub Free Tier'i:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Kuna uus virtuaalmasin peab Kinesis teenusega suhtlema, tuleb talle anda selleks Ôigused. Parim viis seda teha on mÀÀrata IAM Role. SeetÔttu tuleks ekraanil Step 3: Configure Instance Details valida Loo uus IAM Roll:

IAM rolli loomine EC2 jaoks
Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Avanenud aknas valime, et loome uue rolli EC2 jaoks ja liigume jaotisse Permissions:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Õppemudeli puhul ei ole vajalik igasse detaili sĂŒveneda, seetĂ”ttu valime Amazoni eelmÀÀratud poliisid: AmazonKinesisFullAccess ja CloudWatchFullAccess.

Anname sellele rollile mÔne mÔistliku nime, nÀiteks: EC2-KinesisStreams-FullAccess. Tulemuseks peaks olema sama, mis on kujutisel allpool:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
PÀrast uue rolli loomist Àrge unustage seda kinnitada loodavale virtuaalmasina instantsile:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Selles ekraanil ei muuda me enam midagi ja liigume jÀrgmistele akendele.

Karmide ketaste parameetrid vÔib jÀtta vaikevÀÀrtustele, sildid samuti (kuigi hea tava on silte kasutada, vÀhemalt anda instantsile nimi ja mÀrkida keskkond).

NĂŒĂŒd on meil avatud Step 6: Configure Security Group vahekaart, kus tuleb luua uus vĂ”i nĂ€idata olemasolevat Security group, mis vĂ”imaldab SSH (port 22) ĂŒhendust instantsiga. Valige seal Source —> My IP ja saate instantsi kĂ€ivitada.

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Kui see lĂ€heb staatuseks running, saab proovida selle ssh kaudu ĂŒhendust saada.

Kinesis Agendi kasutamiseks, pĂ€rast masina edukat ĂŒhendamist, tuleb terminalis sisestada 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

Loome kausta API vastuste salvestamiseks:

sudo mkdir /var/log/airline_tickets

Enne agendi kÀivitamist tuleb selle konfiguratsioon seadistada:

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Àhtub konfiguratsioonifailist, jÀlgib agent kataloogis /var/log/airline_tickets/ .log laiendiga faile, parsib need ja edastab airline_tickets streaming.

TaaskÀivitame teenuse ja veendume, et see on kÀivitatud ja töötab:

sudo service aws-kinesis-agent restart

NĂŒĂŒd laeme alla Python-skripti, mis pĂ€rib andmeid API-lt:

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

Skript api_caller.py pĂ€rib andmeid Aviasales'ilt ja salvestab saadud vastuse katalooge, mida skannib Kinesis agent. Selle skripti rakendamine on ĂŒsna standardne, on olemas klass TicketsApi, mis vĂ”imaldab asĂŒnkroonselt API-kutsed teha. Selle klassi mÀÀrame pĂ€ise koos tokeniga ja pĂ€ringu parameetrid:

class TicketsApi:
    """API kutsuja klass."""

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

    async def get_data(self, data):
        """Hangi andmed API pÀringust."""
        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('Vastuse olek %s: %s',
                            self.base_url, response.status)
                response_json = await response.json()
            except HTTPError as http_err:
                LOGGER.error('Ops! HTTP viga sĂŒndis: %s', str(http_err))
            except Exception as err:
                LOGGER.error('Ops! SĂŒndis viga: %s', str(err))
            return response_json


def prepare_request(api_token):
    """Tagasta pÀised ja pÀring API pÀringu jaoks."""
    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():
    """KĂ€ivita kood."""
    if len(sys.argv) != 2:
        print('Kasutamine: api_caller.py <teie_api_token>')
        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 on tagastanud %s elementi', len(response['data']))
        try:
            count_rows = log_maker(response)
            LOGGER.info('%s rida on salvestatud faili %s',
                        count_rows,
                        TARGET_FILE)
        except Exception as e:
            LOGGER.error('Ops! PĂ€ringu tulemust ei salvestatud faili. %s',
                         str(e))
    else:
        LOGGER.error('Ops! API pÀring ebaÔnnestus %s!', response)

PÀrast sÀtete Ôigsuse ja agendi töö kontrollimist teeme testkÀivituse skripti api_caller.py jaoks:

sudo .\/api_caller.py TOKEN

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Ja vaatame tulemusi Avari Logfailides ja Monitoring vahekaardil andmevoos airline_tickets:

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

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

Kinesis Data Analytics konfigureerimine

Liigume kogu sĂŒsteemi keskse komponendi juurde — loon uue rakenduse Kinesis Data Analytics nimelt kinesis_analytics_airlines_app:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Kinesis Data Analytics vĂ”imaldab analĂŒĂŒsida andmeid reaalajas Kinesis Streams'i kaudu SQL keeles. See on tĂ€ielikult automaatne teenus (erinevalt Kinesis Streams'ist), mis:

  1. lubab luua uusi vooge (Output Stream) lÀhtudes pÀringutest algandmetele;
  2. pakub tÔrkevoogu, mis ilmnes rakenduste töö kÀigus (Error Stream);
  3. oskab automaatselt mÀÀrata algandmete skeemi (seda saab vajadusel kĂ€sitsi ĂŒmber mÀÀrata).

See teenus ei ole odav — 0,11 USD tunni eest, seega tuleks seda kasutada ettevaatlikult ja kustutada pĂ€rast kasutamise lĂ”petamist.

Ühendame rakenduse andmeallikaga:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Valime voolu, millega plaanime ĂŒhendust luua (airline_tickets):

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
JĂ€rgmiseks tuleb rakendusele lisada uus IAM roll, et see saaks voost lugeda ja voogu kirjutada. Selleks ei pea Access permissions blokis midagi muutma:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
NĂŒĂŒd kĂŒsime andmeskeemi avastamist voos, selleks kliki nuppu «Discover schema». Selle tulemusena uuendatakse (loodakse uus) IAM roll ja kĂ€ivitatakse andmeskeemi avastamine, kasutades andmeid, mis on voogu juba saabunud:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
NĂŒĂŒd on vajalik minna SQL redaktorisse. NĂŒĂŒdse nupu vajutamisel ilmub aken kĂŒsimusega rakenduse kĂ€ivitamise kohta — valige, mida soovite kĂ€ivitada:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
SQL redaktori aknasse lisame sellise 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 te tabelitega, kasutades INSERT-lauseid kirje lisamiseks ja SELECT-lauseid andmete pĂ€rimiseks. Amazon Kinesis Data Analytics'is töötate te voogude (STREAM) ja 'pump' (PUMP) - pidevate sisestuspĂ€ringutega, mis sisestavad andmed ĂŒhest voost teise rakenduses.

Ülaltoodud SQL pĂ€ringus otsitakse Aerofloti pileteid, mille hind on alla viie tuhande rubla. KĂ”ik kirjed, mis vastavad nendele tingimustele, lisatakse voogu DESTINATION_SQL_STREAM.

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Sihtosakonna all valime voolu special_stream ja rippmenĂŒĂŒs In-application stream name DESTINATION_SQL_STREAM:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
KÔigi manipuleerimiste tulemuseks peaks olema midagi sarnast alloleva pildiga:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

SNS teema loomine ja tellimine

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

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Oleme selle teema tellimise vormistamisel, seal mÀÀrame mobiiltelefoni numbri, kuhu saadetakse SMS-teavitused:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

DynamoDB-s tabeli loomine

KÀivitatud andmete salvestamiseks voost airline_tickets loome DynamoDB-s tabeli sama nimega. Peamiseks vÔtmevÀiks on record_id:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

Loomine lambdafunktsioon collector

Loome lambda-funktsioon nimega Collector, mille ĂŒlesandeks on kĂŒsida airline_tickets voogust ja, kui seal leitakse uusi kirjeid, sisestada need kirjed DynamoDB tabelisse. Ilmselgelt peab sellel lambdal lisaks vaikeĂ”igustele olema juurdepÀÀs Kinesi andmevoo lugemisele ja kirjutamisele DynamoDB-sse.

IAM-rolli loomine lambda-funktsioonile collector
Alustuseks loome uue IAM rolli lambdaks nimega Lambda-TicketsProcessingRole:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Testimise nÀitena sobivad eel seadistatud poliitikad AmazonKinesisReadOnlyAccess ja AmazonDynamoDBFullAccess, nagu on kujutatud alloleval pildil:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

See lambda tuleb kÀivitada Kinesi kÀivituse tÔttu uute kirjeid saabumise korral airline_stream voogu, seetÔttu tuleb lisada uus kÀivitaja:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
NĂŒĂŒd tuleb lisada kood ja salvestada lambda.

"""Voogu analĂŒĂŒsides ja DynamoDB tabelisse sisestades."""
import base64
import json
import boto3
from decimal import Decimal

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

class TicketsParser:
    """Voogust info analĂŒĂŒsimine."""

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

    @staticmethod
    def get_json_data(records):
        """Tagastab deserialiseeritud andmed voost."""
        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):
        """Eelprotsessib 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):
        """Tabelisse massilisamine."""
        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('Lisatud on ', len(self.json_data), 'eset')

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

Notifieri lambda-funktsiooni loomine

Teine lambda-funktsioon, mis jÀlgib teist voolu (special_stream) ja saadab teavituse SNS-ile, luuakse sarnaselt. Seega peab see lambda omama lugemisÔigust Kinesis'elt ja saatma sÔnumeid mÀÀratud SNS-teemasse, mis edastatakse seejÀrel kÔigile selle teema tellijatele (e-post, SMS jne).

IAM rolli loomine
Esmalt loome IAM rolli Lambda-KinesisAlarm selle lambda jaoks ja seejÀrel mÀÀrame selle rolli loodavale lambda alarm_notifier.

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

See lambda peaks töötama sissetulevate uute kirjade voolu special_stream pÔhjalise triggerega, seega on vajalik seadistada trigger sarnaselt sellele, kuidas me seda tegime lambda Collector jaoks.

Selle lambda seadistamise mugavuse huvides loome uue keskkonnamuutujate — TOPIC_ARN, kuhu sisestame ANR (Amazon Resource Names) teema Airlines.

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Ja sisestame lambda koodi, see on tÀiesti lihtne:

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 sisu!',
                           Subject='Lennupiletite hÀire')
        print('HĂ€ireteade on edukalt edastatud')
    except Exception as err:
        print('Edastamise viga', str(err))

Tundub, et kĂ€sitsi seadistamine on nĂŒĂŒd lĂ”pule viidud. JÀÀb vaid testida ja veenduda, et oleme kĂ”ik Ă”igesti seadistunud.

Koodi Terraform'i kaudu rakendamine

Vajaliku ettevalmistus

Terraform — vĂ€ga mugav avatud lĂ€htekoodiga tööriist infrastruktuuri juurutamiseks koodist. Sellel on oma sĂŒntaks, mida on lihtne omandada ja palju nĂ€iteid, kuidas ning mida juurutada. Redaktorites like Atom vĂ”i Visual Studio Code on palju mugavaid pluginaid, mis lihtsustavad Terraform'i kasutamist.

Distribuuti saab alla laadida siit. Terraform'i kĂ”igi vĂ”imaluste ĂŒksikasjalik analĂŒĂŒs ĂŒletab selle artikli piire, seega piirdume peamiste punktidega.

Kuidas kÀivitada

Kogu projekti kood asub minu hoidlas. Kloonime endale hoidla. Enne kÀivitamist tuleb veenduda, et AWS CLI on installitud ja seadistatud, kuna Terraform otsib mandaate failist ~/ .aws / credentials.

Heaks praktikaks on enne terve infrastruktuuri juurutamist kasutada kÀsku plan, et nÀha, mida Terraform praegu meie jaoks pilve loob:

terraform.exe plan

KĂŒsimine telefoninumbri sisestamiseks, et saada teateid, on kohustuslik. Selles etapis ei ole selle sisestamine vajalik.

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Programmide tööplaani analĂŒĂŒsides saame alustada ressursside loomist:

terraform.exe apply

PĂ€rast selle kĂ€su saatmist kĂŒsitakse uuesti telefoninumbri sisestamist. Kui kĂŒsitakse, kas tegi tĂ”eliselt toimingut, valige "jah". See vĂ”imaldab tĂ”sta kĂ”iki infrastruktuure, teha vajalikud EC2 seadistused, kĂ€ivitada Lambda funktsioonid jne.

PĂ€rast seda, kui kĂ”ik ressursid on Terraformi koodi kaudu edukalt loodud, peate minema Kinesis Analytics rakenduse ĂŒksikasjadesse (kahjuks ei leidnud ma, kuidas seda otse koodist teha).

KĂ€ivitame rakenduse:

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

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
Aviasales API integreerimine Amazon Kinesis 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 ĂŒhtlaselt.

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

sudo .\/api_caller.py TOKEN

JÀÀnud on oodata SMS sÔnumit teie numbrile:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
SMS - sÔnum jÔuab telefoni peaaegu minuti jooksul:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus
On jÀÀnud vaadata, kas andmed on salvestatud DynamoDB andmebaasi, et teha hiljem pĂ”hjalikumat analĂŒĂŒsi. Tabel airline_tickets sisaldab ligikaudu selliseid andmeid:

Aviasales API integreerimine Amazon Kinesis ja serverless lihtsus

KokkuvÔte

Tehtud töö kĂ€igus loodi veebipĂ”hine andmete töötlemise sĂŒsteem Amazon Kinesisel. Uuriti Kinesis Agendi kasutamise vĂ”imalusi koos Kinesis Data Streams'i ja reaalajas analĂŒĂŒsi Kinesis Analytics'i SQL kĂ€skude abil, samuti Amazon Kinesis’i koostööd teiste AWS teenustega.

Ülaltoodud sĂŒsteemi juurutati kahes variandis: aeglaselt kĂ€sitsi ja kiirelt Terraformi koodi kaudu.

Kogu projekti allikakood on saadaval minu GitHubi hoidlas, soovitan sellega tutvuda.

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

Soovin edu!

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster