Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Hallo, Habr!

Lieben Sie es, mit dem Flugzeug zu reisen? Ich tue das, und während der Selbstisolation habe ich mich zusätzlich darin vertieft, Daten zu Flugtickets von einer bekannten Plattform – Aviasales – zu analysieren.

Heute werden wir die Funktionsweise von Amazon Kinesis untersuchen, ein Streaming-System mit Echtzeitanalysen aufbauen, die NoSQL-Datenbank Amazon DynamoDB als unser primäres Datenspeicher implementieren und Benachrichtigungen über SMS für interessante Flugtickets einrichten.

Alle Details finden Sie im Folgenden! Lassen Sie uns beginnen!

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Einführung

Für unser Beispiel benötigen wir Zugang zu API Aviasales. Der Zugang ist kostenlos und ohne Einschränkungen verfügbar; man muss sich lediglich im Bereich „Entwickler“ registrieren, um seinen API-Token für den Zugriff auf die Daten zu erhalten.

Das Hauptziel dieses Artikels ist es, ein allgemeines Verständnis für die Verwendung von Datenstromübertragung in AWS zu vermitteln. Wir nehmen hierbei in Kauf, dass die über die verwendete API zurückgegebenen Daten nicht unbedingt aktuell sind und aus einem Cache stammen, der auf Suchanfragen von Nutzern der Websites Aviasales.ru und Jetradar.com der letzten 48 Stunden basiert.

Die über die API erhaltenen Daten zu Flugtickets werden vom Kinesis-Agent, der auf der Produzentenmaschine installiert ist, automatisch analysiert und an den entsprechenden Stream über Kinesis Data Analytics gesendet. Die unbearbeitete Version dieses Streams wird direkt in den Speicher geschrieben. Die in DynamoDB implementierte Speicherung von "rohen" Daten ermöglicht eine tiefere Analyse von Tickets mithilfe von BI-Tools wie AWS QuickSight.

Wir betrachten zwei Optionen für die Bereitstellung der gesamten Infrastruktur:

  • Manuell — über die AWS Management Console;
  • Infrastruktur aus Terraform-Code — für die faulen Automatisierer;

Architektur des zu entwickelnden Systems

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Verwendete Komponenten:

  • Aviasales API — Die von dieser API zurückgegebenen Daten werden für alle folgenden Arbeiten verwendet;
  • EC2 Producer Instance — eine gewöhnliche virtuelle Maschine in der Cloud, die den Eingangsdatastream generieren wird:
    • Kinesis Agent — eine Java-Anwendung, die lokal auf der Maschine installiert ist und eine einfache Möglichkeit bietet, Daten an Kinesis (Kinesis Data Streams oder Kinesis Firehose) zu sammeln und zu senden. Der Agent überwacht ständig eine Gruppe von Dateien in den angegebenen Verzeichnissen und sendet neue Daten an Kinesis;
    • API Caller-Skript — Python-Skript, das Anfragen an die API stellt und die Antworten in einen von Kinesis Agent überwachten Ordner speichert;
  • Kinesis Data Streams — ein Echtzeit-Datenstreaming-Dienst mit umfangreichen Skalierungsmöglichkeiten;
  • Kinesis Analytics — ein serverloser Dienst, der die Analyse von Streaming-Daten in Echtzeit vereinfacht. Amazon Kinesis Data Analytics konfiguriert Ressourcen für den Betrieb von Anwendungen und skaliert automatisch, um beliebige Mengen eingehender Daten zu verarbeiten;
  • AWS Lambda — ein Dienst, der das Ausführen von Code ohne Bereitstellung und Konfiguration von Servern ermöglicht. Alle Rechenressourcen skalieren automatisch mit jedem Aufruf;
  • Amazon DynamoDB — eine Schlüssel-Wert-Datenbank und Dokumentendatenbank, die eine Verzögerung von weniger als 10 Millisekunden bei jedem Maßstab bietet. Bei der Verwendung von DynamoDB müssen keine Server bereitgestellt, gepatcht oder verwaltet werden. DynamoDB skaliert Tabellen automatisch, passt die Verfügbarkeit von Ressourcen an und sorgt für hohe Leistung. Es sind keine Verwaltungsmaßnahmen erforderlich;
  • Amazon SNS — ein vollständig verwalteter Messaging-Dienst nach dem Publish/Subscribe-Modell, mit dem sich Mikrodienste, verteilte Systeme und serverlose Anwendungen isolieren lassen. SNS kann verwendet werden, um Informationen an Endbenutzer über mobile Push-Benachrichtigungen, SMS und E-Mails zu versenden.

Erste Vorbereitung

Um einen Datenstrom zu emulieren, habe ich mich entschieden, Informationen über Flugtickets zu verwenden, die von der API Aviasales zurückgegeben werden. In Dokumentation. eine recht umfangreiche Liste verschiedener Methoden, nehmen wir eine davon – "Monatskartenpreis", die die Preise für jeden Tag des Monats zurückgibt, gruppiert nach der Anzahl der Umstiege. Wenn im Anfrage keine Suchmonat übergeben wird, werden die Informationen für den Monat zurückgegeben, der auf den aktuellen folgt.

Also, wir registrieren uns und erhalten unser Token.

Beispielanfrage unten:

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

Der oben beschriebene Ansatz zur Datenabfrage von der API mit dem Token in der Anfrage funktioniert, aber ich bevorzuge es, den Zugriffstoken über den Header zu übermitteln, weshalb wir im Skript api_caller.py genau diese Methode verwenden werden.

Beispielantwort:

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

Im obigen API-Antwortbeispiel wird ein Ticket von Sankt Petersburg nach Phuket gezeigt… Ach, warum träumen…
Da ich aus Kasan komme und Phuket gerade nur ein ‚Traum‘ für uns ist, suchen wir nach Tickets von Sankt Petersburg nach Kasan.

Es wird angenommen, dass Sie bereits ein AWS-Konto haben. Ich möchte besonders darauf hinweisen, dass Kinesis und das Versenden von Benachrichtigungen über SMS nicht im jährlichen Free Tier (kostenlose Nutzung). Aber selbst mit ein paar Dollar im Kopf kann man das vorgeschlagene System aufbauen und damit experimentieren. Und natürlich sollten Sie nicht vergessen, alle Ressourcen zu löschen, nachdem sie nicht mehr benötigt werden.

Glücklicherweise sind DynamoDB und Lambda-Funktionen für uns bedingt kostenlos, wenn wir innerhalb der monatlichen Freigrenzen bleiben. Zum Beispiel für DynamoDB: 25 GB Speicher, 25 WCU/RCU und 100 Millionen Anfragen. Und eine Million Lambda-Funktionsaufrufe pro Monat.

Manuelle Bereitstellung des Systems

Einrichtung von Kinesis Data Streams

Gehen wir zu dem Kinesis Data Streams-Dienst und erstellen zwei neue Streams, jeweils mit einem Shard.

Was ist ein Shard?
Ein Shard ist die grundlegende Einheit für die Datenübertragung im Amazon Kinesis Stream. Ein Shard ermöglicht die Eingabe von Daten mit einer Geschwindigkeit von 1 MB/s und die Ausgabe mit 2 MB/s. Ein Shard unterstützt bis zu 1000 PUT-Anfragen pro Sekunde. Bei der Erstellung eines Datenstreams muss die erforderliche Anzahl an Shards angegeben werden. Beispielsweise kann man einen Datenstream mit zwei Shards erstellen. Dieser Stream ermöglicht die Eingabe von Daten mit einer Geschwindigkeit von 2 MB/s und die Ausgabe mit 4 MB/s, wodurch bis zu 2000 PUT-Anfragen pro Sekunde unterstützt werden.

Je mehr Shards Ihr Stream hat, desto größer ist seine Durchsatzkapazität. Im Grunde genommen werden Streams durch das Hinzufügen von Shards skaliert. Allerdings steigt mit der Anzahl der Shards auch der Preis. Jeder Shard kostet 1,5 Cent pro Stunde und zusätzlich 1,4 Cent für jede Million PUT-Anfragen (PUT payload units).

Lassen Sie uns einen neuen Stream mit dem Namen airline_ticketserstellen. Ein Shard reicht dafür aus:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Jetzt erstellen wir einen weiteren Stream mit dem Namen special_stream:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Producer-Konfiguration

Als Datenproduzent für die Aufgabenbearbeitung reicht ein gewöhnlicher EC2-Instanz aus. Es muss keine leistungsstarke, teure virtuelle Maschine sein; ein Spot-Preis t2.micro ist völlig ausreichend.

Wichtiger Hinweis: Für das Beispiel sollte das Image – Amazon Linux AMI 2018.03.0 verwendet werden, da es weniger Konfigurationen für einen schnellen Start des Kinesis Agent benötigt.

Gehen Sie zum EC2-Service, erstellen Sie eine neue virtuelle Maschine, wählen Sie das gewünschte AMI mit dem Typ t2.micro, das im Free Tier enthalten ist:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Damit die neu erstellte virtuelle Maschine mit dem Kinesis-Service interagieren kann, müssen ihr die entsprechenden Berechtigungen erteilt werden. Der beste Weg, dies zu tun, ist, eine IAM-Rolle zuzuweisen. Daher sollten Sie auf dem Bildschirm Schritt 3: Instanzdetails konfigurieren wählen Neue IAM-Rolle erstellen:

Erstellen einer IAM-Rolle für EC2
Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Im geöffneten Fenster wählen wir aus, dass die neue Rolle für EC2 erstellt wird, und gehen zum Abschnitt Berechtigungen:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Im Übungsbeispiel müssen wir nicht auf alle Feinheiten der granularen Berechtigungszuweisungen für Ressourcen eingehen, daher wählen wir die von Amazon vordefinierten Policen: AmazonKinesisFullAccess und CloudWatchFullAccess.

Wir geben dieser Rolle einen sinnvollen Namen, zum Beispiel: EC2-KinesisStreams-FullAccess. Das Ergebnis sollte dasselbe sein wie im Bild unten angegeben:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Nachdem Sie diese neue Rolle erstellt haben, vergessen Sie nicht, sie dem zu erstellenden virtuellen Maschinen-Instance zuzuordnen:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
An diesem Bildschirm ändern wir nichts weiter und gehen zu den nächsten Fenstern über.

Die Festplattendefinitionen können auf den Standardwerten belassen werden, ebenso die Tags (obwohl es gute Praxis ist, Tags zu verwenden, um dem Instance einen Namen zu geben und die Umgebung zu kennzeichnen).

Jetzt sind wir auf der Registerkarte Schritt 6: Sicherheitsgruppe konfigurieren, wo Sie eine neue Sicherheitsgruppe erstellen oder eine bestehende Sicherheitsgruppe angeben müssen, die Verbindungen über SSH (Port 22) zur Instance erlaubt. Wählen Sie dort Quelle -> Meine IP aus, und Sie können die Instance starten.

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Sobald sie den Status 'running' erreicht hat, können Sie versuchen, sich über SSH mit ihr zu verbinden.

Um mit dem Kinesis-Agent arbeiten zu können, müssen Sie nach erfolgreicher Verbindung zur Maschine die folgenden Befehle im Terminal eingeben:

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

Erstellen wir einen Ordner zum Speichern der API-Antworten:

sudo mkdir /var/log/airline_tickets

Vor dem Start des Agents muss seine Konfiguration eingerichtet werden:

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

Der Inhalt der Datei agent.json sollte wie folgt aussehen:

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

Wie aus der Konfigurationsdatei hervorgeht, wird der Agent in dem Verzeichnis /var/log/airline_tickets/ Dateien mit der Endung .log überwachen, diese parsen und an den Stream airline_tickets übermitteln.

Wir starten den Dienst neu und stellen sicher, dass er läuft:

sudo service aws-kinesis-agent restart

Laden wir nun das Python-Skript herunter, das die Daten vom API abruft:

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

Das Skript api_caller.py fragt Daten von Aviasales ab und speichert die erhaltene Antwort in dem Verzeichnis, das vom Kinesis-Agent gescannt wird. Die Implementierung dieses Skripts ist ziemlich standardisiert, es gibt eine Klasse TicketsApi, die es ermöglicht, das API asynchron abzufragen. In diese Klasse übergeben wir den Header mit dem Token und die Anfrageparameter:

class TicketsApi:
    """API Aufruft Klasse."""

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

    async def get_data(self, data):
        """Holt die Daten aus der API-Abfrage."""
        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('Antwortstatus %s: %s',
                            self.base_url, response.status)
                response_json = await response.json()
            except HTTPError as http_err:
                LOGGER.error('Ups! HTTP-Fehler ist aufgetreten: %s', str(http_err))
            except Exception as err:
                LOGGER.error('Ups! Ein Fehler ist aufgetreten: %s', str(err))
            return response_json


