Ciao, Habr!
Ti piace volare sugli aerei? A me piace molto, ma durante l'autoisolamento ho iniziato anche ad analizzare i dati sui biglietti aerei di una risorsa nota — Aviasales.
Oggi analizzeremo il funzionamento di Amazon Kinesis, costruiremo un sistema di streaming con analisi in tempo reale, utilizzeremo Amazon DynamoDB come principale archivio dati e imposteremo notifiche SMS per biglietti interessanti.
Tutti i dettagli dopo il salto! Andiamo!

Introduzione
Per esempio, avremo bisogno di accesso a . L'accesso è gratuito e illimitato, è necessario solo registrarsi nella sezione «Sviluppatori» per ottenere il proprio token API per accedere ai dati.
L'obiettivo principale di questo articolo è fornire una comprensione generale dell'uso della trasmissione di informazioni in AWS, escludendo il fatto che i dati restituiti dall'API utilizzata non siano rigorosamente aggiornati e vengano forniti da una cache, che è formata sulla base delle ricerche degli utenti sui siti Aviasales.ru e Jetradar.com nelle ultime 48 ore.
I dati sui biglietti aerei ottenuti tramite l'API saranno automaticamente analizzati e trasmessi nel flusso giusto tramite Kinesis Data Analytics. La versione non elaborata di questo flusso sarà scritta direttamente nell'archivio. L'archivio di dati «grezzi» sviluppato in DynamoDB consentirà di eseguire analisi più approfondite sui biglietti tramite strumenti BI, ad esempio AWS Quick Sight.
Esamineremo due opzioni per il deployment dell'intera infrastruttura:
- Manuale — tramite AWS Management Console;
- Infrastruttura da codice Terraform — per gli automatizzatori pigri;
Architettura del sistema in fase di sviluppo

Componenti utilizzati:
- — i dati restituiti da questa API saranno utilizzati per tutto il lavoro successivo;
- — una macchina virtuale comune nel cloud, sulla quale verrà generato il flusso di dati in ingresso:
- — è un'applicazione Java che viene installata localmente sulla macchina e offre un modo semplice per raccogliere e inviare dati a Kinesis (Kinesis Data Streams o Kinesis Firehose). L'agente monitora costantemente un insieme di file nelle directory specificate e invia nuovi dati a Kinesis;
- — uno script Python che invia richieste all'API e memorizza le risposte in una cartella monitorata da Kinesis Agent;
- — servizio di streaming dati in tempo reale con ampie capacità di scalabilità;
- — un servizio senza server che semplifica l'analisi dei flussi di dati in tempo reale. Amazon Kinesis Data Analytics configura le risorse per il funzionamento delle applicazioni e si scala automaticamente per gestire qualsiasi volume di dati in ingresso;
- — un servizio che consente di eseguire codice senza dover riservare e configurare server. Tutte le risorse computazionali si scalano automaticamente per ogni chiamata;
- — un database per coppie "chiave-valore" e documenti, che offre una latenza inferiore a 10 millisecondi operando su qualsiasi scala. Utilizzando DynamoDB, non è necessario distribuire server, applicare patch o gestirli. DynamoDB scala automaticamente le tabelle, regolando la quantità di risorse disponibili e mantenendo elevate prestazioni. Non sono necessarie azioni di amministrazione del sistema;
- — un servizio completamente gestito per l'invio di messaggi secondo il modello "pubblicatore - sottoscrittore" (Pub/Sub), grazie al quale è possibile isolare microservizi, sistemi distribuiti e applicazioni senza server. SNS può essere utilizzato per inviare informazioni agli utenti finali attraverso notifiche push mobili, messaggi SMS e email.
Formazione iniziale
Per emulare un flusso di dati ho deciso di utilizzare le informazioni sui biglietti aerei restituite dall'API Aviasales. In un elenco abbastanza ampio di vari metodi, prendiamo uno di essi — "Calendario dei prezzi per mese", che restituisce i prezzi per ogni giorno del mese, raggruppati in base al numero di scali. Se non si passa il mese di ricerca nella richiesta, verranno restituiti i dati per il mese successivo all'attuale.
Quindi, ci registriamo e otteniamo il nostro token.
Un esempio di richiesta è di seguito:
http://api.travelpayouts.com/v2/prices/month-matrix?currency=rub&origin=LED&destination=HKT&show_to_affiliates=true&token=TOKEN_APIIl metodo descritto sopra per ottenere dati dall'API specificando il token nella richiesta funzionerà, ma preferisco passare il token di accesso tramite l'intestazione, quindi nel codice di api_caller.py utilizzeremo proprio questo metodo.
Esempio di risposta:
{{
"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
}]
}
Nell'esempio di risposta dell'API sopra è mostrato un biglietto da San Pietroburgo a Phuket... Ah, perché sognare...
Poiché vengo da Kazan e Phuket ora è solo un sogno, cerchiamo voli da San Pietroburgo a Kazan.
Si presume che tu abbia già un account AWS. Vorrei sottolineare che Kinesis e l'invio di notifiche tramite SMS non fanno parte dell' . Tuttavia, anche tenendo a mente un paio di dollari, è possibile costruire il sistema proposto e divertirsi con esso. E, naturalmente, non dimenticare di eliminare tutte le risorse dopo che non sono più necessarie.
Fortunatamente, DynamoDb e le funzioni lambda saranno per noi sostanzialmente gratuite, se rimaniamo entro i limiti mensili gratuiti. Ad esempio, per DynamoDB: 25 GB di spazio di archiviazione, 25 WCU/RCU e 100 milioni di richieste. E un milione di invocazioni di funzioni lambda al mese.
Distribuzione manuale del sistema
Impostazione di Kinesis Data Streams
Passiamo al servizio Kinesis Data Streams e creiamo due nuovi flussi con uno shard ciascuno.
Cos'è uno shard?
Uno shard è l'unità principale di trasmissione dei dati nel flusso Amazon Kinesis. Un segmento fornisce la trasmissione dei dati in ingresso a una velocità di 1 MB/s e la trasmissione dei dati in uscita a una velocità di 2 MB/s. Un segmento supporta fino a 1000 registrazioni PUT al secondo. Quando crei un flusso di dati, devi specificare il numero necessario di segmenti. Ad esempio, puoi creare un flusso di dati con due segmenti. Questo flusso di dati garantirà la trasmissione dei dati in ingresso a una velocità di 2 MB/s e la trasmissione dei dati in uscita a una velocità di 4 MB/s, supportando fino a 2000 registrazioni PUT al secondo.
Maggiore è il numero di shard nel tuo flusso, maggiore sarà la sua capacità. In effetti, i flussi si scalano aggiungendo shard. Ma più shard hai, più alto è il costo. Ogni shard costa 1,5 centesimi all'ora e ulteriori 1,4 centesimi per ogni milione di operazioni di aggiunta al flusso (PUT payload units).
Creiamo un nuovo flusso di nome airline_tickets, avrà bisogno di un solo shard:

Ora creiamo un altro flusso di nome special_stream:

Impostazione del produttore
Come produttore di dati per affrontare il compito, è sufficiente utilizzare una normale istanza EC2. Non deve essere una macchina virtuale potente e costosa, va bene anche un t2.micro spot.
Nota importante: per l'esempio, si consiglia di utilizzare l'immagine - Amazon Linux AMI 2018.03.0, con essa ci sono meno impostazioni per un avvio rapido di Kinesis Agent.
Andiamo nel servizio EC2, creiamo una nuova macchina virtuale, scegliamo l'AMI necessaria con tipo t2.micro, che rientra nel Free Tier:

Affinché la nuova macchina virtuale possa interagire con il servizio Kinesis, è necessario concederle i permessi necessari. Il modo migliore per farlo è assegnare un IAM Role. Quindi, nella schermata Step 3: Configura i dettagli dell'istanza, è necessario selezionare Crea un nuovo IAM Role:
Creazione di un IAM Role per EC2

Nella finestra aperta, scegliamo che stiamo creando un nuovo ruolo per EC2 e passiamo alla sezione Permessi:

Nell'esempio didattico non ci soffermeremo su tutte le sfumature della configurazione granulare dei diritti sui risorse, quindi selezioniamo le policy preimpostate da Amazon: AmazonKinesisFullAccess e CloudWatchFullAccess.
Diamo un nome significativo a questo ruolo, ad esempio: EC2-KinesisStreams-FullAccess. Di conseguenza, dovrebbe risultare lo stesso di quanto indicato nell'immagine sottostante:

Dopo aver creato questo nuovo ruolo, non dimentichiamo di associarlo all'istanza della macchina virtuale che stiamo creando:

Non cambiamo nient'altro in questa schermata e passiamo alle schermate successive.
Le impostazioni del disco rigido possono essere lasciate di default, anche i tag (anche se è buona pratica utilizzare i tag, perlomeno dando un nome all’istanza e specificando l’ambiente).
Ora siamo nella scheda Step 6: Configura il gruppo di sicurezza, dove è necessario creare un nuovo gruppo di sicurezza o indicare quello esistente che consente di connettersi tramite ssh (porta 22) all'istanza. Selezionate lì Sorgente —> Il mio IP e potete avviare l'istanza.

Non appena passerà allo stato di running, potrete provare a connettervi tramite ssh.
Per poter lavorare con Kinesis Agent, dopo essere riusciti a connetterci alla macchina, è necessario inserire i seguenti comandi nel terminale:
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
Creiamo una cartella per salvare le risposte dell'API:
sudo mkdir /var/log/airline_ticketsPrima di avviare l'agente, è necessario configurare il suo file di configurazione:
sudo vim /etc/aws-kinesis/agent.jsonIl contenuto del file agent.json dovrebbe avere il seguente aspetto:
{
"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"]
}
]
}
]
}
Come si può vedere dal file di configurazione, l'agente monitorerà nella directory /var/log/airline_tickets/ i file con estensione .log, li parserà e li invierà nel flusso airline_tickets.
Riavviamo il servizio e ci assicuriamo che sia partito e funzionante:
sudo service aws-kinesis-agent restartOra scarichiamo lo script Python che richiederà i dati dall'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
Lo script api_caller.py richiede dati a Aviasales e salva la risposta ricevuta nella directory che l'agente Kinesis esamina. L'implementazione di questo script è abbastanza standard, c'è la classe TicketsApi, che permette di estrarre l'API in modo asincrono. In questa classe passiamo l'intestazione con il token e i parametri della richiesta:
class TicketsApi:
"""Classe per chiamare l'API."""
def __init__(self, headers):
"""Metodo di inizializzazione."""
self.base_url = BASE_URL
self.headers = headers
async def get_data(self, data):
"""Ottieni i dati dalla query 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 della risposta %s: %s',
self.base_url, response.status)
response_json = await response.json()
except HTTPError as http_err:
LOGGER.error('Oops! Si è verificato un errore HTTP: %s', str(http_err))
except Exception as err:
LOGGER.error('Oops! Si è verificato un errore: %s', str(err))
return response_json
def prepare_request(api_token):
"""Restituisci le intestazioni e la query per la richiesta 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():
"""Esegui il codice."""
if len(sys.argv) != 2:
print('Utilizzo: 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('L'API ha restituito %s articoli', len(response['data']))
try:
count_rows = log_maker(response)
LOGGER.info('%s righe sono state salvate in %s',
count_rows,
TARGET_FILE)
except Exception as e:
LOGGER.error('Oops! Il risultato della richiesta non è stato salvato nel file. %s',
str(e))
else:
LOGGER.error('Oops! La richiesta API non è stata riuscita %s!', response)
Per testare la correttezza delle impostazioni e la funzionalità dell'agente, faremo un'esecuzione di prova dello script api_caller.py:
sudo ./api_caller.py TOKEN 
E controlliamo il risultato del lavoro nei log dell'agente e nella scheda Monitoraggio nel flusso dati airline_tickets:
tail -f /var/log/aws-kinesis-agent/aws-kinesis-agent.log 

Come si può vedere, tutto funziona e Kinesis Agent sta inviando con successo i dati nel flusso. Ora configuriamo il consumer.
Configurazione di Kinesis Data Analytics
Passiamo al componente centrale dell'intero sistema — creiamo una nuova applicazione in Kinesis Data Analytics con il nome kinesis_analytics_airlines_app:

Kinesis Data Analytics consente di eseguire analisi dei dati in tempo reale dai Kinesis Streams utilizzando il linguaggio SQL. Si tratta di un servizio completamente scalabile automaticamente (a differenza di Kinesis Streams), che:
- permette di creare nuovi flussi (Output Stream) basati su query sui dati sorgente;
- fornisce un flusso con errori che si sono verificati durante l'esecuzione delle applicazioni (Error Stream);
- può rilevare automaticamente lo schema dei dati in ingresso (può essere sovrascritto manualmente se necessario).
Questo è un servizio costoso — 0.11 USD all'ora, quindi bisogna usarlo con attenzione e rimuoverlo al termine del lavoro.
Colleghiamo l'applicazione alla fonte di dati:

Selezioniamo il flusso a cui intendiamo collegarci (airline_tickets):

Successivamente, è necessario allegare un nuovo Ruolo IAM affinché l'applicazione possa leggere dal flusso e scrivere nel flusso. Per questo basta non modificare nulla nel blocco Access permissions:

Ora richiediamo il rilevamento dello schema dati nel flusso, per questo premiamo il pulsante «Discover schema». Di conseguenza, verrà aggiornata (creata una nuova) ruolo IAM e inizierà il rilevamento dello schema dai dati che sono già arrivati nel flusso:

È ora necessario passare all'editor SQL. Cliccando su questo pulsante, apparirà una finestra con la domanda sul lancio dell'applicazione — scegliamo cosa vogliamo avviare:

Nella finestra dell'editor SQL incolliamo questa semplice query e premiamo Salva e Esegui 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';
Nelle basi di dati relazionali si lavora con tabelle, utilizzando gli operatori INSERT per aggiungere record e l'operatore SELECT per interrogare i dati. In Amazon Kinesis Data Analytics si lavora con flussi (STREAM) e 'pompe' (PUMP) — query continue di inserimento, che trasferiscono dati da un flusso dell'applicazione a un altro flusso.
Nella query SQL sopra presentata avviene la ricerca di biglietti Aeroflot al costo inferiore di cinquemila rubli. Tutti i record che soddisfano queste condizioni saranno posizionati nel flusso DESTINATION_SQL_STREAM.

Nel blocco Destination scegliamo il flusso special_stream, e nel menu a discesa In-application stream name DESTINATION_SQL_STREAM:

Dopo tutte le manipolazioni, dovremmo ottenere qualcosa di simile all'immagine sottostante:

Creazione e iscrizione a un topic SNS
Accediamo al servizio Simple Notification Service e creiamo un nuovo topic con il nome Airlines:

Registriamo l'iscrizione a questo topic, specificando il numero di telefono mobile a cui verranno inviati gli SMS:

Creazione della tabella in DynamoDB
Per memorizzare i dati non elaborati del flusso airline_tickets, creeremo una tabella in DynamoDB con lo stesso nome. Utilizzeremo record_id come chiave primaria:

Creazione della funzione lambda collector
Creeremo una funzione lambda chiamata Collector, il cui compito sarà monitorare il flusso airline_tickets e, in caso di nuove registrazioni, inserire queste registrazioni nella tabella DynamoDB. È evidente che, oltre ai diritti di accesso predefiniti, questa lambda deve avere accesso in lettura al flusso di dati Kinesis e in scrittura a DynamoDB.
Creazione di un ruolo IAM per la funzione lambda collector
Per iniziare, creiamo un nuovo ruolo IAM per la lambda con il nome Lambda-TicketsProcessingRole:

Per un esempio di test possono andare bene le policy preimpostate AmazonKinesisReadOnlyAccess e AmazonDynamoDBFullAccess, come mostrato nell'immagine sottostante:


Questa lambda deve essere attivata da un trigger di Kinesis quando vengono aggiunte nuove registrazioni al flusso airline_stream, quindi è necessario aggiungere un nuovo trigger:


Rimane solo da inserire il codice e salvare la lambda.
"""Analisi del flusso e inserimento nella tabella DynamoDB."""
import base64
import json
import boto3
from decimal import Decimal
DYNAMO_DB = boto3.resource('dynamodb')
TABLE_NAME = 'airline_tickets'
class TicketsParser:
"""Analizza le informazioni dal flusso."""
def __init__(self, table_name, records):
"""Metodo di inizializzazione."""
self.table = DYNAMO_DB.Table(table_name)
self.json_data = TicketsParser.get_json_data(records)
@staticmethod
def get_json_data(records):
"""Restituisce i dati deserializzati dal flusso."""
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):
"""Pre-elabora i dati 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):
"""Inserimento batch nella tabella."""
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('Sono stati aggiunti ', len(self.json_data), 'elementi')
def lambda_handler(event, context):
"""Analizza il flusso e inserisci nella tabella DynamoDB."""
print('Evento ricevuto:', event)
parser = TicketsParser(TABLE_NAME, event['Records'])
parser.run()
Creazione della funzione lambda notifier
La seconda funzione lambda, che monitorerà il secondo flusso (special_stream) e invierà una notifica a SNS, viene creata in modo analogo. Pertanto, questa lambda deve avere accesso in lettura da Kinesis e inviare messaggi al topic SNS specificato, il quale sarà successivamente instradato a tutti gli abbonati di questo topic (email, SMS, ecc.).
Creazione di un ruolo IAM
Innanzitutto creiamo il ruolo IAM Lambda-KinesisAlarm per questa lambda, e poi assegniamo questo ruolo alla lambda alarm_notifier che stiamo creando:


Questa lambda deve attivarsi quando arrivano nuovi record nel flusso special_stream, quindi è necessario configurare il trigger in modo simile a quanto fatto per la lambda Collector.
Per facilitare la configurazione di questa lambda, introduciamo una nuova variabile d'ambiente — TOPIC_ARN, dove inseriamo l'ARN (Amazon Resource Names) del topic Airlines:

E inseriamo il codice della lambda, che è piuttosto semplice:
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='Ciao! Ho trovato delle cose interessanti!',
Subject='Allerta biglietti aerei')
print('Il messaggio di allerta è stato consegnato con successo')
except Exception as err:
print('Fallimento della consegna', str(err))
Sembra che ora la configurazione manuale del sistema sia completata. Resta solo da testare e assicurarsi che abbiamo impostato tutto correttamente.
Distribuzione del codice Terraform
Preparazione necessaria
— uno strumento open-source molto utile per il deployment dell'infrastruttura dal codice. Ha una propria sintassi che è facile da apprendere e numerosi esempi su come e cosa distribuire. Nell'editor Atom o Visual Studio Code ci sono molti plugin utili che facilitano il lavoro con Terraform.
È possibile scaricare il pacchetto . Un'analisi dettagliata di tutte le funzionalità di Terraform va oltre l'ambito di questo articolo, quindi ci limiteremo ai punti principali.
Come avviare
Il codice completo del progetto si trova . Cloniamo il repository. Prima di avviare, è necessario assicurarsi che AWS CLI sia installato e configurato, poiché Terraform cercherà le credenziali nel file ~/\.aws/credentials.
È una buona pratica eseguire il comando plan prima di distribuire l'intera infrastruttura, per vedere cosa Terraform sta per creare nel cloud:
terraform.exe planVerrà richiesto di inserire il numero di telefono per ricevere notifiche. A questo punto non è obbligatorio fornirlo.

Dopo aver analizzato il piano di lavoro del programma, possiamo avviare la creazione delle risorse:
terraform.exe applyDopo aver inviato questo comando, verrà nuovamente richiesto di inserire il numero di telefono, digita “yes” quando verrà mostrato il messaggio riguardante l'esecuzione effettiva delle azioni. Questo permetterà di avviare l'intera infrastruttura, effettuare tutte le configurazioni necessarie per EC2, distribuire funzioni lambda, ecc.
Dopo che tutte le risorse sono state create con successo tramite il codice Terraform, è necessario accedere ai dettagli dell'applicazione Kinesis Analytics (purtroppo non ho trovato un modo per farlo direttamente dal codice).
Avviamo l'applicazione:

Dopo questo, è necessario specificare esplicitamente il nome del flusso in-app, selezionandolo dall'elenco a discesa:


Ora tutto è pronto per funzionare.
Test del funzionamento dell'applicazione
Indipendentemente da come hai distribuito il sistema, manualmente o tramite codice Terraform, funzionerà allo stesso modo.
Accediamo via SSH alla macchina virtuale EC2, dove è installato Kinesis Agent e avviamo lo script api_caller.py
sudo ./api_caller.py TOKENOra dobbiamo aspettare l'SMS sul vostro numero:

L'SMS - il messaggio arriva sul telefono praticamente dopo 1 minuto:

Resta da vedere se i record sono stati salvati nel database DynamoDB per un'analisi successiva più dettagliata. La tabella airline_tickets contiene dati simili a questi:

Conclusione
Durante il lavoro svolto è stato costruito un sistema di elaborazione dati online basato su Amazon Kinesis. Sono state esaminate le opzioni di utilizzo di Kinesis Agent in combinazione con Kinesis Data Streams e l'analisi in tempo reale di Kinesis Analytics tramite comandi SQL, oltre all'interazione di Amazon Kinesis con altri servizi AWS.
Il sistema descritto sopra è stato implementato in due modi: uno manuale piuttosto lungo e uno veloce utilizzando il codice Terraform.
Tutto il codice sorgente del progetto è disponibile , vi invito a prenderne visione.
Sarò lieto di discutere l'articolo, attendo i vostri commenti. Spero in una critica costruttiva.
Ti auguro successo!
Fonte: habr.com
