Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Ciao, Habr!

E tu ami volare in aerei? Io lo adoro, ma durante il lockdown ho iniziato anche ad analizzare i dati sui biglietti aerei di un noto sito — Aviasales.

Oggi analizzeremo il funzionamento di Amazon Kinesis, costruiremo un sistema di streaming con analisi in tempo reale, utilizzeremo Amazon DynamoDB come database NoSQL principale e imposteremo avvisi via SMS per i biglietti interessanti.

Tutti i dettagli sotto il tag! Cominciamo!

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Introduzione

Per esempio, avremo bisogno dell'accesso a API Aviasales. L'accesso è gratuito e senza limiti, è necessario solo registrarsi nella sezione "Sviluppatori" per ottenere il proprio token API per accedere ai dati.

L'obiettivo principale di questo articolo è fornire una comprensione generale dell'utilizzo dello streaming di informazioni in AWS; escludiamo il fatto che i dati restituiti dall'API utilizzata non siano necessariamente aggiornati e vengano trasmessi da una cache che viene formata sulla base delle ricerche degli utenti sui siti Aviasales.ru e Jetradar.com negli ultimi 48 ore.

I dati sugli aerei ottenuti tramite API, Kinesis-agent, installato sulla macchina di produzione, parserà automaticamente e trasmetterà al flusso necessario tramite Kinesis Data Analytics. La versione non elaborata di questo flusso sarà scritta direttamente nel repository. La storage di 'raw' data distribuita in DynamoDB consentirà un'analisi più approfondita dei biglietti tramite strumenti BI, come AWS Quick Sight.

Esamineremo due opzioni per il deployment dell'intera infrastruttura:

  • Manuale — tramite AWS Management Console;
  • Infrastruttura basata su codice Terraform — per gli automatizzatori pigri;

Architettura del sistema in fase di sviluppo

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Componenti utilizzati:

  • Aviasales API — i dati restituiti da questa API saranno utilizzati per tutto il lavoro successivo;
  • EC2 Producer Instance — una normale macchina virtuale nel cloud, sulla quale verrà generato il flusso di dati in ingresso:
    • Kinesis Agent — è un'applicazione Java, installata localmente sulla macchina, che fornisce un modo semplice per raccogliere e inviare dati a Kinesis (Kinesis Data Streams o Kinesis Firehose). L'agente monitora costantemente un insieme di file in directory specificate e invia nuovi dati a Kinesis;
    • Script API Caller — script Python che esegue richieste all'API e salva le risposte in una cartella monitorata da Kinesis Agent;
  • Kinesis Data Streams — servizio di streaming dati in tempo reale con ampie capacità di scalabilità;
  • Kinesis Analytics — servizio serverless che semplifica l'analisi dei dati in streaming in tempo reale. Amazon Kinesis Data Analytics configura automaticamente le risorse richieste dalle applicazioni e scala per gestire qualsiasi volume di dati in ingresso;
  • AWS Lambda — servizio che consente di eseguire codice senza l'onere di provisioning e gestione dei server. Tutte le risorse di calcolo si ridimensionano automaticamente per ogni invocazione;
  • Amazon DynamoDB — database a coppie di "chiave-valore" e documenti che offre una latenza inferiore a 10 millisecondi a qualsiasi scala. Con DynamoDB non è necessario distribuire server, installare patch o gestirli. DynamoDB scala automaticamente le tabelle, regolando la quantità di risorse disponibili e mantenendo elevate prestazioni. Nessuna gestione del sistema è richiesta;
  • Amazon SNS — servizio di invio messaggi completamente gestito secondo il modello «publish-subscribe» (Pub/Sub), che consente di isolare microservizi, sistemi distribuiti e applicazioni serverless. SNS può essere utilizzato per inviare informazioni agli utenti finali tramite notifiche push mobili, messaggi SMS e email.

Preparazione iniziale

Per simulare il flusso di dati, ho deciso di utilizzare le informazioni sui biglietti aerei fornite dall'API di Aviasales. In documentazione un elenco piuttosto ampio di diversi metodi, prendiamo uno di essi — «Calendario dei prezzi mensile», che restituisce i prezzi per ogni giorno del mese, raggruppati per numero di scali. Se non si specifica il mese di ricerca nella richiesta, verranno restituite le informazioni per il mese successivo all'attuale.

