Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Здравейте, Хабр!

Обичате ли да летите със самолети? Аз обожавам, но по време на самоизолацията започнах и да анализирам данни за самолетни билети от един известен ресурс - Aviasales.

Днес ще разгледаме работата на Amazon Kinesis, ще построим стрийминг система с анализ в реално време, ще инсталираме NoSQL база данни Amazon DynamoDB като основно хранилище и ще настроим известия чрез SMS за интересни билети.

Всички подробности по-долу! Да започнем!

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Въведение

За примера ще ни е необходим достъп до API Aviasales. Достъпът до него е безплатен и неограничен, необходимо е само да се регистрирате в секцията „Разработчици“, за да получите своя API токен за достъп до данните.

Основната цел на тази статия е да предостави общо разбиране за използването на потоково предаване на информация в AWS, като изключваме факта, че данните, предоставени от използвания API, не са строго актуални и се предават от кеш, който се образува на базата на търсенията на потребителите на сайтовете Aviasales.ru и Jetradar.com през последните 48 часа.

Данните за самолетни билети, получени чрез API, Kinesis-agent, инсталиран на машината-производител, автоматично ще се парсват и предават в необходимия поток чрез Kinesis Data Analytics. Необработената версия на този поток ще се записва директно в хранилището. Разгърнатото в DynamoDB хранилище на „сурови“ данни ще позволи извършването на по-дълбок анализ на билетите чрез BI инструменти, например AWS Quick Sight.

Ще разгледаме два варианта за деплой на цялата инфраструктура:

  • Ръчен — чрез AWS Management Console;
  • Инфраструктура на код Terraform — за ленивите автоматизатори;

Архитектура на разработваната система

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Използвани компоненти:

  • Aviasales API — данните, предоставени от този API, ще се използват за цялата последваща работа;
  • EC2 Producer Instance — обикновена виртуална машина в облака, на която ще се генерира входящият поток от данни:
    • Kinesis Agent — това е Java-приложение, инсталирано локално на машината, което предоставя прост начин за събиране и изпращане на данни в Kinesis (Kinesis Data Streams или Kinesis Firehose). Агентът постоянно следи набор от файлове в зададени директории и изпраща нови данни в Kinesis;
    • Скрипт API Caller — Python-скрипт, който прави заявки към API и съхранява отговора в папка, която Kinesis Agent наблюдава;
  • Kinesis Data Streams — услуга за потоково предаване на данни в реално време с широки възможности за мащабиране;
  • Kinesis Analytics — безсерверна услуга, улесняваща анализа на потокови данни в реално време. Amazon Kinesis Data Analytics конфигурира ресурсите за работа на приложения и автоматично мащабира обработката на каквито и да е обеми входящи данни;
  • AWS Lambda — услуга, позволяваща изпълнение на код без нужда от резервиране и настройка на сървъри. Всички изчислителни мощности автоматично се мащабират под всеки призив;
  • Amazon DynamoDB — база данни от тип "ключ-стойност" и документи, която осигурява закъснение от по-малко от 10 милисекунди при работа в какъвто и да е мащаб. При използване на DynamoDB не е необходимо да се разпределят несвързани сървъри, да се инсталират пачове или да се управляват. DynamoDB автоматично мащабира таблиците, коригирайки обема налични ресурси и поддържайки висока производителност. Не са необходими действия за администриране на системата;
  • Amazon SNS — напълно управлявана услуга за изпращане на съобщения по модела "издател — абонат" (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 не са включени в годишния Free Tier (безплатно използване). Но дори и с това, ако имате на ум няколко долара, спокойно можете да изградите предложената система и да поиграете с нея. И, разбира се, не забравяйте да изтривате всички ресурси след като те вече не са нужни.

С幸ие, 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 шард:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Сега нека да създадем още един поток с име special_stream:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Настройка на продюсера

Като продюсер на данните за разглеждане на задачата е достатъчно да използвате обикновен EC2 инстанс. Не трябва да е мощна скъпа виртуална машина, спотовият t2.micro ще е напълно подходящ.

Важно забележка: за примера следва да се използва image — Amazon Linux AMI 2018.03.0, с него има по-малко настройки за бързо стартиране на Kinesis Agent.

Преминаваме към услугата EC2, създаваме нова виртуална машина, избираме необходимия AMI с тип t2.micro, който попада в Free Tier:

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

Създаване на IAM роля за EC2
Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
В отвореното прозорче избираме, че новата роля е за EC2 и преминаваме в раздела Permissions:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
На учебния пример не е нужно да вдаваме в детайлите на гранулярната настройка на правата за ресурси, затова избираме предварително настроените политики от Амазон: AmazonKinesisFullAccess и CloudWatchFullAccess.

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
След създаването на тази нова роля, не забравяйте да я прикрепите към създаваната инстанция на виртуалната машина:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
На този екран не променяме нищо повече и преминаваме към следващите прозорци.

Настройките на твърдия диск могат да останат по подразбиране, таговете също (въпреки че е добра практика да се използват тагове, поне за даване на име на инстанцията и указване на средата).

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Когато тя премине в статус 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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
И наблюдаваме резултата от работата в логовете на агента и на раздела Мониторинг в потока данни airline_tickets:

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Както виждате, всичко работи и Kinesis Agent успешно изпраща данни в потока. Сега ще настроим consumer.

Настройка на Kinesis Data Analytics

Преминаваме към централния компонент на цялата система — ще създадем ново приложение в Kinesis Data Analytics с името kinesis_analytics_airlines_app:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Kinesis Data Analytics позволява реално време анализ на данни от Kinesis Streams с помощта на SQL. Това е напълно авторазширяем сервиз (в отличие от Kinesis Streams), който:

  1. позволява създаването на нови потоци (Output Stream) на базата на заявки към изходните данни;
  2. предоставя поток с грешки, които са възникнали по време на работа на приложенията (Error Stream);
  3. умее автоматично да определя схемата на входните данни (която може да бъде ръчно променяна при необходимост).

Това не е евтин сервиз — 0.11 USD на час работа, затова трябва да го използвате внимателно и да го изтривате при приключване на работа.

Свързваме приложението към източника на данни:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Избиране на потока, към който ще се свържем (airline_tickets):

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
След това, е необходимо да прикачим нова IAM роля, за да може приложението да чете от потока и да пише в потока. За целта не е нужно да променяте нищо в блока Access permissions:

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

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
В прозореца на 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.

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
В блока Destination избираме потока special_stream, а в падащото меню In-application stream name DESTINATION_SQL_STREAM:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
В резултат на всички манипулации трябва да се получи нещо подобно на изображението по-долу:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Създаване и абонамент за тема SNS

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Оформяме абонамента за тази тема, в него указваме номер на мобилен телефон, на който ще получаваме SMS известия:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Създаване на таблица в DynamoDB

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Създаване на лямбда-функция collector

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

Създаване на IAM роля за лямбда-функцията collector
Първо ще създадем нова IAM роля за лямбда с име Lambda-TicketsProcessingRole:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
За тестовия пример напълно подходящи са предварително настроените политики AmazonKinesisReadOnlyAccess и AmazonDynamoDBFullAccess, както е показано на изображението по-долу:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Остава да поставим кода и да запазим лямбдата.

"""Извличане на потока и вмъкване в таблицата 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:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Тази лямбда трябва да работи с триггер при добавяне на нови записи в потока special_stream, затова е необходимо да настроим триггера по аналогия с това, което направихме за лямбда Collector.

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

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
И вмъкваме кода на лямбдата, който е съвсем прост:

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 код

Необходима подготовка

Terraform — много удобен open-source инструмент за разгръщане на инфраструктура от код. Той има свой собствен синтаксис, който е лесен за усвояване и много примери за това как и какво да разгръщате. В редактора Atom или Visual Studio Code има много удобни плъгини, които улесняват работата с Terraform.

Дистрибуцията може да бъде изтеглена оттук. Подробният преглед на всички възможности на Terraform излиза извън обхвата на тази статия, затова ще се ограничим до основните моменти.

Как да стартирате

Целият код на проекта е наличен в моят репозиторий. Клонирайте репозитория при себе си. Преди да стартирате, трябва да се уверите, че имате инсталиран и настроен AWS CLI, тъй като Terraform ще търси идентификационни данни в файла ~/ .aws/credentials.

Добра практика е преди разгръщането на цялата инфраструктура да се стартира командата plan, за да се види какво ще създаде Terraform в облака:

terraform.exe plan

Ще бъде предложено да въведете телефонен номер, за да получавате уведомления. На този етап не е задължително да го въвеждате.

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
След анализ на плана за работа на програмата, можем да стартираме създаването на ресурси:

terraform.exe apply

След изпращането на тази команда отново ще бъде поставен въпрос за въвеждане на телефонен номер, въвеждаме 'yes', когато се покаже въпрос за фактическото изпълнение на действията. Това ще позволи да се издигне цялата инфраструктура, да се проведе необходимата настройка на EC2, да се разгръщат лямбда функции и т.н.

След като всички ресурси са успешно създадени чрез кода на Terraform, трябва да влезете в детайлите на приложението Kinesis Analytics (за съжаление, не успях да намеря как да направя това веднага от кода).

Стартиране на приложението:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
След това е необходимо явно да зададете името на in-application stream, като изберете от падащия списък:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Сега всичко е готово за работа.

Тестване на работата на приложението

Независимо как сте разгръщали системата, ръчно или чрез Terraform код, тя ще работи по същия начин.

Свързваме се по SSH с виртуалната машина EC2, на която е инсталиран Kinesis Agent, и стартираме скрипта api_caller.py.

sudo ./api_caller.py TOKEN

Трябва само да изчакаме SMS на вашия номер:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
SMS – съобщението пристига на телефона почти за 1 минута:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless
Остава да проверим дали записите са запазени в базата данни DynamoDB за по-подробен анализ по-късно. Таблицата airline_tickets съдържа приблизително следните данни:

Интеграция на Aviasales API с Amazon Kinesis и лесен serverless

Заключение

В хода на работата беше изградена система за онлайн обработка на данни на базата на Amazon Kinesis. Бяха разгледани варианти за използването на Kinesis Agent заедно с Kinesis Data Streams и реално-времева аналитика с Kinesis Analytics с помощта на SQL команди, както и взаимодействието на Amazon Kinesis с други услуги на AWS.

По-гореописаната система разположихме по два начина: достатъчно бавен ръчен метод и бърз чрез кода на Terraform.

Целият изходен код на проекта е достъпен в моя репозиторий в GitHub, предлагам да се запознаете с него.

С удоволствие ще обсъдя статията, очаквам вашите коментари. Надявам се на конструктивна критика.

Желая успехи!

Източник: habr.com

Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри 🔥 Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри | ProHoster