Integracja Aviasales API z Amazon Kinesis i prostota serverless

Cześć, Habr!

Czy lubisz latać samolotami? Uwielbiam, ale podczas izolacji polubiłem także analizowanie danych o biletach lotniczych z jednego znanego serwisu — Aviasales.

Dziś przeanalizujemy działanie Amazon Kinesis, zbudujemy system strumieniowy z analizą w czasie rzeczywistym, wdrożymy bazę danych NoSQL Amazon DynamoDB jako główne miejsce przechowywania danych i skonfigurujemy powiadomienia SMS o interesujących biletach.

Wszystkie szczegóły poniżej! Zaczynamy!

Integracja Aviasales API z Amazon Kinesis i prostota serverless

Wprowadzenie

Na przykład będziemy potrzebować dostępu do API Aviasales. Dostęp do niego jest darmowy i bez ograniczeń, wystarczy zarejestrować się w sekcji „Dla programistów”, aby uzyskać swój token API do dostępu do danych.

Głównym celem tego artykułu jest przedstawienie ogólnego zrozumienia korzystania z przesyłania strumieniowego informacji w AWS, pomijamy fakt, że dane zwracane przez używane API nie są ściśle aktualne i są przesyłane z pamięci podręcznej, która jest tworzona na podstawie wyszukiwań użytkowników serwisów Aviasales.ru i Jetradar.com w ciągu ostatnich 48 godzin.

Dane o biletach lotniczych uzyskane przez API Kinesis-agent, zainstalowany na maszynie-producenta, będą automatycznie parsowane i przesyłane do odpowiedniego strumienia za pośrednictwem Kinesis Data Analytics. Surowa wersja tego strumienia będzie zapisywana bezpośrednio w magazynie. Rozbudowane w DynamoDB magazyn „surowych” danych umożliwi przeprowadzenie głębszej analizy biletów za pomocą narzędzi BI, na przykład AWS Quick Sight.

Rozważymy dwie opcje wdrożenia całej infrastruktury:

  • Ręczna — przez AWS Management Console;
  • Infrastruktura z kodu Terraform — dla leniwych automatyków;

Architektura opracowywanego systemu

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Używane komponenty:

  • Aviasales API — dane zwracane przez to API będą używane do całej późniejszej pracy;
  • EC2 Producer Instance — zwykła maszyna wirtualna w chmurze, na której będzie generowany wejściowy strumień danych:
    • Kinesis Agent — to aplikacja Java instalowana lokalnie na maszynie, która zapewnia prosty sposób zbierania i wysyłania danych do Kinesis (Kinesis Data Streams lub Kinesis Firehose). Agent permanentnie monitoruje zestaw plików w określonych katalogach i wysyła nowe dane do Kinesis;
    • Skrypt API Caller — skrypt w Pythonie, który wykonuje zapytania do API i zapisuje odpowiedzi w folderze monitorowanym przez Kinesis Agent;
  • Kinesis Data Streams — usługa przesyłania danych w czasie rzeczywistym z szerokimi możliwościami skalowania;
  • Kinesis Analytics — usługa bezserwerowa, która upraszcza analizę danych strumieniowych w czasie rzeczywistym. Amazon Kinesis Data Analytics dostosowuje zasoby do działania aplikacji i automatycznie skaluje się w celu obsługi dowolnych ilości przychodzących danych;
  • AWS Lambda — usługa, która pozwala uruchamiać kod bez konieczności rezerwacji i konfiguracji serwerów. Cała moc obliczeniowa automatycznie skaluje się dla każdego wywołania;
  • Amazon DynamoDB — baza danych typu „klucz-wartość” i dokumentów, która zapewnia opóźnienie poniżej 10 milisekund przy pracy w każdym rozmiarze. Korzystając z DynamoDB, nie trzeba rozdzielać żadnych serwerów, instalować poprawek ani nimi zarządzać. DynamoDB automatycznie skaluje tabele, dostosowując ilość dostępnych zasobów i zapewniając wysoką wydajność. Żadne działania związane z administracją systemu nie są wymagane;
  • Amazon SNS — w pełni zarządzana usługa wysyłania wiadomości oparta na modelu „wydawca — subskrybent” (Pub/Sub), która umożliwia izolację mikrousług, rozproszonych systemów oraz aplikacji bezserwerowych. SNS może być używane do powiadamiania użytkowników końcowych za pomocą mobilnych powiadomień push, wiadomości SMS i e-maili.

Pierwsze przygotowanie

Aby zasymulować strumień danych, postanowiłem użyć informacji o biletach lotniczych zwracanych przez API Aviasales. W dokumentacji dość obszernym zestawie różnych metod, weźmy jedną z nich — „Kalendarz cen na miesiąc”, która zwraca ceny za każdy dzień miesiąca, pogrupowane według liczby przesiadek. Jeśli nie uwzględnimy w zapytaniu miesiąca poszukiwań, zwrócone zostaną informacje za miesiąc następujący po bieżącym.

Tak więc rejestrujemy się, otrzymujemy swój token.

Przykład zapytania poniżej:

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

Powyższa metoda uzyskiwania danych z API z uwzględnieniem tokena w zapytaniu będzie działać, ale bardziej podoba mi się przekazywanie tokena dostępu przez nagłówek, dlatego w skrypcie api_caller.py będziemy korzystać właśnie z tej metody.

Przykład odpowiedzi:

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

W powyższym przykładzie odpowiedzi API pokazano bilet z Sankt Petersburga do Phuket… Ech, po co marzyć…
Ponieważ jestem z Kazania, a Phuket teraz "tylko nam się śni", poszukajmy biletów z Petersburga do Kazania.

Zakłada się, że masz już konto w AWS. Od razu chcę zwrócić szczególną uwagę, że Kinesis i wysyłanie powiadomień przez SMS nie są objęte rocznym Free Tier (bezpłatne korzystanie). Ale nawet mimo to, mając na uwadze kilka dolarów, można w pełni zbudować zaproponowany system i poeksperymentować z nim. Oczywiście nie należy zapominać o usunięciu wszystkich zasobów po tym, jak przestaną być potrzebne.

Na szczęście DynamoDb i funkcje lambda będą dla nas warunkowo bezpłatne, jeśli zmieszczą się w miesięcznych limitach bezpłatnych. Na przykład dla DynamoDB: 25 GB pamięci, 25 WCU/RCU i 100 mln zapytań. I milion wywołań funkcji lambda miesięcznie.

Ręczne wdrożenie systemu

Konfiguracja Kinesis Data Streams

Przechodzimy do usługi Kinesis Data Streams i tworzymy dwa nowe strumienie z jednym shardem każdy.

Czym jest shard?
Shard to podstawowa jednostka przesyłania danych strumienia Amazon Kinesis. Jeden segment zapewnia przesyłanie danych wejściowych z prędkością 1 MB/s i przesyłanie danych wyjściowych z prędkością 2 MB/s. Jeden segment obsługuje do 1000 zapisów PUT na sekundę. Przy tworzeniu strumienia danych należy określić wymaganą liczbę segmentów. Na przykład można stworzyć strumień danych z dwoma segmentami. Ten strumień danych zapewni przesyłanie danych wejściowych z prędkością 2 MB/s i przesyłanie danych wyjściowych z prędkością 4 MB/s, obsługując do 2000 zapisów PUT na sekundę.

Im więcej shardów w Twoim strumieniu, tym większa jego przepustowość. W zasadzie tak się skalują strumienie — poprzez dodawanie shardów. Ale im więcej masz shardów, tym wyższa cena. Każdy shard kosztuje 1,5 centa za godzinę i dodatkowo 1,4 centa za każdego miliona operacji dodawania do strumienia (jednostki obciążenia PUT).

Stwórzmy nowy strumień o nazwie airline_tickets, wystarczy mu 1 shard:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Teraz stwórzmy kolejny strumień o nazwie special_stream:

Integracja Aviasales API z Amazon Kinesis i prostota serverless

Konfiguracja producenta

Jako producent danych do analizy zadania wystarczy użyć zwykłego instancji EC2. Nie musi to być mocna, droga maszyna wirtualna, wystarczy spotowy t2.micro.

Ważna uwaga: w przykładzie należy użyć obrazu — Amazon Linux AMI 2018.03.0, z nim jest mniej ustawień do szybkiego uruchomienia Kinesis Agent.

Przechodzimy do usługi EC2, tworzymy nową maszynę wirtualną, wybieramy odpowiedni AMI typu t2.micro, który wchodzi w skład Free Tier:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Aby nowo utworzona maszyna wirtualna mogła współpracować z usługą Kinesis, należy jej nadać odpowiednie uprawnienia. Najlepszym sposobem, aby to zrobić, jest przypisanie Roli IAM. Dlatego na ekranie Krok 3: Konfiguracja szczegółów instancji należy wybrać Utwórz nową Rolę IAM:

Tworzenie Roli IAM dla EC2
Integracja Aviasales API z Amazon Kinesis i prostota serverless
W otwartym oknie wybieramy, że nową rolę tworzymy dla EC2 i przechodzimy do sekcji Uprawnienia:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
W przypadku przykładu edukacyjnego nie musimy wnikać w wszystkie szczegóły granularnej konfiguracji uprawnień do zasobów, dlatego wybierzemy wstępnie skonfigurowane polityki Amazon: AmazonKinesisFullAccess i CloudWatchFullAccess.

Nadamy tej roli jakieś sensowne imię, na przykład: EC2-KinesisStreams-FullAccess. Ostatecznie powinno to wyglądać tak samo jak na poniższym obrazku:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Po utworzeniu tej nowej roli nie zapomnijmy przypiąć jej do tworzonej instancji maszyny wirtualnej:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Na tym ekranie nic więcej nie zmieniamy i przechodzimy do następnych okien.

Ustawienia dysku twardego możemy pozostawić domyślne, tagi również (choć dobrą praktyką jest używać tagów, przynajmniej nadając instancji nazwę i wskazując środowisko).

Teraz jesteśmy na zakładce Krok 6: Konfiguracja grupy zabezpieczeń, gdzie należy utworzyć nową lub wskazać istniejącą grupę zabezpieczeń, pozwalającą na połączenie przez ssh (port 22) z instancją. Wybierz tam Źródło —> Mój IP i możesz uruchomić instancję.

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Gdy tylko przejdzie w status uruchomione, można spróbować połączyć się z nią przez ssh.

Aby uzyskać możliwość pracy z Kinesis Agent, po pomyślnym połączeniu z maszyną, należy wpisać następujące polecenia w terminalu:

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

Utwórzmy folder na odpowiedzi API:

sudo mkdir /var/log/airline_tickets

Przed uruchomieniem agenta należy skonfigurować jego plik konfiguracyjny:

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

Zawartość pliku agent.json powinna mieć następujący wygląd:

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

Jak widać z pliku konfiguracyjnego, agent będzie monitorował w katalogu /var/log/airline_tickets/ pliki z rozszerzeniem .log, analizował je i przekazywał do strumienia airline_tickets.

Restartujemy usługę i upewniamy się, że została uruchomiona i działa:

sudo service aws-kinesis-agent restart

Teraz pobierzemy skrypt Python, który będzie pobierał dane z 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

Skrypt api_caller.py pobiera dane z Aviasales i zapisuje otrzymaną odpowiedź w katalogu, który skanuje agent Kinesis. Wdrożenie tego skryptu jest dość standardowe, istnieje klasa TicketsApi, która pozwala asynchronicznie wywoływać API. Do tej klasy przekazujemy nagłówek z tokenem i parametry zapytania:

class TicketsApi:
    """Klasa wywołująca API."""

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

    async def get_data(self, data):
        """Pobiera dane z zapytania 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('Status odpowiedzi %s: %s',
                            self.base_url, response.status)
                response_json = await response.json()
            except HTTPError as http_err:
                LOGGER.error('Ups! Wystąpił błąd HTTP: %s', str(http_err))
            except Exception as err:
                LOGGER.error('Ups! Wystąpił błąd: %s', str(err))
            return response_json


def prepare_request(api_token):
    """Zwraca nagłówki i zapytanie dla żądania 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():
    """Uruchom kod."""
    if len(sys.argv) != 2:
        print('Użycie: 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 zwróciło %s elementów', len(response['data']))
        try:
            count_rows = log_maker(response)
            LOGGER.info('%s wierszy zapisano do %s',
                        count_rows,
                        TARGET_FILE)
        except Exception as e:
            LOGGER.error('Ups! Wynik zapytania nie został zapisany do pliku. %s',
                         str(e))
    else:
        LOGGER.error('Ups! Żądanie API nie powiodło się %s!', response)

Aby przetestować poprawność ustawień i operacyjność agenta, wykonamy testowe uruchomienie skryptu api_caller.py:

sudo ./api_caller.py TOKEN

Integracja Aviasales API z Amazon Kinesis i prostota serverless
I sprawdzamy wyniki pracy w logach Agenta oraz na zakładce Monitoring w strumieniu danych airline_tickets:

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

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Integracja Aviasales API z Amazon Kinesis i prostota serverless
Jak widać, wszystko działa i Kinesis Agent pomyślnie przesyła dane do strumienia. Teraz skonfigurujemy konsumenta.

Konfiguracja Kinesis Data Analytics

Przejdźmy do centralnego komponentu całego systemu — stworzymy nową aplikację w Kinesis Data Analytics o nazwie kinesis_analytics_airlines_app:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Kinesis Data Analytics umożliwia analizowanie danych w czasie rzeczywistym z Kinesis Streams za pomocą języka SQL. To całkowicie automatycznie skalowalna usługa (w odróżnieniu od Kinesis Streams), która:

  1. pozwala tworzyć nowe strumienie (Output Stream) na podstawie zapytań do danych źródłowych;
  2. oferuje strumień z błędami, które wystąpiły podczas działania aplikacji (Error Stream);
  3. potrafi automatycznie rozpoznać schemat danych wejściowych (można go ręcznie nadpisać w razie potrzeby).

To nie jest tani serwis — 0.11 USD za godzinę pracy, dlatego należy go używać ostrożnie i usuwać po zakończeniu pracy.

Podłączmy aplikację do źródła danych:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Wybieramy strumień, do którego zamierzamy się podłączyć (airline_tickets):

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Następnie należy dołączyć nową rolę IAM, aby aplikacja mogła czytać ze strumienia i zapisywać do strumienia. W tym celu wystarczy nic nie zmieniać w sekcji Uprawnienia dostępu:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Teraz zażądajmy odkrycia schematu danych w strumieniu, w tym celu klikamy przycisk „Discover schema”. W rezultacie zaktualizuje się (zostanie utworzona nowa) rola IAM i rozpocznie się odkrywanie schematu z danych, które już napłynęły do strumienia:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Teraz musimy przejść do edytora SQL. Po naciśnięciu tego przycisku otworzy się okno z pytaniem o uruchomienie aplikacji — wybieramy, co chcemy uruchomić:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Do okna edytora SQL wstawiamy takie proste zapytanie i klikamy Zapisz i uruchom 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';

W relacyjnych bazach danych pracujesz z tabelami, używając operatorów INSERT do dodawania zapisów oraz operatora SELECT do zapytania danych. W Amazon Kinesis Data Analytics pracujesz z strumieniami (STREAM) i „pompami” (PUMP) — ciągłymi zapytaniami o wstawianie, które wstawiają dane z jednego strumienia w aplikacji do innego strumienia.

W przedstawionym powyżej zapytaniu SQL następuje wyszukiwanie biletów Aeroflot po cenie poniżej pięciu tysięcy rubli. Wszystkie zapisy, które spełniają te warunki, zostaną umieszczone w strumieniu DESTINATION_SQL_STREAM.

Integracja Aviasales API z Amazon Kinesis i prostota serverless
W bloku Destination wybieramy strumień special_stream, a w rozwijanej liście In-application stream name DESTINATION_SQL_STREAM:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
W wyniku wszystkich manipulacji powinno powstać coś podobnego do poniższego obrazu:

Integracja Aviasales API z Amazon Kinesis i prostota serverless

Tworzenie i subskrypcja tematu SNS

Przechodzimy do usługi Simple Notification Service i tam tworzymy nowy temat o nazwie Airlines:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Zatwierdzamy subskrypcję do tego tematu, w której podajemy numer telefonu komórkowego, na który będą przychodziły powiadomienia SMS:

Integracja Aviasales API z Amazon Kinesis i prostota serverless

Tworzenie tabeli w DynamoDB

Aby przechowywać surowe dane z strumienia airline_tickets, stworzymy tabelę w DynamoDB o tej samej nazwie. Jako klucz główny użyjemy record_id:

Integracja Aviasales API z Amazon Kinesis i prostota serverless

Tworzenie funkcji lambda collector

Stworzymy funkcję lambda o nazwie Collector, której zadaniem będzie monitorowanie strumienia airline_tickets i, w przypadku znalezienia nowych rekordów, wstawienie tych rekordów do tabeli DynamoDB. Oczywiście oprócz domyślnych uprawnień, ta lambda musi mieć dostęp do odczytu strumienia danych Kinesis oraz zapisu do DynamoDB.

Tworzenie roli IAM dla funkcji lambda collector
Na początek stwórzmy nową rolę IAM dla lambdy o nazwie Lambda-TicketsProcessingRole:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Dla przykładu testowego odpowiednie będą wstępnie skonfigurowane polityki AmazonKinesisReadOnlyAccess oraz AmazonDynamoDBFullAccess, jak pokazano na poniższym obrazie:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Integracja Aviasales API z Amazon Kinesis i prostota serverless

Ta lambda ma być uruchamiana przez wyzwalacz z Kinesis po pojawieniu się nowych rekordów w strumieniu airline_stream, więc trzeba dodać nowy wyzwalacz:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Integracja Aviasales API z Amazon Kinesis i prostota serverless
Pozostało wstawić kod i zapisać lambdę.

"""Analizowanie strumienia i wstawianie do tabeli DynamoDB."""
import base64
import json
import boto3
from decimal import Decimal

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

class TicketsParser:
    """Analiza informacji ze strumienia."""

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

    @staticmethod
    def get_json_data(records):
        """Zwraca zdeserializowane dane ze strumienia."""
        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):
        """Wstępne przetwarzanie danych 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):
        """Wstawianie wsadowe do tabeli."""
        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('Dodano ', len(self.json_data), 'elementów')

def lambda_handler(event, context):
    """Analizowanie strumienia i wstawianie do tabeli DynamoDB."""
    print('Otrzymano zdarzenie:', event)
    parser = TicketsParser(TABLE_NAME, event['Records'])
    parser.run()

Tworzenie funkcji lambda notifier

Druga funkcja lambda, która będzie monitorować drugi strumień (special_stream) i wysyłać powiadomienia do SNS, jest tworzona w podobny sposób. W związku z tym, ta lambda musi mieć dostęp do odczytu z Kinesis i wysyłania wiadomości do określonego tematu SNS, który następnie zostanie przesłany przez serwis SNS do wszystkich subskrybentów tego tematu (e-mail, SMS itp.).

Tworzenie roli IAM
Najpierw tworzymy rolę IAM Lambda-KinesisAlarm dla tej lambdy, a następnie przypisujemy tę rolę do tworzonej lambdy alarm_notifier:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Integracja Aviasales API z Amazon Kinesis i prostota serverless

Ta lambda powinna działać w odpowiedzi na nowe wpisy w strumieniu special_stream, więc musimy skonfigurować wyzwalacz w podobny sposób, jak to zrobiliśmy dla lambdy Collector.

Dla ułatwienia konfiguracji tej lambdy wprowadzimy nową zmienną środowiskową — TOPIC_ARN, gdzie umieścimy ANR (Amazon Resource Names) tematu Airlines:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
I wstawiamy kod lambdy, który nie jest skomplikowany:

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='Cześć! Znalazłem coś interesującego!',
                           Subject='Alarm biletów lotniczych')
        print('Wiadomość alarmowa została pomyślnie dostarczona')
    except Exception as err:
        print('Błąd dostarczenia', str(err))

Wygląda na to, że ręczna konfiguracja systemu jest zakończona. Pozostaje tylko przetestować i upewnić się, że wszystko zostało poprawnie skonfigurowane.

Wdrożenie z kodu Terraform

Wymagana konfiguracja

Terraform — bardzo wygodne narzędzie open-source do wdrażania infrastruktury z kodu. Ma własną składnię, którą łatwo opanować, oraz wiele przykładów, jak i co wdrożyć. W edytorze Atom lub Visual Studio Code jest wiele przydatnych wtyczek, które ułatwiają pracę z Terraform.

Dostępny jest do pobrania dystrybucja stąd. Szczegółowa analiza wszystkich możliwości Terraform wykracza poza ramy tego artykułu, dlatego ograniczymy się do najważniejszych kwestii.

Jak uruchomić

Pełny kod projektu znajduje się w moim repozytorium. Klonujemy repozytorium do siebie. Przed uruchomieniem należy upewnić się, że masz zainstalowane i skonfigurowane AWS CLI, ponieważ Terraform będzie szukał danych logowania w pliku ~/ .aws /credentials.

Dobrą praktyką przed wdrożeniem całej infrastruktury jest uruchomienie polecenia plan, aby zobaczyć, co Terraform teraz utworzy w chmurze:

terraform.exe plan

Zostanie poproszony o podanie numeru telefonu do wysyłania powiadomień. Na tym etapie wprowadzanie go nie jest konieczne.

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Analizując plan działania programu, możemy rozpocząć tworzenie zasobów:

terraform.exe apply

Po wysłaniu tego polecenia ponownie pojawi się prośba o podanie numeru telefonu, wprowadzamy „tak”, gdy zostanie zadane pytanie o rzeczywiste wykonanie działań. To pozwoli uruchomić całą infrastrukturę, przeprowadzić niezbędną konfigurację EC2, wdrożyć funkcje lambda itp.

Po tym, jak wszystkie zasoby zostaną pomyślnie utworzone za pomocą kodu Terraform, należy przejść do szczegółów aplikacji Kinesis Analytics (niestety, nie znalazłem, jak to zrobić bezpośrednio z kodu).

Uruchamiamy aplikację:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Następnie należy jawnie określić nazwę strumienia wewnątrz aplikacji, wybierając z rozwijanego menu:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Integracja Aviasales API z Amazon Kinesis i prostota serverless
Teraz wszystko jest gotowe do pracy.

Testowanie działania aplikacji

Bez względu na to, w jaki sposób wdrożyłeś system, ręcznie czy za pomocą kodu Terraform, będzie działać tak samo.

Logujemy się przez SSH na maszynę wirtualną EC2, na której zainstalowany jest Kinesis Agent i uruchamiamy skrypt api_caller.py

sudo ./api_caller.py TOKEN

Teraz musimy poczekać na SMS-a na Twój numer:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
SMS - wiadomość przychodzi na telefon praktycznie w ciągu 1 minuty:

Integracja Aviasales API z Amazon Kinesis i prostota serverless
Trzeba jeszcze sprawdzić, czy zapisy w bazie danych DynamoDB zostały zapisane do dalszej, bardziej szczegółowej analizy. Tabela airline_tickets zawiera mniej więcej takie dane:

Integracja Aviasales API z Amazon Kinesis i prostota serverless

Podsumowanie

W trakcie wykonanej pracy zbudowano system przetwarzania danych w czasie rzeczywistym oparty na Amazon Kinesis. Rozważono różne opcje wykorzystania Kinesis Agent w połączeniu z Kinesis Data Streams oraz analizą w czasie rzeczywistym za pomocą Kinesis Analytics przy użyciu komend SQL, a także interakcję Amazon Kinesis z innymi usługami AWS.

Opisaną powyżej system uruchomiliśmy w dwóch sposobach: dość długim ręcznym i szybkim za pomocą kodu Terraform.

Cały kod źródłowy projektu jest dostępny w moim repozytorium na GitHubie, zachęcam do zapoznania się z nim.

Chętnie przedyskutuję artykuł, czekam na Twoje komentarze. Liczę na konstruktywną krytykę.

Życzę powodzenia!

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster