Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Salut, Habr !

Aimez-vous prendre l'avion ? J'adore cela, mais pendant le confinement, j'ai Ă©galement appris Ă  analyser les donnĂ©es sur les billets d'avion d'une ressource bien connue — Aviasales.

Aujourd'hui, nous allons examiner Amazon Kinesis, créer un systÚme de streaming avec une analyse en temps réel, mettre en place une base de données NoSQL Amazon DynamoDB comme principal stockage de données et configurer des alertes par SMS pour les billets intéressants.

Tous les détails ci-dessous ! Allons-y !

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Introduction

Pour l'exemple, nous aurons besoin d'accÚs à l'API Aviasales. L'accÚs est gratuit et sans restrictions, il suffit de s'inscrire dans la section « Développeurs » pour obtenir votre jeton API pour accéder aux données.

L'objectif principal de cet article est de donner une compréhension générale de l'utilisation de la diffusion de données dans AWS. Nous omettons le fait que les données renvoyées par l'API utilisée ne sont pas strictement à jour et proviennent d'un cache, qui est généré en fonction des recherches des utilisateurs des sites Aviasales.ru et Jetradar.com au cours des 48 derniÚres heures.

Les données sur les billets d'avion obtenues via l'API seront automatiquement analysées par Kinesis-agent, installé sur la machine productrice, et transmises au bon flux via Kinesis Data Analytics. La version brute de ce flux sera directement écrite dans le stockage. Le stockage "brut" déployé dans DynamoDB permettra d'effectuer une analyse plus approfondie des billets à l'aide d'outils BI, par exemple AWS Quick Sight.

Nous examinerons deux options de déploiement de l'ensemble de l'infrastructure :

  • Manuel — via la console de gestion AWS ;
  • Infrastructure en code Terraform — pour les automateurs paresseux ;

Architecture du systÚme en cours de développement

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Composants utilisés :

  • API Aviasales — les donnĂ©es renvoyĂ©es par cette API seront utilisĂ©es pour tout le travail ultĂ©rieur ;
  • Instance EC2 Producer — une machine virtuelle classique dans le cloud, qui gĂ©nĂ©rera le flux de donnĂ©es d'entrĂ©e :
    • Kinesis Agent — il s'agit d'une application Java installĂ©e localement sur la machine, qui fournit un moyen simple de collecter et d'envoyer des donnĂ©es Ă  Kinesis (Kinesis Data Streams ou Kinesis Firehose). L'agent surveille en permanence un ensemble de fichiers dans les rĂ©pertoires spĂ©cifiĂ©s et envoie les nouvelles donnĂ©es Ă  Kinesis ;
    • Script API Caller — un script Python qui effectue des requĂȘtes Ă  l'API et stocke les rĂ©ponses dans un dossier surveillĂ© par le Kinesis Agent ;
  • Kinesis Data Streams — service de diffusion de donnĂ©es en temps rĂ©el avec de vastes possibilitĂ©s d'Ă©volutivitĂ© ;
  • Kinesis Analytics — un service sans serveur qui simplifie l'analyse des flux de donnĂ©es en temps rĂ©el. Amazon Kinesis Data Analytics configure les ressources pour le fonctionnement des applications et s'adapte automatiquement pour traiter tout volume de donnĂ©es entrant;
  • AWS Lambda — un service qui permet d'exĂ©cuter du code sans provisionner ni configurer des serveurs. Toutes les capacitĂ©s de calcul s'ajustent automatiquement pour chaque appel;
  • Amazon DynamoDB — une base de donnĂ©es de paires « clĂ©-valeur » et de documents, qui garantit un temps de rĂ©ponse infĂ©rieur Ă  10 millisecondes Ă  n'importe quelle Ă©chelle. Avec DynamoDB, il n'est pas nĂ©cessaire de provisionner des serveurs, d'appliquer des correctifs ou de gĂ©rer quoi que ce soit. DynamoDB ajuste automatiquement les tables en modifiant la quantitĂ© de ressources disponibles tout en maintenant des performances Ă©levĂ©es. Aucune action d'administration n'est requise;
  • Amazon SNS — un service de messagerie entiĂšrement gĂ©rĂ© basĂ© sur le modĂšle « Ă©diteur-abonnĂ© » (Pub/Sub), permettant d'isoler les microservices, les systĂšmes distribuĂ©s et les applications sans serveur. SNS peut ĂȘtre utilisĂ© pour envoyer des informations aux utilisateurs finaux via des notifications push mobiles, des SMS et des emails.