Quindi, registriamoci e otteniamo il nostro token.

Un esempio di richiesta è riportato qui sotto:

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

Il metodo sopra descritto per ottenere dati dall'API specificando il token nella richiesta funzionerà, ma a me piace di più passare il token di accesso tramite l'intestazione, quindi nel file api_caller.py utilizzeremo proprio questo metodo.

Esempio di risposta:

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

Nell'esempio di risposta API sopra è mostrato un biglietto da San Pietroburgo a Phuket… Eh, ma a cosa serve sognare…
Poiché sono di Kazan e Phuket è attualmente 'solo un sogno', cerchiamo biglietti da San Pietroburgo a Kazan.

Si presume che tu abbia già un account AWS. Voglio subito prestare particolare attenzione al fatto che Kinesis e l'invio di notifiche tramite SMS non rientrano nel piano annuale. Free Tier (utilizzo gratuito). Ma anche nonostante ciò, considerando qualche dollaro, è possibile costruire il sistema proposto e giocare con esso. E, naturalmente, non dimenticare di rimuovere tutte le risorse dopo che non servono più.

Fortunatamente, DynamoDB e le funzioni Lambda saranno per noi sostanzialmente gratuiti, se rimaniamo entro i limiti mensili gratuiti. Ad esempio, per DynamoDB: 25 GB di spazio di archiviazione, 25 WCU/RCU e 100 milioni di richieste. E un milione di invocazioni di funzioni Lambda al mese.

Distribuzione manuale del sistema

Impostazione di Kinesis Data Streams

Passiamo al servizio Kinesis Data Streams e creiamo due nuovi stream con un shard ciascuno.

Che cos'è uno shard?
Uno shard è l'unità fondamentale di trasmissione dei dati dello stream di Amazon Kinesis. Un segmento fornisce una capacità di ingresso di dati pari a 1 MB/s e un'uscita di dati pari a 2 MB/s. Un segmento supporta fino a 1000 operazioni PUT al secondo. Quando si crea uno stream di dati, è necessario specificare il numero di segmenti desiderato. Ad esempio, è possibile creare uno stream di dati con due segmenti. Questo stream offrirà una capacità di ingresso di dati pari a 2 MB/s e un'uscita di dati pari a 4 MB/s, supportando fino a 2000 operazioni PUT al secondo.

Maggiore è il numero di shard nel tuo stream, maggiore sarà la sua capacità di throughput. In linea di principio, gli stream si scalano in questo modo: aggiungendo shard. Tuttavia, più shard hai, più alto sarà il costo. Ogni shard costa 1,5 centesimi all'ora e ulteriori 1,4 centesimi per ogni milione di operazioni di aggiunta allo stream (PUT payload units).

Creiamo un nuovo stream chiamato airline_tickets, avrà bisogno di un solo shard:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Ora creiamo un altro stream chiamato special_stream:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Configurazione del produttore

Come produttore di dati per l'analisi del task, è sufficiente utilizzare una normale istanza EC2. Non è necessario un server virtuale potente e costoso, va benissimo un t2.micro spot.

Nota importante: per l'esempio, si dovrebbe usare l'immagine — Amazon Linux AMI 2018.03.0, poiché richiede meno configurazioni per un avvio rapido dell'Agent Kinesis.

Passiamo al servizio EC2, creiamo una nuova macchina virtuale, scegliamo l'AMI desiderata con tipo t2.micro, che rientra nel Free Tier:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Per permettere alla nuova macchina virtuale di interagire con il servizio Kinesis, è necessario conferirle i diritti necessari. Il modo migliore per farlo è assegnare un IAM Role. Pertanto, nella schermata Step 3: Configure Instance Details, dobbiamo selezionare Crea nuovo IAM Role:

Creazione di un IAM Role per EC2
Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Nella finestra aperta, scegliamo di creare un nuovo ruolo per EC2 e passiamo alla sezione Permessi:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Nel contesto di questa formazione, non è necessario entrare nei dettagli della configurazione granulare dei permessi sulle risorse, quindi selezioniamo le policy preimpostate da Amazon: AmazonKinesisFullAccess e CloudWatchFullAccess.

Diamo un nome significativo a questo ruolo, ad esempio: EC2-KinesisStreams-FullAccess. Dovrebbe risultare come illustrato nell'immagine qui sotto:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Dopo aver creato questo nuovo ruolo, non dimentichiamo di associarlo all'istanza della macchina virtuale che stiamo creando:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Non apportiamo ulteriori modifiche su questa schermata e passiamo alle finestre successive.

Le impostazioni del disco rigido possono rimanere predefinite, anche le etichette (anche se è buona pratica utilizzare le etichette, almeno per dare un nome all'istanza e specificare l'ambiente).

Ora siamo nella scheda Passaggio 6: Configura il Gruppo di Sicurezza, dove è necessario creare un nuovo gruppo o specificare quello esistente che consente di connettersi tramite ssh (porta 22) all'istanza. Seleziona lì Sorgente --> Il mio IP e puoi avviare l'istanza.

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Non appena passerà allo stato in esecuzione, puoi provare a connetterti tramite ssh.

Per poter lavorare con Kinesis Agent, dopo esserti connesso con successo alla macchina, è necessario inserire i seguenti comandi nel terminale:

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

Creiamo una cartella per salvare le risposte dell'API:

sudo mkdir /var/log/airline_tickets

Prima di avviare l'agente, è necessario configurarlo:

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

Il contenuto del file agent.json deve avere il seguente aspetto:

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

Come visibile nel file di configurazione, l'agente monitorerà nella directory /var/log/airline_tickets/ i file con estensione .log, li parserà e li trasferirà nel flusso airline_tickets.

Riavviamo il servizio e ci assicuriamo che sia avviato e funzionante:

sudo service aws-kinesis-agent restart

Ora scarichiamo lo script Python che richiederà dati all'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

Lo script api_caller.py richiede dati da Aviasales e salva la risposta ottenuta nella directory scansionata dall'agente Kinesis. L'implementazione di questo script è abbastanza standard, c'è una classe TicketsApi, che consente di interrogare l'API in modo asincrono. In questa classe passiamo l'intestazione con il token e i parametri della richiesta:

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 ocurred: %s', str(err))
            return response_json


def prepare_request(api_token):
    """Return the headers and query fot 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 <your_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 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)

Per testare la correttezza delle impostazioni e il funzionamento dell'agente, eseguiremo un test dello script api_caller.py:

sudo ./api_caller.py TOKEN

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
E controlliamo il risultato nei log dell'agente e nella scheda Monitoring nel flusso di dati airline_tickets:

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

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Come si può vedere, tutto funziona e Kinesis Agent invia correttamente i dati nel flusso. Ora configuriamo il consumer.

Configurazione di Kinesis Data Analytics

Passiamo al componente centrale dell'intero sistema: creiamo una nuova applicazione in Kinesis Data Analytics con il nome kinesis_analytics_airlines_app:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Kinesis Data Analytics consente di eseguire analisi dei dati in tempo reale da Kinesis Streams utilizzando il linguaggio SQL. Si tratta di un servizio completamente scalabile (a differenza di Kinesis Streams) che:

  1. consente di creare nuovi flussi (Output Stream) basati su query sui dati di origine;
  2. fornisce un flusso di errori che si sono verificati durante l'esecuzione delle applicazioni (Error Stream);
  3. è in grado di determinare automaticamente lo schema dei dati in ingresso (che può essere ridefinito manualmente se necessario).

Questo è un servizio piuttosto costoso: 0,11 USD all'ora di funzionamento, quindi è consigliabile usarlo con cautela e rimuoverlo al termine dell'utilizzo.

Colleghiamo l'applicazione alla fonte di dati:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Scegliamo il flusso a cui ci connettiamo (airline_tickets):

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
In seguito, è necessario allegare un nuovo ruolo IAM affinché l'applicazione possa leggere e scrivere nel flusso. A tal fine, non è necessario modificare nulla nel blocco Access permissions:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Ora richiediamo la scoperta dello schema dei dati nel flusso, per fare ciò premiamo il pulsante «Discover schema». Di conseguenza, verrà aggiornata (creata una nuova) IAM Role e verrà avviata la scoperta dello schema dai dati già arrivati nel flusso:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Ora è necessario andare nell'editor SQL. Facendo clic su questo pulsante, apparirà una finestra con la domanda riguardante l'avvio dell'applicazione — scegliamo cosa vogliamo avviare:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Nella finestra dell'editor SQL inseriremo una semplice query e clicchiamo su 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';

Nei database relazionali, lavori con tabelle utilizzando l'operatore INSERT per aggiungere record e l'operatore SELECT per interrogare i dati. In Amazon Kinesis Data Analytics, lavori con flussi (STREAM) e "pompe" (PUMP) — richieste continue di inserimento che trasferiscono dati da un flusso all'altro all'interno dell'applicazione.

Nella query SQL sopra riportata, viene effettuata una ricerca dei biglietti Aeroflot a un prezzo inferiore a cinquemila rubli. Tutti i record che soddisfano questi criteri verranno inseriti nel flusso DESTINATION_SQL_STREAM.

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Nella sezione Destinazione, selezioniamo il flusso special_stream e nel menu a tendina del nome del flusso applicativo DESTINATION_SQL_STREAM:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Il risultato di tutte queste operazioni dovrebbe somigliare all'immagine sottostante:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Creazione e iscrizione a un argomento SNS

Accediamo al servizio Simple Notification Service e creiamo un nuovo argomento con il nome Airlines:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Effettuiamo l'iscrizione a questo argomento, specificando il numero di telefono cellulare a cui verranno inviate le notifiche SMS:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Creazione di una tabella in DynamoDB

Per memorizzare i dati non elaborati del flusso airline_tickets, creeremo una tabella in DynamoDB con lo stesso nome. Utilizzeremo record_id come chiave primaria:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Creazione della funzione Lambda collector

Creeremo una funzione Lambda chiamata Collector, il cui compito sarà quello di interrogare il flusso airline_tickets e, nel caso vengano trovati nuovi record, inserire questi record nella tabella DynamoDB. È evidente che, oltre ai permessi di default, questa Lambda deve avere accesso alla lettura del flusso di dati Kinesis e alla scrittura in DynamoDB.

Creazione di un ruolo IAM per la funzione Lambda collector
Cominciamo creando un nuovo ruolo IAM per la Lambda chiamato Lambda-TicketsProcessingRole:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Per un esempio di test andranno benissimo le policy preimpostate AmazonKinesisReadOnlyAccess e AmazonDynamoDBFullAccess, come mostrato nell'immagine qui sotto:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Questa Lambda deve essere attivata tramite un trigger da Kinesis quando nuovi record entrano nel flusso airline_stream, quindi bisogna aggiungere un nuovo trigger:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Rimane da inserire il codice e salvare la Lambda.

"""Analisi del flusso e inserimento nella tabella DynamoDB."""
import base64
import json
import boto3
from decimal import Decimal

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

class TicketsParser:
    """Analizza le informazioni dal flusso."""

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

    @staticmethod
    def get_json_data(records):
        """Restituisce i dati deserializzati dal flusso."""
        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-processa i dati 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):
        """Inserimento in batch nella tabella."""
        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('Sono stati aggiunti ', len(self.json_data), 'elementi')

def lambda_handler(event, context):
    """Analizza il flusso e inserisci nella tabella DynamoDB."""
    print('Evento ricevuto:', event)
    parser = TicketsParser(TABLE_NAME, event['Records'])
    parser.run()

Creazione della funzione lambda notifier

La seconda funzione lambda, che monitorerà il secondo flusso (special_stream) e invierà notifiche a SNS, viene creata in modo simile. Pertanto, questa lambda deve avere accesso in lettura da Kinesis e inviare messaggi al topic SNS specificato, che successivamente il servizio SNS invierà a tutti gli iscritti a quel topic (email, SMS, ecc.).

Creazione del ruolo IAM
Iniziamo creando il ruolo IAM Lambda-KinesisAlarm per questa lambda e poi assegniamo questo ruolo alla lambda che stiamo creando, alarm_notifier:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Questa lambda deve funzionare su un trigger quando nuove registrazioni arrivano nel flusso special_stream, quindi è necessario configurare il trigger in modo simile a come abbiamo fatto per la lambda Collector.

Per facilitare la configurazione di questa lambda, introduciamo una nuova variabile d'ambiente — TOPIC_ARN, dove inseriamo l'ARN (Amazon Resource Names) del topic Airlines:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
E incolliamo il codice della lambda, che è piuttosto semplice:

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='Ciao! Ho trovato qualcosa di interessante!',
                           Subject='Allerta biglietti aerei')
        print('Il messaggio di allerta è stato consegnato con successo')
    except Exception as err:
        print('Errore di consegna', str(err))

Sembra che la configurazione manuale del sistema sia completa. Resta solo da testare e assicurarsi che tutto sia impostato correttamente.

Deploy da codice Terraform

Preparazione necessaria

Terraform — è uno strumento open-source molto pratico per il deployment dell'infrastruttura da codice. Ha una propria sintassi, facile da apprendere, e molteplici esempi su come e cosa distribuire. Nel editor Atom o Visual Studio Code ci sono molti plugin utili che semplificano il lavoro con Terraform.

Il distributore può essere scaricato da qui. Un'analisi dettagliata di tutte le funzionalità di Terraform va oltre il contesto di questo articolo, quindi ci limiteremo ai punti principali.

Come eseguire

Il codice completo del progetto è disponibile nel mio repository. Dobbiamo clonare il repository. Prima di eseguire, assicurati di avere installato e configurato AWS CLI, poiché Terraform cercherà le credenziali nel file ~/.aws/credentials.

È una buona pratica eseguire il comando plan prima di deployare l'intera infrastruttura, per vedere cosa Terraform creerà attualmente nel cloud:

terraform.exe plan

Verrà richiesto di inserire il numero di telefono per ricevere le notifiche. In questa fase non è necessario fornirlo.

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Analizzando il piano di lavoro del programma, possiamo avviare la creazione delle risorse:

terraform.exe apply

Dopo aver inviato questo comando, apparirà di nuovo la richiesta di inserire il numero di telefono; digita 'yes' quando viene mostrata la domanda sull'effettiva esecuzione delle azioni. Questo consentirà di sollevare l'intera infrastruttura, effettuare tutte le configurazioni necessarie per EC2, distribuire le funzioni Lambda, ecc.

Dopo che tutte le risorse sono state create con successo tramite codice Terraform, è necessario accedere ai dettagli dell'applicazione Kinesis Analytics (sfortunatamente, non ho trovato come farlo direttamente dal codice).

Avviamo l'applicazione:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Dopo di ciò, è necessario definire esplicitamente il nome dello stream in-app, scegliendo dall'elenco a discesa:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Ora tutto è pronto per funzionare.

Testare il funzionamento dell'applicazione

Indipendentemente da come hai distribuito il sistema, manualmente o tramite codice Terraform, funzionerà allo stesso modo.

Accediamo tramite SSH alla macchina virtuale EC2, dove è installato Kinesis Agent, e avviamo lo script api_caller.py.

sudo ./api_caller.py TOKEN

Resta solo da attendere l'SMS al tuo numero:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
L'SMS - il messaggio arriva sul telefono praticamente dopo 1 minuto:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless
Resta da verificare se ci sono registrazioni nel database DynamoDB per analisi più dettagliate in seguito. La tabella airline_tickets contiene dati simili a questi:

Integrazione dell'API Aviasales con Amazon Kinesis e semplicità serverless

Conclusione

Nel corso del lavoro svolto, è stato costruito un sistema di elaborazione dei dati online basato su Amazon Kinesis. Sono state esplorate opzioni per utilizzare Kinesis Agent insieme a Kinesis Data Streams e l'analisi in tempo reale con Kinesis Analytics tramite comandi SQL, oltre all'interazione di Amazon Kinesis con altri servizi AWS.

Il sistema descritto è stato implementato in due modi: uno lungo e manuale e uno veloce dal codice Terraform.

Tutto il codice sorgente del progetto è disponibile nel mio repository su GitHub, vi invito a darci un'occhiata.

Sono pronto a discutere l'articolo e attendo i vostri commenti. Spero in una critica costruttiva.

Vi auguro successo!

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster