0Pricing
Cloud & IT Cert Prep · Lezione

Kinesis Streams, Firehose e analisi in tempo reale

Acquisire dati in streaming con Kinesis Data Streams, distribuirli in S3 o Redshift con Firehose e analizzarli in tempo reale con Managed Service for Apache Flink

Kinesis Streams, Firehose e analisi in tempo reale è una lezione Cloud & IT Cert Prep gratuita su CoddyKit. Questa è la lezione 4 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento Cloud & IT Cert Prep, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso Cloud & IT Cert Prep include 4 lezioni in totale.

Panoramica della famiglia Kinesis

Amazon Kinesis è una famiglia di servizi per la raccolta, l'elaborazione e l'analisi di dati in streaming in tempo reale. I tre servizi principali sono: Kinesis Data Streams (elaborazione personalizzata a bassa latenza), Kinesis Data Firehose (distribuzione completamente gestita verso S3/Redshift/OpenSearch) e Managed Service for Apache Flink (in precedenza Kinesis Data Analytics), per SQL e l'elaborazione Flink in tempo reale. Ogni servizio è destinato a una fase diversa della pipeline di streaming.

Architettura di Kinesis Data Streams

Un Kinesis Data Stream è un log durevole e ordinato, suddiviso in shard. Ogni shard offre una velocità di scrittura di 1 MB/s e una velocità di lettura di 2 MB/s. I record di dati vengono conservati per impostazione predefinita per 24 ore, con possibilità di estendere il periodo a 7 o 365 giorni. I produttori scrivono i record in uno stream; i consumer — Lambda, applicazioni KCL, Firehose o Flink — leggono da uno o più shard in parallelo. Una volta scritti, i record sono immutabili.

# Create a Kinesis Data Stream with 4 shards
aws kinesis create-stream \
  --stream-name clickstream \
  --shard-count 4

# Put a record into the stream
aws kinesis put-record \
  --stream-name clickstream \
  --partition-key 'user-123' \
  --data 'eyJldmVudCI6ICJjbGljayJ9'

Shard, throughput e scalabilità

Il numero di shard determina il throughput totale di uno stream. Può suddividere uno shard per raddoppiare il throughput oppure unire due shard per ridurre i costi. Utilizzi Enhanced Fan-Out per assegnare a ogni consumer registrato un throughput di lettura dedicato di 2 MB/s, indipendente dagli altri consumer, eliminando il throttling delle letture quando più applicazioni utilizzano lo stesso stream. Monitori GetRecords.IteratorAgeMilliseconds per rilevare il ritardo dei consumer.

# Split shard to increase throughput
aws kinesis split-shard \
  --stream-name clickstream \
  --shard-to-split shardId-000000000001 \
  --new-starting-hash-key 170141183460469231731687303715884105728

# Register an enhanced fan-out consumer
aws kinesis register-stream-consumer \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/clickstream \
  --consumer-name analytics-app

Kinesis Data Firehose: distribuzione gestita

Kinesis Data Firehose è un servizio completamente gestito che acquisisce, trasforma e distribuisce dati in streaming verso destinazioni tra cui S3, Amazon Redshift, Amazon OpenSearch Service, Splunk ed endpoint HTTP. Non è necessario gestire shard: Firehose esegue automaticamente la scalabilità. Configuri una dimensione del buffer (1–128 MB) e un intervallo del buffer (60–900 secondi); Firehose distribuisce i dati non appena viene raggiunto per primo uno dei due limiti.

# Create a Firehose delivery stream to S3
aws firehose create-delivery-stream \
  --delivery-stream-name clickstream-to-s3 \
  --s3-destination-configuration '{
    "RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
    "BucketARN": "arn:aws:s3:::my-data-lake-123",
    "Prefix": "landing/clickstream/year=!{timestamp:yyyy}/month=!{timestamp:MM}/",
    "BufferingHints": {"SizeInMBs": 64, "IntervalInSeconds": 300},
    "CompressionFormat": "GZIP"
  }'

Trasformazione dei dati Firehose con Lambda

Firehose può richiamare una funzione Lambda per ogni batch di record prima della distribuzione, così da trasformare, arricchire o filtrare i dati durante il transito. Tra i casi d'uso comuni rientrano la conversione da JSON a Parquet (tramite lo schema Glue), il mascheramento dei campi contenenti informazioni personali e l'eliminazione degli eventi di scarso valore. I record per i quali la trasformazione non riesce possono essere inviati facoltativamente a un prefisso S3 separato per la rielaborazione, evitando così la perdita di dati.

# Lambda transform function signature for Firehose
def lambda_handler(event, context):
    output = []
    for record in event['records']:
        import base64, json
        payload = json.loads(base64.b64decode(record['data']))
        # Drop events with no user_id
        if not payload.get('user_id'):
            output.append({'recordId': record['recordId'], 'result': 'Dropped', 'data': record['data']})
        else:
            output.append({'recordId': record['recordId'], 'result': 'Ok', 'data': record['data']})
    return {'records': output}

Managed Service for Apache Flink

Amazon Managed Service for Apache Flink (in precedenza Kinesis Data Analytics) esegue applicazioni Apache Flink su un'infrastruttura completamente gestita. Lo utilizzi per analisi in tempo reale con stato: aggregazioni su finestre scorrevoli, rilevamento di anomalie, individuazione di modelli nelle sequenze di eventi e join tra dati in streaming e tabelle di riferimento. Scriva il codice Flink in Java, Python o Scala; Flink gestisce i checkpoint e lo stato con garanzia exactly-once.

# Flink SQL-style tumbling window (conceptual)
# Count page views per URL every 5 minutes
CREATE TABLE clickstream (
  url STRING,
  event_time TIMESTAMP(3),
  WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH ('connector' = 'kinesis', 'stream' = 'clickstream', ...);

SELECT
  url,
  COUNT(*) AS views,
  TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS window_start
FROM clickstream
GROUP BY url, TUMBLE(event_time, INTERVAL '5' MINUTE);

Scegliere tra Streams e Firehose

Per l'esame SAA-C03, è importante conoscere i criteri di scelta. Utilizzi Kinesis Data Streams quando sono necessarie una latenza inferiore al secondo, più consumer che leggono simultaneamente o una logica di elaborazione personalizzata con pieno controllo sulla conservazione e sulla rilettura dei dati. Utilizzi Firehose quando deve semplicemente distribuire in modo affidabile dati in streaming verso S3, Redshift o OpenSearch, con poco codice, trasformazione durante il transito opzionale e scalabilità automatica, accettando una latenza maggiore (60 o più secondi).

Kinesis o SQS: la scelta classica d'esame

Una domanda comune dell'esame SAA-C03 chiede di scegliere tra Kinesis e SQS. Differenze principali: Kinesis mantiene l'ordine all'interno di uno shard, supporta più consumer che leggono simultaneamente gli stessi dati e conserva i record per consentirne la rilettura. SQS rimuove i messaggi dopo il consumo (senza possibilità di rilettura), la modalità FIFO garantisce un ordinamento rigoroso ed è più adatta al disaccoppiamento dei microservizi. Se lo scenario menziona analisi in tempo reale o rilettura, scelga Kinesis.

Producer Kinesis: SDK e KPL

Kinesis Producer Library (KPL) è un client ad alte prestazioni per la scrittura in Kinesis Data Streams dalle applicazioni. KPL aggrega automaticamente più record di piccole dimensioni in una singola chiamata API (fino a 1 MB) e gestisce i tentativi ripetuti con back-off. Questo riduce notevolmente il costo PUT per record e aumenta il throughput per shard. Utilizzi KPL per produttori ad alto volume, come flussi di clic web, sensori IoT o pipeline di log.

# Basic KPL usage (Java pseudocode, no backticks)
KinesisProducer producer = new KinesisProducer();
byte[] data = 'hello world'.getBytes();
ByteBuffer buf = ByteBuffer.wrap(data);
// addUserRecord handles aggregation and retry internally
ListenableFuture future = producer.addUserRecord('clickstream', 'partitionKey', buf);
Futures.addCallback(future, new FutureCallback() { ... });

Modello Firehose verso Redshift

Un'architettura comune consiste nell'utilizzare Firehose come pipeline gestita dagli stream Kinesis o dai produttori direttamente verso Amazon Redshift per il data warehousing. Firehose scrive innanzitutto i dati in un bucket S3 intermedio di staging, quindi esegue un comando COPY per caricare i dati in Redshift. Questo è il modo più efficiente per caricare in blocco dati in streaming in Redshift: gli inserimenti diretti riga per riga in Redshift sarebbero estremamente lenti a causa dell'overhead per riga.

# Firehose Redshift destination (CLI snippet)
--redshift-destination-configuration '{
  "RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
  "ClusterJDBCURL": "jdbc:redshift://cluster.xyz.us-east-1.redshift.amazonaws.com:5439/sales",
  "CopyCommand": {
    "DataTableName": "clickevents",
    "CopyOptions": "JSON 'auto'"
  },
  "Username": "firehose_user",
  "Password": "{{resolve:secretsmanager:redshift-pw}}",
  "S3Configuration": {
    "RoleARN": "...",
    "BucketARN": "arn:aws:s3:::firehose-staging"
  }
}'

Monitoraggio di Kinesis con CloudWatch

Metriche CloudWatch principali per Kinesis: IncomingBytes e IncomingRecords per misurare il throughput dei produttori, GetRecords.IteratorAgeMilliseconds per misurare il ritardo dei consumer (un valore elevato indica che i consumer non riescono a tenere il passo) e WriteProvisionedThroughputExceeded per rilevare quando i produttori raggiungono i limiti degli shard. Imposti allarmi CloudWatch sull'età dell'iteratore e sul superamento del throughput assegnato per attivare la scalabilità automatica degli shard tramite Application Auto Scaling.

# CloudWatch alarm on high consumer lag
aws cloudwatch put-metric-alarm \
  --alarm-name kinesis-high-lag \
  --metric-name GetRecords.IteratorAgeMilliseconds \
  --namespace AWS/Kinesis \
  --dimensions Name=StreamName,Value=clickstream \
  --statistic Maximum \
  --period 60 \
  --threshold 60000 \
  --comparison-operator GreaterThanThreshold \
  --evaluation-periods 3 \
  --alarm-actions arn:aws:sns:us-east-1:123456789012:ops-alerts

Verifica rapida

Verifichi la Sua comprensione dei concetti di AWS Solutions Architect (SAA-C03) trattati in questa lezione.

Riepilogo della lezione

In questa lezione ha imparato che: Kinesis Data Streams fornisce uno streaming durevole, ordinato e suddiviso in shard, con più consumer e possibilità di rilettura, Kinesis Data Firehose offre una distribuzione completamente gestita e senza codice verso S3, Redshift e OpenSearch e Managed Service for Apache Flink consente analisi in tempo reale con stato sui flussi. Ora esploreremo i 7 R della strategia di migrazione al cloud.

Domande Frequenti

La lezione «Kinesis Streams, Firehose e analisi in tempo reale» è gratuita?

Sì — il testo completo di «Kinesis Streams, Firehose e analisi in tempo reale» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso Cloud & IT Cert Prep, passa a CoddyKit PRO. Il corso Cloud & IT Cert Prep include 4 lezioni in totale.

Cosa imparerò in «Kinesis Streams, Firehose e analisi in tempo reale»?

Acquisire dati in streaming con Kinesis Data Streams, distribuirli in S3 o Redshift con Firehose e analizzarli in tempo reale con Managed Service for Apache Flink Eserciti Cloud & IT Cert Prep con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.

Ho bisogno di esperienza per iniziare Cloud & IT Cert Prep?

Non è richiesta alcuna esperienza precedente. Cloud & IT Cert Prep su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 4 di 4.

Quanto tempo richiede la lezione «Kinesis Streams, Firehose e analisi in tempo reale»?

La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.

Posso scrivere ed eseguire codice in questa lezione Cloud & IT Cert Prep?

Sì. Ogni lezione Cloud & IT Cert Prep include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.

Tutte le lezioni di questo corso

  1. Creare un data lake su S3
  2. AWS Glue: ETL e catalogo dei dati
  3. Amazon Athena: SQL serverless su S3
  4. Kinesis Streams, Firehose e analisi in tempo reale
← Torna a Cloud & IT Cert Prep