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!

Wprowadzenie
Na przykład będziemy potrzebować dostępu do . 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

Używane komponenty:
- — dane zwracane przez to API będą używane do całej późniejszej pracy;
- — zwykła maszyna wirtualna w chmurze, na której będzie generowany wejściowy strumień danych:
- — 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 w Pythonie, który wykonuje zapytania do API i zapisuje odpowiedzi w folderze monitorowanym przez Kinesis Agent;
- — usługa przesyłania danych w czasie rzeczywistym z szerokimi możliwościami skalowania;
- — 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;
- — 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;
- — 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;
- — 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 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_APIPowyż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 . 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:

Teraz stwórzmy kolejny strumień o nazwie special_stream:

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:

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

W otwartym oknie wybieramy, że nową rolę tworzymy dla EC2 i przechodzimy do sekcji Uprawnienia:

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:

Po utworzeniu tej nowej roli nie zapomnijmy przypiąć jej do tworzonej instancji maszyny wirtualnej:

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ę.

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_ticketsPrzed uruchomieniem agenta należy skonfigurować jego plik konfiguracyjny:
sudo vim /etc/aws-kinesis/agent.jsonZawartość 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 restartTeraz 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 
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 

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:

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:
- pozwala tworzyć nowe strumienie (Output Stream) na podstawie zapytań do danych źródłowych;
- oferuje strumień z błędami, które wystąpiły podczas działania aplikacji (Error Stream);
- 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:

Wybieramy strumień, do którego zamierzamy się podłączyć (airline_tickets):

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:

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:

Teraz musimy przejść do edytora SQL. Po naciśnięciu tego przycisku otworzy się okno z pytaniem o uruchomienie aplikacji — wybieramy, co chcemy uruchomić:

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.

W bloku Destination wybieramy strumień special_stream, a w rozwijanej liście In-application stream name DESTINATION_SQL_STREAM:

W wyniku wszystkich manipulacji powinno powstać coś podobnego do poniższego obrazu:

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

Zatwierdzamy subskrypcję do tego tematu, w której podajemy numer telefonu komórkowego, na który będą przychodziły powiadomienia SMS:

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:

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:

Dla przykładu testowego odpowiednie będą wstępnie skonfigurowane polityki AmazonKinesisReadOnlyAccess oraz AmazonDynamoDBFullAccess, jak pokazano na poniższym obrazie:


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:


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:


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:

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
— 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 . 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ę . 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 planZostanie poproszony o podanie numeru telefonu do wysyłania powiadomień. Na tym etapie wprowadzanie go nie jest konieczne.

Analizując plan działania programu, możemy rozpocząć tworzenie zasobów:
terraform.exe applyPo 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ę:

Następnie należy jawnie określić nazwę strumienia wewnątrz aplikacji, wybierając z rozwijanego menu:


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 TOKENTeraz musimy poczekać na SMS-a na Twój numer:

SMS - wiadomość przychodzi na telefon praktycznie w ciągu 1 minuty:

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:

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 , 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