Préparation initiale

Pour simuler un flux de donnĂ©es, j'ai dĂ©cidĂ© d'utiliser les informations sur les billets d'avion retournĂ©es par l'API Aviasales. Dans documentation une liste assez variĂ©e de diffĂ©rents mĂ©thodes, prenons-en une — « Calendrier des prix du mois », qui retourne les prix pour chaque jour du mois, regroupĂ©s par nombre d'escales. Si aucun mois de recherche n'est spĂ©cifiĂ© dans la requĂȘte, les informations pour le mois suivant le mois en cours seront retournĂ©es.

Alors, inscrivons-nous et obtenons notre jeton.

Exemple de requĂȘte ci-dessous :

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

La mĂ©thode dĂ©crite ci-dessus pour obtenir des donnĂ©es de l'API avec mention du jeton dans la requĂȘte fonctionnera, mais je prĂ©fĂšre passer le jeton d'accĂšs via l'en-tĂȘte, donc dans le script api_caller.py, nous utiliserons ce moyen.

Exemple de réponse :

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

Dans l'exemple de rĂ©ponse de l'API ci-dessus, un billet est montrĂ© de Saint-PĂ©tersbourg Ă  Phuket
 Ah, pourquoi rĂȘver

Étant donnĂ© que je viens de Kazan et que Phuket est maintenant un rĂȘve lointain, cherchons des billets de Saint-PĂ©tersbourg Ă  Kazan.

On suppose que vous avez dĂ©jĂ  un compte AWS. Je tiens Ă  souligner que Kinesis et l'envoi de notifications via SMS ne sont pas inclus dans l'annĂ©e gratuite. Free Tier (utilisation gratuite). Mais mĂȘme avec quelques dollars en tĂȘte, il est tout Ă  fait possible de construire le systĂšme proposĂ© et de jouer avec. Et bien sĂ»r, n'oubliez pas de supprimer toutes les ressources une fois qu'elles ne sont plus nĂ©cessaires.

Heureusement, DynamoDb et les fonctions Lambda seront conditionnellement gratuits pour nous si nous restons dans les limites mensuelles gratuites. Par exemple, pour DynamoDB : 25 Go de stockage, 25 WCU/RCU et 100 millions de requĂȘtes. Et un million d'appels de fonctions Lambda par mois.

Déploiement manuel du systÚme

Configuration des Kinesis Data Streams

Accédons au service Kinesis Data Streams et créons deux nouveaux flux avec un shard chacun.

Qu'est-ce qu'un shard ?
Un shard est l'unité fondamentale de transmission des données dans le flux Amazon Kinesis. Un segment permet de transmettre les données entrantes à une vitesse de 1 Mo/s et les données sortantes à une vitesse de 2 Mo/s. Un segment prend en charge jusqu'à 1000 enregistrements PUT par seconde. Lors de la création d'un flux de données, il faut indiquer le nombre de segments souhaités. Par exemple, on peut créer un flux de données avec deux segments. Ce flux de données permettra une transmission des données entrantes à 2 Mo/s et des données sortantes à 4 Mo/s, en prenant en charge jusqu'à 2000 enregistrements PUT par seconde.

Plus il y a de shards dans votre flux, plus sa capacité d'acheminement est grande. En principe, les flux se scalent de cette maniÚre : en ajoutant des shards. Mais plus vous avez de shards, plus le prix est élevé. Chaque shard coûte 1,5 cent par heure et en plus 1,4 cent pour chaque million d'opérations d'ajout au flux (unités de charge utile PUT).

Créons un nouveau flux nommé airline_tickets, un seul shard suffira :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Créons maintenant un autre flux nommé special_stream:

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Configuration du producteur

En tant que producteur de donnĂ©es pour traiter la tĂąche, un simple instance EC2 suffira. Ce ne doit pas ĂȘtre une machine virtuelle coĂ»teuse et puissante, une t2.micro spot conviendra parfaitement.

Remarque importante : pour l'exemple, utilisez l'image - Amazon Linux AMI 2018.03.0, qui nécessite moins de configurations pour un démarrage rapide du Kinesis Agent.

Accédez au service EC2, créez une nouvelle machine virtuelle et choisissez l'AMI approprié de type t2.micro, qui fait partie de l'offre Free Tier :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Pour que la nouvelle machine virtuelle puisse interagir avec le service Kinesis, il est nĂ©cessaire de lui donner les droits nĂ©cessaires. Le meilleur moyen de le faire est d'attribuer un rĂŽle IAM. Par consĂ©quent, sur l'Ă©cran Étape 3 : Configurer les dĂ©tails de l'instance, sĂ©lectionnez CrĂ©er un nouveau rĂŽle IAM:

Création d'un rÎle IAM pour EC2
Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Dans la fenĂȘtre qui s'ouvre, sĂ©lectionnez que vous crĂ©ez un nouveau rĂŽle pour EC2 et passez Ă  la section Permissions :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Dans cet exemple d'apprentissage, il n'est pas nécessaire de se plonger dans tous les détails de la configuration granulaire des droits sur les ressources, donc nous allons choisir des politiques prédéfinies par Amazon : AmazonKinesisFullAccess et CloudWatchFullAccess.

Attribuons un nom significatif à ce rÎle, par exemple : EC2-KinesisStreams-FullAccess. En résultat, cela devrait ressembler à ce qui est indiqué sur l'image ci-dessous :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
AprÚs avoir créé ce nouveau rÎle, n'oubliez pas de l'attacher à l'instance de machine virtuelle que vous créez :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Nous ne modifions rien d'autre Ă  cet Ă©cran et passons aux fenĂȘtres suivantes.

Les paramÚtres du disque dur peuvent rester par défaut, les balises aussi (bien que, il soit préférable d'utiliser des balises, au moins pour donner un nom à l'instance et indiquer l'environnement).

Nous sommes maintenant sur l'onglet Étape 6 : Configurer le groupe de sĂ©curitĂ©, oĂč vous devez crĂ©er un nouveau groupe de sĂ©curitĂ© ou indiquer celui que vous avez, qui permet de se connecter via ssh (port 22) Ă  l'instance. SĂ©lectionnez lĂ  Source —> Mon IP et vous pouvez lancer l'instance.

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
DĂšs qu'il passe au statut de fonctionnement, vous pouvez essayer de vous y connecter via ssh.

Pour pouvoir travailler avec Kinesis Agent, aprĂšs vous ĂȘtre connectĂ© avec succĂšs Ă  la machine, vous devez entrer les commandes suivantes dans le 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

Créons un dossier pour sauvegarder les réponses de l'API :

sudo mkdir /var/log/airline_tickets

Avant de lancer l'agent, vous devez configurer son fichier de configuration :

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

Le contenu du fichier agent.json doit avoir la forme suivante :

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

Comme indiqué dans le fichier de configuration, l'agent surveillera les fichiers avec l'extension .log dans le répertoire /var/log/airline_tickets/, les analysera et les transmettra au flux airline_tickets.

Redémarrons le service et vérifions qu'il s'est lancé et fonctionne :

sudo service aws-kinesis-agent restart

Téléchargeons maintenant le script Python qui interrogera les données de l'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

Le script api_caller.py interroge les donnĂ©es d'Aviasales et enregistre la rĂ©ponse reçue dans le rĂ©pertoire surveillĂ© par l'agent Kinesis. La mise en Ɠuvre de ce script est assez standard, il y a une classe TicketsApi qui permet d'appeler l'API de maniĂšre asynchrone. Dans cette classe, nous passons l'en-tĂȘte avec le token et les paramĂštres de la requĂȘte :

class TicketsApi:
    """Classe d'appel d'API."""

    def __init__(self, headers):
        """Méthode d'initialisation."""
        self.base_url = BASE_URL
        self.headers = headers

    async def get_data(self, data):
        """Obtenir les donnĂ©es de la requĂȘte 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('Statut de la réponse %s : %s',
                            self.base_url, response.status)
                response_json = await response.json()
            except HTTPError as http_err:
                LOGGER.error('Oups ! Une erreur HTTP est survenue : %s', str(http_err))
            except Exception as err:
                LOGGER.error('Oups ! Une erreur est survenue : %s', str(err))
            return response_json


def prepare_request(api_token):
    """Retourne les en-tĂȘtes et la requĂȘte pour la demande 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():
    """Lancer le 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('L’API a renvoyĂ© %s Ă©lĂ©ments', len(response['data']))
        try:
            count_rows = log_maker(response)
            LOGGER.info('%s lignes ont été enregistrées dans %s',
                        count_rows,
                        TARGET_FILE)
        except Exception as e:
            LOGGER.error('Oups ! Le rĂ©sultat de la requĂȘte n’a pas Ă©tĂ© enregistrĂ© dans le fichier. %s',
                         str(e))
    else:
        LOGGER.error('Oups ! La requĂȘte API a Ă©chouĂ© %s !', response)

Pour tester la configuration et le bon fonctionnement de l'agent, effectuons un test du script api_caller.py :

sudo ./api_caller.py TOKEN

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Et observons le résultat dans les journaux de l'Agent et dans l'onglet Monitoring du flux de données airline_tickets :

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

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Comme vous pouvez le constater, tout fonctionne et l'agent Kinesis envoie avec succÚs des données au flux. Configurez maintenant le consommateur.

Configuration de Kinesis Data Analytics

Passons au composant central de tout le systĂšme — crĂ©ons une nouvelle application dans Kinesis Data Analytics nommĂ©e kinesis_analytics_airlines_app :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Kinesis Data Analytics permet d'effectuer une analyse de données en temps réel à partir de Kinesis Streams en utilisant le langage SQL. Il s'agit d'un service entiÚrement évolutif (contrairement à Kinesis Streams) qui :

  1. permet de crĂ©er de nouveaux flux (Output Stream) sur la base des requĂȘtes aux donnĂ©es d'origine ;
  2. fournit un flux d'erreurs survenant pendant l'exécution des applications (Error Stream) ;
  3. peut automatiquement dĂ©terminer le schĂ©ma des donnĂ©es d'entrĂ©e (celui-ci peut ĂȘtre redĂ©fini manuellement si nĂ©cessaire).

C'est un service coĂ»teux — 0,11 USD par heure d'utilisation, il convient donc de l'utiliser avec prĂ©caution et de le supprimer Ă  la fin de l'utilisation.

Connectons l'application à la source de données :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Sélectionnez le flux auquel nous prévoyons de nous connecter (airline_tickets) :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Ensuite, il est nécessaire d'attacher un nouveau rÎle IAM afin que l'application puisse lire et écrire dans le flux. Il suffit de ne rien changer dans le bloc des autorisations d'accÚs :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Nous allons maintenant demander la découverte du schéma des données dans le flux, en cliquant sur le bouton « Discover schema ». Cela mettra à jour (ou créera un nouveau) rÎle IAM et lancera la découverte du schéma à partir des données qui ont déjà été envoyées dans le flux :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Nous devons maintenant accĂ©der Ă  l'Ă©diteur SQL. En cliquant sur ce bouton, une fenĂȘtre s'ouvrira avec une question sur le dĂ©marrage de l'application — sĂ©lectionnons ce que nous voulons lancer :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Dans la fenĂȘtre de l'Ă©diteur SQL, insĂ©rons cette simple requĂȘte et appuyons sur 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';

Dans les bases de donnĂ©es relationnelles, vous travaillez avec des tables, en utilisant des opĂ©rateurs INSERT pour ajouter des enregistrements et l'opĂ©rateur SELECT pour interroger des donnĂ©es. Dans Amazon Kinesis Data Analytics, vous travaillez avec des flux (STREAM) et des « pompes » (PUMP) — des requĂȘtes d'insertion continues qui insĂšrent des donnĂ©es d'un flux dans l'application Ă  un autre flux.

Dans la requĂȘte SQL prĂ©sentĂ©e ci-dessus, nous recherchons des billets Aeroflot coĂ»tant moins de cinq mille roubles. Tous les enregistrements rĂ©pondant Ă  ces critĂšres seront placĂ©s dans le flux DESTINATION_SQL_STREAM.

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Dans le bloc Destination, sélectionnez le flux special_stream, et dans le menu déroulant Nom du flux dans l'application, choisissez DESTINATION_SQL_STREAM :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
À la suite de toutes ces manipulations, vous devriez obtenir quelque chose ressemblant à l'image ci-dessous :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Création et abonnement à un sujet SNS

Accédez au service Simple Notification Service et créez un nouveau sujet nommé Airlines :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Abonnez-vous à ce sujet, en indiquant le numéro de téléphone mobile sur lequel les notifications SMS seront envoyées :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Création d'une table dans DynamoDB

Pour stocker les donnĂ©es brutes de leur flux airline_tickets, crĂ©ons une table dans DynamoDB du mĂȘme nom. Nous utiliserons record_id comme clĂ© primaire :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Création d'une fonction lambda collector

CrĂ©ons une fonction lambda appelĂ©e Collector, dont la tĂąche sera d'interroger le flux airline_tickets et, en cas de nouvelles entrĂ©es trouvĂ©es, d'insĂ©rer ces entrĂ©es dans la table DynamoDB. Évidemment, en plus des droits par dĂ©faut, cette lambda doit avoir accĂšs en lecture au flux de donnĂ©es Kinesis et en Ă©criture Ă  DynamoDB.

Création d'un rÎle IAM pour la fonction lambda collector
Pour commencer, créons un nouveau rÎle IAM pour la lambda nommé Lambda-TicketsProcessingRole :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Pour cet exemple de test, les politiques prĂȘtes Ă  l'emploi AmazonKinesisReadOnlyAccess et AmazonDynamoDBFullAccess conviendront parfaitement, comme le montre l'image ci-dessous :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Cette lambda doit ĂȘtre dĂ©clenchĂ©e par Kinesis lors de l'arrivĂ©e de nouvelles entrĂ©es dans le flux airline_stream, il faut donc ajouter un nouveau dĂ©clencheur :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Il ne reste plus qu'à insérer le code et à sauvegarder la lambda.

"""Analyse le flux et insérer dans la table DynamoDB."""
import base64
import json
import boto3
from decimal import Decimal

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

class TicketsParser:
    """Analyse les informations du flux."""

    def __init__(self, table_name, records):
        """Méthode d'initialisation."""
        self.table = DYNAMO_DB.Table(table_name)
        self.json_data = TicketsParser.get_json_data(records)

    @staticmethod
    def get_json_data(records):
        """Retourner les données désérialisées depuis le flux."""
        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):
        """Prétraiter les données 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):
        """Insertion en lot dans la 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('Ajouté ', len(self.json_data), 'éléments')

def lambda_handler(event, context):
    """Analyser le flux et insérer dans la table DynamoDB."""
    print('ÉvĂ©nement reçu :', event)
    parser = TicketsParser(TABLE_NAME, event['Records'])
    parser.run()

Création de la fonction lambda notifier

La deuxiÚme fonction lambda, qui surveillera le deuxiÚme flux (special_stream) et enverra une notification à SNS, est créée de maniÚre similaire. Par conséquent, cette lambda doit avoir accÚs à la lecture depuis Kinesis et à l'envoi de messages dans le sujet SNS spécifié, qui sera par la suite envoyé par le service SNS à tous les abonnés de ce sujet (email, SMS, etc.).

Création du rÎle IAM
Tout d'abord, nous créons le rÎle IAM Lambda-KinesisAlarm pour cette lambda, puis nous assignons ce rÎle à la lambda alarm_notifier que nous avons créée :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Cette lambda doit fonctionner sur un dĂ©clencheur lorsqu'il y a de nouvelles entrĂ©es dans le flux special_stream, nous devons donc configurer le dĂ©clencheur de la mĂȘme maniĂšre que nous l'avons fait pour la lambda Collector.

Pour faciliter la configuration de cette lambda, introduisons une nouvelle variable d'environnement — TOPIC_ARN, oĂč nous plaçons l'ARN (Amazon Resource Names) du sujet Airlines :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Et insérons le code de la lambda, il n'est pas compliqué du tout :

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='Salut ! J'ai trouvé quelque chose d'intéressant !',
                           Subject='Alerte des billets d'avion')
        print('Le message d’alerte a Ă©tĂ© livrĂ© avec succĂšs')
    except Exception as err:
        print('Échec de la livraison', str(err))

Il semble que la configuration manuelle du systÚme soit terminée. Il ne reste plus qu'à tester et s'assurer que tout est correctement configuré.

Déploiement à partir du code Terraform

Préparation nécessaire

Terraform — un outil open-source trĂšs pratique pour dĂ©ployer une infrastructure Ă  partir de code. Il a sa propre syntaxe, qui est facile Ă  maĂźtriser, ainsi que de nombreux exemples de ce que l'on peut dĂ©ployer. Dans l'Ă©diteur Atom ou Visual Studio Code, il existe de nombreux plugins pratiques qui facilitent le travail avec Terraform.

Le distribution peut ĂȘtre tĂ©lĂ©chargĂ©e d'ici. Une analyse dĂ©taillĂ©e de toutes les fonctionnalitĂ©s de Terraform dĂ©passe le cadre de cet article, limitons-nous donc aux points principaux.

Comment lancer

Le code complet du projet se trouve dans mon dĂ©pĂŽt. Clonons le dĂ©pĂŽt. Avant de lancer, assurez-vous d'avoir installĂ© et configurĂ© AWS CLI, car Terraform recherchera les informations d’identification dans le fichier ~/.aws/credentials.

Une bonne pratique avant de déployer toute l'infrastructure est d'exécuter la commande plan pour voir ce que Terraform va créer dans le cloud :

terraform.exe plan

Il vous sera demandĂ© de saisir un numĂ©ro de tĂ©lĂ©phone pour y envoyer des notifications. À ce stade, son entrĂ©e n'est pas obligatoire.

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
AprÚs avoir analysé le plan d'action du programme, nous pouvons commencer à créer des ressources :

terraform.exe apply

AprÚs avoir envoyé cette commande, une nouvelle demande de saisie du numéro de téléphone apparaßtra, tapez « yes » lorsque la question sur l'exécution réelle des actions sera posée. Cela permettra de déployer toute l'infrastructure, de configurer correctement EC2, de déployer des fonctions lambda, etc.

AprÚs que toutes les ressources soient créées avec succÚs via le code Terraform, il est nécessaire d'accéder aux détails de l'application Kinesis Analytics (malheureusement, je n'ai pas trouvé comment le faire directement à partir du code).

Lançons l'application :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Ensuite, il est nécessaire de spécifier explicitement le nom du flux en application, en choisissant dans le menu déroulant :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Maintenant, tout est prĂȘt Ă  fonctionner.

Test de fonctionnement de l'application

Peu importe comment vous avez dĂ©ployĂ© le systĂšme, manuellement ou via le code Terraform, il fonctionnera de la mĂȘme maniĂšre.

Connectez-vous via SSH Ă  la machine virtuelle EC2 oĂč le Kinesis Agent est installĂ© et exĂ©cutez le script api_caller.py

sudo ./api_caller.py TOKEN

Vous devez maintenant attendre un SMS sur votre numéro :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Le SMS — le message arrive sur le tĂ©lĂ©phone presque dans la minute :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless
Il reste à vérifier si les enregistrements ont été sauvegardés dans la base de données DynamoDB pour une analyse plus détaillée ultérieure. La table airline_tickets contient environ les données suivantes :

Intégration de l'API Aviasales avec Amazon Kinesis et simplicité serverless

Conclusion

Au cours de ce travail, un systÚme de traitement de données en ligne basé sur Amazon Kinesis a été construit. Nous avons examiné les options d'utilisation du Kinesis Agent en conjonction avec Kinesis Data Streams et l'analyse en temps réel avec Kinesis Analytics à l'aide de commandes SQL, ainsi que l'interaction d'Amazon Kinesis avec d'autres services AWS.

Nous avons déployé le systÚme décrit ci-dessus de deux maniÚres : de maniÚre suffisamment longue manuelle et rapidement par code Terraform.

L'intégralité du code source du projet est disponible dans mon dépÎt sur GitHub, je vous invite à le consulter.

Je suis prĂȘt Ă  discuter de l'article avec plaisir et j'attends vos commentaires. J'espĂšre une critique constructive.

Je vous souhaite du succĂšs !

Source : habr.com

Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS đŸ”„ Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster