Здравейте, Хабр!
Обичате ли да летите със самолети? Аз обожавам, но по време на самоизолацията започнах и да анализирам данни за самолетни билети от един известен ресурс - Aviasales.
Днес ще разгледаме работата на Amazon Kinesis, ще построим стрийминг система с анализ в реално време, ще инсталираме NoSQL база данни Amazon DynamoDB като основно хранилище и ще настроим известия чрез SMS за интересни билети.
Всички подробности по-долу! Да започнем!

Въведение
За примера ще ни е необходим достъп до . Достъпът до него е безплатен и неограничен, необходимо е само да се регистрирате в секцията „Разработчици“, за да получите своя API токен за достъп до данните.
Основната цел на тази статия е да предостави общо разбиране за използването на потоково предаване на информация в AWS, като изключваме факта, че данните, предоставени от използвания API, не са строго актуални и се предават от кеш, който се образува на базата на търсенията на потребителите на сайтовете Aviasales.ru и Jetradar.com през последните 48 часа.
Данните за самолетни билети, получени чрез API, Kinesis-agent, инсталиран на машината-производител, автоматично ще се парсват и предават в необходимия поток чрез Kinesis Data Analytics. Необработената версия на този поток ще се записва директно в хранилището. Разгърнатото в DynamoDB хранилище на „сурови“ данни ще позволи извършването на по-дълбок анализ на билетите чрез BI инструменти, например AWS Quick Sight.
Ще разгледаме два варианта за деплой на цялата инфраструктура:
- Ръчен — чрез AWS Management Console;
- Инфраструктура на код Terraform — за ленивите автоматизатори;
Архитектура на разработваната система

Използвани компоненти:
- — данните, предоставени от този API, ще се използват за цялата последваща работа;
- — обикновена виртуална машина в облака, на която ще се генерира входящият поток от данни:
- — това е Java-приложение, инсталирано локално на машината, което предоставя прост начин за събиране и изпращане на данни в Kinesis (Kinesis Data Streams или Kinesis Firehose). Агентът постоянно следи набор от файлове в зададени директории и изпраща нови данни в Kinesis;
- — Python-скрипт, който прави заявки към API и съхранява отговора в папка, която Kinesis Agent наблюдава;
- — услуга за потоково предаване на данни в реално време с широки възможности за мащабиране;
- — безсерверна услуга, улесняваща анализа на потокови данни в реално време. Amazon Kinesis Data Analytics конфигурира ресурсите за работа на приложения и автоматично мащабира обработката на каквито и да е обеми входящи данни;
- — услуга, позволяваща изпълнение на код без нужда от резервиране и настройка на сървъри. Всички изчислителни мощности автоматично се мащабират под всеки призив;
- — база данни от тип "ключ-стойност" и документи, която осигурява закъснение от по-малко от 10 милисекунди при работа в какъвто и да е мащаб. При използване на DynamoDB не е необходимо да се разпределят несвързани сървъри, да се инсталират пачове или да се управляват. DynamoDB автоматично мащабира таблиците, коригирайки обема налични ресурси и поддържайки висока производителност. Не са необходими действия за администриране на системата;
- — напълно управлявана услуга за изпращане на съобщения по модела "издател — абонат" (Pub/Sub), с което може да се изолират микросервиси, разпределени системи и безсървърни приложения. SNS може да се използва за разпространение на информация до крайните потребители чрез мобилни push уведомления, SMS съобщения и електронни имейли.
Начална подготовка
За симулация на поток от данни реших да използвам информация за самолетни билети, връщани от API Aviasales. В доста обширен списък от различни методи, ще вземем един от тях — "Календар на цените за месеца", който връща цените за всеки ден от месеца, групирани по брой на прелетите. Ако не се предаде месец за търсене в заявката, ще бъде върната информация за месеца, следващ текущия.
Така че, регистрираме се, получаваме своя токен.
Пример за заявка по-долу:
http://api.travelpayouts.com/v2/prices/month-matrix?currency=rub&origin=LED&destination=HKT&show_to_affiliates=true&token=TOKEN_APIПо-гореописаният начин за получаване на данни от API с указание за токена в заявката ще работи, но предпочитам да предавам токена за достъп чрез заглавка, затова в скрипта api_caller.py ще използваме именно този метод.
Примерен отговор:
{{
"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
}]
}
В примера с отговора на API по-горе е показан билет от Санкт Петербург до Пхук… Ах, защо да мечтаем…
Тъй като съм от Казан, а Пукет в момента е "само в нашите мечти", нека да потърсим билети от Санкт Петербург до Казан.
Предполага се, че вече имате акаунт в AWS. Искам веднага да обърна особено внимание, че Kinesis и изпращането на известия чрез SMS не са включени в годишния . Но дори и с това, ако имате на ум няколко долара, спокойно можете да изградите предложената система и да поиграете с нея. И, разбира се, не забравяйте да изтривате всички ресурси след като те вече не са нужни.
С幸ие, DynamoDB и Lambda функции ще бъдат условно безплатни за нас, ако останем в рамките на месечните безплатни лимити. Например, за DynamoDB: 25 GB хранилище, 25 WCU/RCU и 100 милиона заявки. И един милион извиквания на Lambda функции на месец.
Ръчно внедряване на системата
Настройка на Kinesis Data Streams
Преминаваме към услугата Kinesis Data Streams и създаваме два нови потока с по един шард на всеки.
Какво е шард?
Шардът е основната единица за предаване на данни в потока Amazon Kinesis. Един сегмент осигурява предаване на входящи данни със скорост 1 MB/s и предаване на изходящи данни със скорост 2 MB/s. Един сегмент поддържа до 1000 PUT записи в секунда. При създаване на поток от данни трябва да посочите необходимия брой сегменти. Например, можете да създадете поток от данни с два сегмента. Този поток от данни ще осигури предаване на входящи данни със скорост 2 MB/s и предаване на изходящи данни със скорост 4 MB/s с поддръжка на до 2000 PUT записи в секунда.
Колкото повече шардове има във вашия поток - толкова по-голяма е пропускателната му способност. Принципно, потоците се мащабират по този начин - като добавяте шардове. Но колкото повече шардове имате, толкова по-висока е цената. Всеки шард струва 1.5 цента на час и допълнително 1.4 цента за всеки милион операции по добавяне в потока (PUT payload единици).
Нека създадем нов поток с име airline_tickets, за него ще бъде напълно достатъчен 1 шард:

Сега нека да създадем още един поток с име special_stream:

Настройка на продюсера
Като продюсер на данните за разглеждане на задачата е достатъчно да използвате обикновен EC2 инстанс. Не трябва да е мощна скъпа виртуална машина, спотовият t2.micro ще е напълно подходящ.
Важно забележка: за примера следва да се използва image — Amazon Linux AMI 2018.03.0, с него има по-малко настройки за бързо стартиране на Kinesis Agent.
Преминаваме към услугата EC2, създаваме нова виртуална машина, избираме необходимия AMI с тип t2.micro, който попада в Free Tier:

За да може новосъздадената виртуална машина да взаимодейства с услугата Kinesis, е необходимо да ѝ се предоставят права. Най-добрият начин да го направим е да назначим IAM роля. Затова на екрана Step 3: Configure Instance Details трябва да изберем Създайте нова IAM роля:
Създаване на IAM роля за EC2

В отвореното прозорче избираме, че новата роля е за EC2 и преминаваме в раздела Permissions:

На учебния пример не е нужно да вдаваме в детайлите на гранулярната настройка на правата за ресурси, затова избираме предварително настроените политики от Амазон: AmazonKinesisFullAccess и CloudWatchFullAccess.
Даваме някакво смислено име на тази роля, например: EC2-KinesisStreams-FullAccess. В резултат, трябва да получите същото, каквото е показано на картинката по-долу:

След създаването на тази нова роля, не забравяйте да я прикрепите към създаваната инстанция на виртуалната машина:

На този екран не променяме нищо повече и преминаваме към следващите прозорци.
Настройките на твърдия диск могат да останат по подразбиране, таговете също (въпреки че е добра практика да се използват тагове, поне за даване на име на инстанцията и указване на средата).
Сега сме на таба Step 6: Configure Security Group, където е необходимо да създадем нова или да посочим вече съществуваща Security group, позволяваща свързване чрез ssh (порт 22) с инстанцията. Изберете Source —> My IP и можете да стартирате инстанцията.

Когато тя премине в статус running, можете да опитате да се свържете с нея чрез ssh.
За да получите възможност за работа с Kinesis Agent, след успешната свързаност с машината, е необходимо да въведете следните команди в терминала:
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
Създайте папка за запазване на отговорите от API:
sudo mkdir /var/log/airline_ticketsПреди стартиране на агента, е необходимо да настроите неговия конфиг:
sudo vim /etc/aws-kinesis/agent.jsonСъдържанието на файла agent.json трябва да има следния вид:
{
"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"]
}
]
}
]
}
Както изглежда от конфигурационния файл, агентът ще следи в директорията /var/log/airline_tickets/ файлове с разширение .log, ще ги парсва и предава в потока airline_tickets.
Рестартираме услугата и се уверяваме, че тя е стартирана и работи:
sudo service aws-kinesis-agent restartСега ще свалим Python скрипт, който ще искане данни от 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
Скриптът api_caller.py иска данни от Aviasales и запазва получения отговор в директорията, която сканира Kinesis агент. Реализацията на този скрипт е доста стандартна, има клас TicketsApi, който позволява асинхронно извикване на API. В този клас предаваме заглавката с токена и параметрите на заявката:
class TicketsApi:
"""Клас за извикване на API."""
def __init__(self, headers):
"""Инициализационен метод."""
self.base_url = BASE_URL
self.headers = headers
async def get_data(self, data):
"""Получаване на данните от заявката към 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('Статус на отговора %s: %s',
self.base_url, response.status)
response_json = await response.json()
except HTTPError as http_err:
LOGGER.error('Упс! Възникна HTTP грешка: %s', str(http_err))
except Exception as err:
LOGGER.error('Упс! Възникна грешка: %s', str(err))
return response_json
def prepare_request(api_token):
"""Връща заглавките и заявката за 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():
"""Стартиране на кода."""
if len(sys.argv) != 2:
print('Използване: 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 е върнал %s елемента', len(response['data']))
try:
count_rows = log_maker(response)
LOGGER.info('%s реда са записани в %s',
count_rows,
TARGET_FILE)
except Exception as e:
LOGGER.error('Упс! Резултатът от заявката не беше записан в файла. %s',
str(e))
else:
LOGGER.error('Упс! Заявката към API не беше успешна %s!', response)
За тестване на правилността на настройките и работоспособността на агента, ще направим тестово стартиране на скрипта api_caller.py:
sudo ./api_caller.py TOKEN 
И наблюдаваме резултата от работата в логовете на агента и на раздела Мониторинг в потока данни airline_tickets:
tail -f /var/log/aws-kinesis-agent/aws-kinesis-agent.log 

Както виждате, всичко работи и Kinesis Agent успешно изпраща данни в потока. Сега ще настроим consumer.
Настройка на Kinesis Data Analytics
Преминаваме към централния компонент на цялата система — ще създадем ново приложение в Kinesis Data Analytics с името kinesis_analytics_airlines_app:

Kinesis Data Analytics позволява реално време анализ на данни от Kinesis Streams с помощта на SQL. Това е напълно авторазширяем сервиз (в отличие от Kinesis Streams), който:
- позволява създаването на нови потоци (Output Stream) на базата на заявки към изходните данни;
- предоставя поток с грешки, които са възникнали по време на работа на приложенията (Error Stream);
- умее автоматично да определя схемата на входните данни (която може да бъде ръчно променяна при необходимост).
Това не е евтин сервиз — 0.11 USD на час работа, затова трябва да го използвате внимателно и да го изтривате при приключване на работа.
Свързваме приложението към източника на данни:

Избиране на потока, към който ще се свържем (airline_tickets):

След това, е необходимо да прикачим нова IAM роля, за да може приложението да чете от потока и да пише в потока. За целта не е нужно да променяте нищо в блока Access permissions:

Сега ще поискаме откритие на схемата на данни в потока, за целта натискаме бутона «Discover schema». В резултат на това ще се актуализира (създаде нова) IAM роля и ще стартира откритие на схемата от данните, които вече са постъпили в потока:

Сега е необходимо да преминем в SQL редактора. При натискане на този бутон ще се появи прозорец с въпрос за стартиране на приложението — избираме какво искаме да стартираме:

В прозореца на SQL редактора ще поставим такъв прост запит и натискаме 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';
В релационни бази данни работите с таблици, използвайки оператори INSERT за добавяне на записи и оператор SELECT за заявка на данни. В Amazon Kinesis Data Analytics работите с потоци (STREAM) и «помпи» (PUMP) — непрекъснати запроси за вмъкване, които вмъкват данни от един поток в приложението в друг поток.
В представеното по-горе SQL запитване се търсят билети на Аерофлот с цена под пет хиляди рубли. Всички записи, които отговарят на тези условия, ще бъдат поставени в потока DESTINATION_SQL_STREAM.

В блока Destination избираме потока special_stream, а в падащото меню In-application stream name DESTINATION_SQL_STREAM:

В резултат на всички манипулации трябва да се получи нещо подобно на изображението по-долу:

Създаване и абонамент за тема SNS
Преминаваме към услугата Simple Notification Service и създаваме нова тема с име Airlines:

Оформяме абонамента за тази тема, в него указваме номер на мобилен телефон, на който ще получаваме SMS известия:

Създаване на таблица в DynamoDB
За съхранение на суровите данни от потока airline_tickets, ще създадем таблица в DynamoDB с такова име. Като първичен ключ ще използваме record_id:

Създаване на лямбда-функция collector
Създаваме лямбда-функция с името Collector, чиято задача ще бъде да опрашва потока airline_tickets и, в случай на нови записи, да вмъква тези записи в таблицата DynamoDB. Очевидно, че освен правата по подразбиране, тази лямбда трябва да има достъп до четене на потока данни Kinesis и запис в DynamoDB.
Създаване на IAM роля за лямбда-функцията collector
Първо ще създадем нова IAM роля за лямбда с име Lambda-TicketsProcessingRole:

За тестовия пример напълно подходящи са предварително настроените политики AmazonKinesisReadOnlyAccess и AmazonDynamoDBFullAccess, както е показано на изображението по-долу:


Тази лямбда трябва да се задейства от тригера на Kinesis, когато новите записи попаднат в потока airline_stream, затова трябва да добавим нов тригер:


Остава да поставим кода и да запазим лямбдата.
"""Извличане на потока и вмъкване в таблицата DynamoDB."""
import base64
import json
import boto3
from decimal import Decimal
DYNAMO_DB = boto3.resource('dynamodb')
TABLE_NAME = 'airline_tickets'
class TicketsParser:
"""Извличане на информация от потока."""
def __init__(self, table_name, records):
"""Метод за инициализация."""
self.table = DYNAMO_DB.Table(table_name)
self.json_data = TicketsParser.get_json_data(records)
@staticmethod
def get_json_data(records):
"""Връща десериализирани данни от потока."""
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):
"""Предварителна обработка на 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):
"""Пакетно вмъкване в таблицата."""
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('Добавени са ', len(self.json_data), 'елемента')
def lambda_handler(event, context):
"""Извлечете потока и вмъкнете в таблицата DynamoDB."""
print('Получено събитие:', event)
parser = TicketsParser(TABLE_NAME, event['Records'])
parser.run()
Създаване на лямбда-функция notifier
Втората лямбда-функция, която ще следи втория поток (special_stream) и да изпраща уведомление в SNS, се създава по аналогичен начин. Следователно, тази лямбда трябва да има достъп за четене от Kinesis и да изпраща съобщения в определен SNS-топик, който по-късно услугата SNS ще изпрати до всички абонати на този топик (email, SMS и т.н.).
Създаване на IAM роля
Първо създаваме IAM роля Lambda-KinesisAlarm за тази лямбда, а след това назначаваме тази роля на създаваната лямбда alarm_notifier:


Тази лямбда трябва да работи с триггер при добавяне на нови записи в потока special_stream, затова е необходимо да настроим триггера по аналогия с това, което направихме за лямбда Collector.
За удобство при настройка на тази лямбда, ще въведем нова променлива на средата — TOPIC_ARN, където ще поместим ANR (Amazon Resource Names) на топика Airlines:

И вмъкваме кода на лямбдата, който е съвсем прост:
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='Здравейте! Намерих интересни неща!',
Subject='Сигнал за самолетни билети')
print('Сигналното съобщение беше успешно доставено')
except Exception as err:
print('Неуспешно доставяне', str(err))
Изглежда, че ръчната настройка на системата приключи. Остава само да тестваме и да се уверим, че всичко е настроено правилно.
Деплой от Terraform код
Необходима подготовка
— много удобен open-source инструмент за разгръщане на инфраструктура от код. Той има свой собствен синтаксис, който е лесен за усвояване и много примери за това как и какво да разгръщате. В редактора Atom или Visual Studio Code има много удобни плъгини, които улесняват работата с Terraform.
Дистрибуцията може да бъде изтеглена . Подробният преглед на всички възможности на Terraform излиза извън обхвата на тази статия, затова ще се ограничим до основните моменти.
Как да стартирате
Целият код на проекта е наличен . Клонирайте репозитория при себе си. Преди да стартирате, трябва да се уверите, че имате инсталиран и настроен AWS CLI, тъй като Terraform ще търси идентификационни данни в файла ~/ .aws/credentials.
Добра практика е преди разгръщането на цялата инфраструктура да се стартира командата plan, за да се види какво ще създаде Terraform в облака:
terraform.exe planЩе бъде предложено да въведете телефонен номер, за да получавате уведомления. На този етап не е задължително да го въвеждате.

След анализ на плана за работа на програмата, можем да стартираме създаването на ресурси:
terraform.exe applyСлед изпращането на тази команда отново ще бъде поставен въпрос за въвеждане на телефонен номер, въвеждаме 'yes', когато се покаже въпрос за фактическото изпълнение на действията. Това ще позволи да се издигне цялата инфраструктура, да се проведе необходимата настройка на EC2, да се разгръщат лямбда функции и т.н.
След като всички ресурси са успешно създадени чрез кода на Terraform, трябва да влезете в детайлите на приложението Kinesis Analytics (за съжаление, не успях да намеря как да направя това веднага от кода).
Стартиране на приложението:

След това е необходимо явно да зададете името на in-application stream, като изберете от падащия списък:


Сега всичко е готово за работа.
Тестване на работата на приложението
Независимо как сте разгръщали системата, ръчно или чрез Terraform код, тя ще работи по същия начин.
Свързваме се по SSH с виртуалната машина EC2, на която е инсталиран Kinesis Agent, и стартираме скрипта api_caller.py.
sudo ./api_caller.py TOKENТрябва само да изчакаме SMS на вашия номер:

SMS – съобщението пристига на телефона почти за 1 минута:

Остава да проверим дали записите са запазени в базата данни DynamoDB за по-подробен анализ по-късно. Таблицата airline_tickets съдържа приблизително следните данни:

Заключение
В хода на работата беше изградена система за онлайн обработка на данни на базата на Amazon Kinesis. Бяха разгледани варианти за използването на Kinesis Agent заедно с Kinesis Data Streams и реално-времева аналитика с Kinesis Analytics с помощта на SQL команди, както и взаимодействието на Amazon Kinesis с други услуги на AWS.
По-гореописаната система разположихме по два начина: достатъчно бавен ръчен метод и бърз чрез кода на Terraform.
Целият изходен код на проекта е достъпен , предлагам да се запознаете с него.
С удоволствие ще обсъдя статията, очаквам вашите коментари. Надявам се на конструктивна критика.
Желая успехи!
Източник: habr.com