def prepare_request(api_token):
    """Gibt die Header und Abfrage für die API-Anforderung zurück."""
    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():
    """Führt den Code aus."""
    if len(sys.argv) != 2:
        print('Verwendung: 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 hat %s Artikel zurückgegeben', len(response['data']))
        try:
            count_rows = log_maker(response)
            LOGGER.info('%s Zeilen wurden in %s gespeichert',
                        count_rows,
                        TARGET_FILE)
        except Exception as e:
            LOGGER.error('Ups! Das Anfrageergebnis wurde nicht in die Datei gespeichert. %s',
                         str(e))
    else:
        LOGGER.error('Ups! API-Anforderung war nicht erfolgreich %s!', response)

Um die Richtigkeit der Einstellungen und die Funktionsfähigkeit des Agents zu testen, führen wir einen Testlauf des Skripts api_caller.py durch:

sudo ./api_caller.py TOKEN

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Und wir beobachten das Ergebnis in den Logs des Agents sowie im Monitoring-Bereich des Datenstreams airline_tickets:

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

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Wie zu sehen ist, funktioniert alles und der Kinesis Agent sendet erfolgreich Daten in den Stream. Jetzt richten wir den Consumer ein.

Konfiguration von Kinesis Data Analytics

Kommen wir zu dem zentralen Bestandteil des gesamten Systems — wir erstellen eine neue Anwendung in Kinesis Data Analytics mit dem Namen kinesis_analytics_airlines_app:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Kinesis Data Analytics ermöglicht die Durchführung von Echtzeitanalysen von Daten aus Kinesis Streams mithilfe von SQL. Es handelt sich um einen vollständig skalierbaren Service (im Gegensatz zu Kinesis Streams), der:

  1. es ermöglicht, neue Streams (Output Stream) basierend auf Abfragen der Quelldaten zu erstellen;
  2. einen Fehlerstream bereitstellt, der während der Ausführung der Anwendungen aufgetreten ist (Error Stream);
  3. die Eingabedaten automatisch erkennt (diese kann bei Bedarf manuell überschrieben werden).

Dieser Service ist nicht günstig — 0,11 USD pro Stunde, daher sollte er vorsichtig verwendet und nach Abschluss der Arbeit gelöscht werden.

Wir verbinden die Anwendung mit der Datenquelle:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Wählen Sie den Stream aus, mit dem Sie sich verbinden möchten (airline_tickets):

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Anschließend müssen Sie eine neue IAM-Rolle anfügen, damit die Anwendung aus dem Stream lesen und in den Stream schreiben kann. Dazu müssen Sie im Abschnitt Zugriffsberechtigungen nichts ändern:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Jetzt fordern wir die Entdeckung des Datenschemas im Stream an, indem wir auf die Schaltfläche „Schema entdecken“ klicken. Dadurch wird eine neue IAM-Rolle erstellt und die Entdeckung des Schemas aus den bereits eingegangenen Daten gestartet:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Jetzt ist es notwendig, in den SQL-Editor zu wechseln. Wenn Sie auf diese Schaltfläche klicken, erscheint ein Fenster mit der Frage zum Starten der Anwendung – wählen Sie aus, was Sie starten möchten:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Im SQL-Editor fügen wir die folgende einfache Abfrage ein und klicken auf Speichern und SQL ausführen:

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';

In relationalen Datenbanken arbeiten Sie mit Tabellen und verwenden die INSERT-Anweisung, um Datensätze hinzuzufügen, sowie die SELECT-Anweisung, um Daten abzufragen. In Amazon Kinesis Data Analytics arbeiten Sie mit Streams (STREAM) und Pumpen (PUMP) — kontinuierlichen Einfügeanfragen, die Daten von einem Stream in eine andere Anwendung übertragen.

In der oben dargestellten SQL-Anfrage wird nach Aeroflot-Tickets gesucht, die weniger als fünftausend Rubel kosten. Alle Datensätze, die unter diese Bedingungen fallen, werden in den Stream DESTINATION_SQL_STREAM eingefügt.

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Im Abschnitt Destination wählen wir den Stream special_stream aus und im Dropdown-Menü In-application stream name wählen wir DESTINATION_SQL_STREAM:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Nach all diesen Manipulationen sollte etwas Ähnliches wie das Bild unten entstehen:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Erstellung und Abonnierung eines SNS-Themen

Wir gehen zum Service Simple Notification Service und erstellen dort ein neues Thema mit dem Namen Airlines:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Wir erstellen ein Abonnement für dieses Thema und geben dabei die Handy-Nummer an, an die die SMS-Benachrichtigungen gesendet werden sollen:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Erstellung einer Tabelle in DynamoDB

Um die unbearbeiteten Daten des Streams airline_tickets zu speichern, erstellen wir eine Tabelle in DynamoDB mit dem gleichen Namen. Als Primärschlüssel verwenden wir record_id:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Erstellung der Lambda-Funktion collector

Wir werden eine Lambda-Funktion mit dem Namen Collector erstellen, deren Aufgabe es ist, den Stream airline_tickets abzufragen und, falls neue Einträge vorhanden sind, diese in die DynamoDB-Tabelle einzufügen. Es ist offensichtlich, dass diese Lambda-Funktion neben den Standardberechtigungen auch Zugriff auf das Lesen des Kinesis-Datenstroms und das Schreiben in die DynamoDB haben muss.

Erstellung einer IAM-Rolle für die Lambda-Funktion Collector
Zunächst erstellen wir eine neue IAM-Rolle für die Lambda, die den Namen Lambda-TicketsProcessingRole trägt:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Für dieses Testbeispiel sind die vorkonfigurierten Richtlinien AmazonKinesisReadOnlyAccess und AmazonDynamoDBFullAccess ausreichend, wie im Bild unten dargestellt:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Diese Lambda muss durch einen Trigger von Kinesis gestartet werden, wenn neue Einträge in den Stream airline_stream gelangen, daher müssen wir einen neuen Trigger hinzufügen:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Jetzt müssen wir den Code einfügen und die Lambda-Funktion speichern.

"""Den Stream analysieren und in die DynamoDB-Tabelle einfügen."""
import base64
import json
import boto3
from decimal import Decimal

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

class TicketsParser:
    """Informationen aus dem Stream analysieren."""

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

    @staticmethod
    def get_json_data(records):
        """Gibt die deserialisierten Daten aus dem Stream zurück."""
        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):
        """Vorverarbeitung der JSON-Daten."""
        new_item = {
            'record_id': json_item.get('record_id'),
            'cost': Decimal(json_item.get('cost')),
            'trip_class': json_item.get('trip_class'),
            'show_to_affiliates': json_item.get('show_to_affiliates'),
            'origin': json_item.get('origin'),
            'number_of_changes': int(json_item.get('number_of_changes')),
            'gate': json_item.get('gate'),
            'found_at': json_item.get('found_at'),
            'duration': int(json_item.get('duration')),
            'distance': int(json_item.get('distance')),
            'destination': json_item.get('destination'),
            'depart_date': json_item.get('depart_date'),
            'actual': json_item.get('actual')
        }
        return new_item

    def run(self):
        """Batch-Einfügen in die Tabelle."""
        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('Es wurden ', len(self.json_data), 'Einträge hinzugefügt')

def lambda_handler(event, context):
    """Den Stream analysieren und in die DynamoDB-Tabelle einfügen."""
    print('Event erhalten:', event)
    parser = TicketsParser(TABLE_NAME, event['Records'])
    parser.run()

Erstellung der Lambda-Funktion Notifier

Die zweite Lambda-Funktion, die den zweiten Stream (special_stream) überwachen und Benachrichtigungen an SNS senden wird, wird auf ähnliche Weise erstellt. Daher muss diese Lambda Funktion Lesezugriff auf Kinesis und die Berechtigung zum Senden von Nachrichten an das festgelegte SNS-Topic haben, das die Benachrichtigungen dann an alle Abonnenten dieses Topics (E-Mail, SMS usw.) weiterleitet.

Erstellung einer IAM-Rolle
Zunächst erstellen wir die IAM-Rolle Lambda-KinesisAlarm für diese Lambda-Funktion und weisen dann diese Rolle der zu erstellenden Lambda-Funktion alarm_notifier zu:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Diese Lambda-Funktion soll durch einen Trigger aktiviert werden, wenn neue Einträge im Stream special_stream ankommen. Daher muss der Trigger ähnlich eingerichtet werden, wie wir es für die Lambda-Funktion Collector gemacht haben.

Zur Vereinfachung der Konfiguration dieser Lambda-Funktion fügen wir eine neue Umgebungsvariable hinzu – TOPIC_ARN, in der wir die Amazon Resource Names (ARN) des Topics Airlines speichern:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Und fügen den Lambda-Code ein, der ganz einfach ist:

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='Hallo! Ich habe etwas Interessantes gefunden!',
                           Subject='Alarm für Flugtickets')
        print('Alarmnachricht wurde erfolgreich gesendet')
    except Exception as err:
        print('Zustellfehler', str(err))

Es scheint, dass die manuelle Systemkonfiguration hier abgeschlossen ist. Jetzt müssen wir nur noch testen und sicherstellen, dass alles richtig eingerichtet wurde.

Deployment aus dem Terraform-Code

Notwendige Vorbereitung

Terraform — ein sehr praktisches Open-Source-Tool zum Bereitstellen von Infrastruktur aus Code. Es hat eine eigene Syntax, die leicht zu erlernen ist, und viele Beispiele dafür, wie und was bereitgestellt werden kann. In den Editoren Atom oder Visual Studio Code gibt es zahlreiche nützliche Plugins, die die Arbeit mit Terraform erleichtern.

Das Distributionspaket kann heruntergeladen werden hier herunter. Eine ausführliche Analyse aller Funktionen von Terraform sprengt den Rahmen dieses Artikels, daher beschränken wir uns auf die wesentlichen Punkte.

Wie man es startet

Der vollständige Code des Projekts befindet sich in meinem Repository. Klonen Sie das Repository auf Ihr System. Vor dem Start sollten Sie sicherstellen, dass das AWS CLI installiert und konfiguriert ist, da Terraform die Anmeldeinformationen in der Datei ~/.aws/credentials sucht.

Eine gute Praxis ist es, vor dem Deployment der gesamten Infrastruktur den Befehl plan auszuführen, um zu sehen, was Terraform aktuell in der Cloud erstellen wird:

terraform.exe plan

Es wird verlangt, eine Telefonnummer einzugeben, um Benachrichtigungen zu senden. Zu diesem Zeitpunkt ist die Eingabe jedoch nicht erforderlich.

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Nach der Analyse des Arbeitsplans des Programms können wir die Ressourcenerstellung starten:

terraform.exe apply

Nach dem Absenden dieses Befehls wird erneut nach der Telefonnummer gefragt. Geben Sie "yes" ein, wenn die Frage zur tatsächlichen Durchführung der Aktionen angezeigt wird. Damit wird die gesamte Infrastruktur hochgefahren und alle erforderlichen Konfigurationen für EC2 vorgenommen, Lambda-Funktionen bereitgestellt usw.

Sobald alle Ressourcen erfolgreich über den Terraform-Code erstellt wurden, müssen Sie die Details der Kinesis Analytics-Anwendung aufrufen (leider habe ich nicht herausgefunden, wie das direkt aus dem Code zu machen ist).

Starten Sie die Anwendung:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Danach müssen Sie den in-application Stream-Namen explizit angeben, indem Sie aus dem Dropdown-Menü auswählen:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Jetzt ist alles bereit zur Nutzung.

Testen der Anwendung

Unabhängig davon, ob Sie das System manuell oder über den Terraform-Code bereitgestellt haben, wird es gleich funktionieren.

Loggen Sie sich per SSH auf die EC2-virtuelle Maschine ein, auf der der Kinesis-Agent installiert ist, und führen Sie das Skript api_caller.py aus.

sudo ./api_caller.py TOKEN

Warten Sie auf die SMS an Ihre Nummer:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Die SMS kommt fast innerhalb von 1 Minute auf Ihr Telefon an:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless
Es bleibt zu überprüfen, ob die Datensätze in der DynamoDB-Datenbank für eine spätere, detailliertere Analyse gespeichert wurden. Die Tabelle airline_tickets enthält ungefähre Daten wie folgt:

Integration der Aviasales-API mit Amazon Kinesis und der Einfachheit von Serverless

Fazit

Im Rahmen der Arbeit wurde ein System zur Online-Datenverarbeitung auf Basis von Amazon Kinesis aufgebaut. Es wurden Optionen für den Einsatz des Kinesis Agents in Verbindung mit Kinesis Data Streams und der Echtzeitanalyse Kinesis Analytics unter Verwendung von SQL-Befehlen betrachtet, ebenso wie die Interaktion von Amazon Kinesis mit anderen AWS-Diensten.

Das oben beschriebene System haben wir auf zwei Weisen bereitgestellt: einmal langsam manuell und schnell über Terraform-Code.

Der gesamte Quellcode des Projekts ist verfügbar in meinem Repository auf GitHub, ich lade Sie ein, sich damit vertraut zu machen.

Ich freue mich darauf, den Artikel zu diskutieren und warte auf Ihre Kommentare. Ich hoffe auf konstruktive Kritik.

Viel Erfolg!

Quelle: habr.com

Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen 🔥 Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen | ProHoster