0Pricing
AWS Solutions Architect · Lezione

Kinesis Data Streams per l'elaborazione degli eventi in tempo reale

Produrre e consumare flussi di eventi ad alto throughput con Kinesis Data Streams, gestire gli shard per il throughput e utilizzare Lambda come consumer

Kinesis Data Streams per l'elaborazione degli eventi in tempo reale è una lezione AWS Solutions Architect gratuita su CoddyKit. Questa è la lezione 3 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 AWS Solutions Architect, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso AWS Solutions Architect include 4 lezioni in totale.

Concetti fondamentali di Kinesis Data Streams

Kinesis Data Streams (KDS) è un servizio durevole e ordinato per lo streaming di dati in tempo reale. I dati sono organizzati in uno stream composto da uno o più shard. Ogni shard è una sequenza ordinata di record di dati. I produttori inseriscono i record negli shard utilizzando una chiave di partizione che determina quale shard riceverà il record. I consumatori leggono i record dagli shard elaborandoli nell'ordine in cui sono arrivati all'interno di ciascuno shard.

Capacità degli shard e limiti di throughput

Ogni shard supporta un throughput di scrittura di 1 MB/s o 1.000 record/s e un throughput di lettura di 2 MB/s (condiviso tra tutti i consumatori standard di quello shard). La capacità totale dello stream aumenta linearmente con il numero di shard. Utilizzi la formula: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Se i produttori raggiungono il limite di scrittura, vengono visualizzati errori ProvisionedThroughputExceededException: per risolvere il problema, suddivida gli shard o distribuisca le chiavi di partizione in modo più uniforme.

# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed

aws kinesis create-stream \
  --stream-name iot-telemetry \
  --shard-count 5

Chiavi di partizione e distribuzione dei dati

La chiave di partizione è una stringa che Kinesis sottopone ad hashing (MD5) per determinare quale shard riceverà un record. Una chiave di partizione scelta correttamente distribuisce i record in modo uniforme tra gli shard (prevenzione degli shard sovraccarichi). Per l'IoT, utilizzi l'ID del dispositivo. Per i flussi di clic, utilizzi l'ID della sessione o l'ID dell'utente. Eviti le chiavi con bassa cardinalità (ad esempio, il nome del Paese con soli 5 valori), poiché causano shard sovraccarichi: uno shard riceve una quantità sproporzionata di traffico di scrittura mentre gli altri rimangono inattivi.

import boto3, json

client = boto3.client('kinesis', region_name='us-east-1')

# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
    StreamName='iot-telemetry',
    Data=json.dumps(event),
    PartitionKey='sensor-42'  # high-cardinality -> even distribution
)

Consumatori standard e Enhanced Fan-Out

I consumatori standard condividono il throughput di lettura di 2 MB/s per shard utilizzando GetRecords con il polling. Se ha 3 consumatori su uno shard e ciascuno necessita di 2 MB/s, subiranno il throttling condividendo un totale di 2 MB/s. Enhanced Fan-Out (EFO) assegna a ogni consumatore registrato un canale di lettura dedicato da 2 MB/s tramite una connessione persistente HTTP/2 con push (SubscribeToShard). EFO aggiunge un costo per consumatore-shard-ora, ma elimina completamente la contesa in lettura.

# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --consumer-name real-time-analytics

# List registered consumers
aws kinesis list-stream-consumers \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry

Lambda come consumatore di Kinesis

Lambda si integra nativamente con Kinesis Data Streams tramite un Event Source Mapping. Lambda esegue il polling dello stream, legge batch di record e invoca la Sua funzione. Configuri BatchSize (da 1 a 10.000 record), StartingPosition (TRIM_HORIZON per il record più vecchio, LATEST per quello più recente) e BisectBatchOnFunctionError per suddividere i batch non riusciti. Parallelisation Factor (da 1 a 10) consente a Lambda di avviare più invocazioni simultanee per shard, così da gestire stream che avanzano rapidamente.

# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
  --function-name ProcessIoTEvents \
  --event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --batch-size 100 \
  --starting-position LATEST \
  --parallelization-factor 5 \
  --bisect-batch-on-function-error true \
  --destination-config '{
    "OnFailure": {
      "Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
    }
  }'

Kinesis Client Library (KCL)

Kinesis Client Library (KCL) è un framework applicativo per creare consumatori Kinesis affidabili in Java (o in altri linguaggi tramite MultiLangDaemon). KCL gestisce l'enumerazione degli shard, la gestione dei lease (distribuendo gli shard tra le istanze worker), il salvataggio dei checkpoint dei progressi in DynamoDB e la gestione corretta della suddivisione e dell'unione degli shard. Ogni worker KCL elabora uno o più shard e KCL riequilibra automaticamente gli shard quando i worker aumentano o si verificano errori. KCL è preferibile al polling diretto tramite SDK per le applicazioni consumer in produzione.

# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
#   'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();

Conservazione e riproduzione dei dati

Kinesis Data Streams conserva i record per 24 ore per impostazione predefinita (estendibili a 7 giorni o fino a 365 giorni con la conservazione a lungo termine, con un costo aggiuntivo). A differenza di SQS, i record consumati non vengono eliminati dopo il consumo, ma rimangono disponibili fino alla scadenza della conservazione. Ciò consente a più consumatori di leggere gli stessi record in modo indipendente e permette la riproduzione reimpostando il checkpoint di un consumatore su un numero di sequenza precedente: una funzione preziosa per correggere bug o riempire i dati di nuovi servizi.

# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
  --stream-name iot-telemetry \
  --retention-period-hours 168

# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
  --stream-name iot-telemetry \
  --shard-id shardId-000000000000 \
  --shard-iterator-type TRIM_HORIZON \
  --query 'ShardIterator' --output text)

aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100

Modalità di capacità On-Demand

Kinesis Data Streams supporta due modalità di capacità. Modalità Provisioned: gestisce manualmente il numero di shard e paga per ogni shard-ora. Modalità On-Demand: Kinesis adatta automaticamente la capacità degli shard al throughput in ingresso (fino a 200 MB/s in scrittura e 400 MB/s in lettura per impostazione predefinita) e paga in base ai GB di dati scritti e recuperati. La modalità On-Demand è ideale per modelli di traffico variabili o imprevedibili, quando non desidera gestire il dimensionamento degli shard.

# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --stream-mode-details StreamMode=ON_DEMAND

Garanzie sull'ordinamento all'interno degli shard

Kinesis garantisce l'ordinamento all'interno di uno shard: i record con la stessa chiave di partizione vengono sempre indirizzati allo stesso shard e letti nell'ordine in cui sono stati scritti. Tuttavia, non vi è alcuna garanzia di ordinamento tra shard diversi. Se l'applicazione richiede un ordinamento globale di tutti i record, utilizzi un singolo shard (limitando il throughput a 1 MB/s) oppure la riprogetti in modo che l'ordinamento sia necessario solo all'interno di un gruppo con la stessa chiave di partizione (ad esempio, ordinamento per dispositivo). Questa è una distinzione comune negli esami rispetto a SQS FIFO, che fornisce deduplicazione e ordinamento rigorosi.

Sicurezza: crittografia e VPC

Kinesis Data Streams crittografa i record lato server quando sono inattivi utilizzando AWS KMS (una chiave CMK o gestita da AWS) quando la crittografia lato server è abilitata. Tutti i dati in transito sono crittografati con TLS. Per le applicazioni in esecuzione in una VPC che non devono inviare i dati dello stream tramite Internet pubblico, utilizzi un VPC Interface Endpoint (PrivateLink) per Kinesis, in modo che il traffico rimanga interamente all'interno della rete backbone AWS: un aspetto importante per i carichi di lavoro soggetti a normative.

# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
  --stream-name iot-telemetry \
  --encryption-type KMS \
  --key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123

Kinesis Data Streams, SQS e Kafka a confronto

Per l'esame SAA-C03, confronti Kinesis Data Streams con le alternative. Kinesis e SQS: Kinesis conserva l'ordine all'interno di uno shard e supporta più consumatori che leggono gli stessi dati; SQS elimina i messaggi dopo il consumo. Kinesis e MSK (Kafka): MSK è Apache Kafka gestito; lo utilizzi quando ha bisogno della compatibilità con il protocollo Kafka, di configurazioni avanzate degli argomenti o sta effettuando una migrazione da Kafka on-premises. Utilizzi Kinesis per lo streaming nativo AWS, con un'integrazione più stretta con Lambda, Firehose e Flink. Scelga Kinesis a meno che la domanda non menzioni esplicitamente Kafka o requisiti di compatibilità con Kafka.

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 stream suddivisi in shard, ordinati e durevoli, con conservazione predefinita di 24 ore e possibilità di riproduzione, Enhanced Fan-Out assegna a ogni consumatore 2 MB/s dedicati per shard, eliminando la contesa in lettura e la modalità On-Demand adatta automaticamente gli shard al traffico imprevedibile. Ora analizzeremo i pattern di choreography e orchestration nell'architettura basata sugli eventi.

Domande Frequenti

La lezione «Kinesis Data Streams per l'elaborazione degli eventi in tempo reale» è gratuita?

Sì — il testo completo di «Kinesis Data Streams per l'elaborazione degli eventi 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 AWS Solutions Architect, passa a CoddyKit PRO. Il corso AWS Solutions Architect include 4 lezioni in totale.

Cosa imparerò in «Kinesis Data Streams per l'elaborazione degli eventi in tempo reale»?

Produrre e consumare flussi di eventi ad alto throughput con Kinesis Data Streams, gestire gli shard per il throughput e utilizzare Lambda come consumer Eserciti AWS Solutions Architect 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 AWS Solutions Architect?

Non è richiesta alcuna esperienza precedente. AWS Solutions Architect su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 3 di 4.

Quanto tempo richiede la lezione «Kinesis Data Streams per l'elaborazione degli eventi 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 AWS Solutions Architect?

Sì. Ogni lezione AWS Solutions Architect 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. EventBridge: bus degli eventi e regole
  2. Step Functions: orchestrazione dei flussi di lavoro serverless
  3. Kinesis Data Streams per l'elaborazione degli eventi in tempo reale
  4. Pattern di choreography e orchestration
← Torna a AWS Solutions Architect