Dati in streaming

Sviluppare applicazioni su AWS

Ricardo Sueiras

Principal Technologist

Fondamenti dei dati in streaming

 

dati in streaming

  • I record sono elaborati in continuo, all'arrivo.
  • I produttori scrivono in uno stream.
  • I consumer leggono in modo indipendente e in parallelo.
  • Assorbe picchi di traffico senza attendere job batch.
Sviluppare applicazioni su AWS

Di cosa sono fatti i dati in streaming?

 

componenti di un sistema di streaming

  • Produttori: applicazioni che emettono record.
  • Stream: buffer durevoli e ordinati che disaccoppiano produttori e consumer.
  • Consumer: applicazioni che leggono ed elaborano record.
  • Conservazione: definisce per quanto i record restano rigiocabili.
Sviluppare applicazioni su AWS

Shard

 

shard

 

  • Gli shard sono unità di throughput e parallelismo.
  • Le chiavi di partizione distribuiscono i record tra shard.
  • L'ordine è garantito solo entro uno shard.
  • I consumer usano il check-pointing per tracciare l'avanzamento e recuperare errori.
Sviluppare applicazioni su AWS

Servizi AWS gestiti per dati in streaming

  • Amazon Kinesis Data Streams.
  • Amazon Data Firehose.
  • Mattoni fondamentali per applicazioni di streaming data.

 

servizi aws

Sviluppare applicazioni su AWS

Kinesis Data Streams

  • Streaming a bassa latenza.
  • Conservazione configurabile: 24 ore (default) fino a 365 giorni.
  • Stream durevoli e rigiocabili.
  • Più consumer leggono in parallelo.
  • Gli sviluppatori gestiscono lo scaling tramite shard.

 

kinesis data streams

Sviluppare applicazioni su AWS

Quando usare Kinesis Data Streams

  • Sceglilo quando serve latenza real-time.
  • Ti serve rigiocare i record.
  • Hai più consumer indipendenti.
  • Serve logica di trasformazione personalizzata.

 

kinesis data streams

Sviluppare applicazioni su AWS

Scrivere record

 

librerie client

  • Le chiavi di partizione controllano placement e ordine negli shard.
  • PutRecords fa batch fino a 500 record o 5 MB per richiesta.
  • Gli errori sono riportati per singolo record.
  • FailedRecordCount aiuta a individuare i record falliti.
Sviluppare applicazioni su AWS

Consumo dei record: classico

 

classico

 

  • I consumer fanno polling on demand sull'API Kinesis.
  • Limiti di scrittura: 1 MB/sec o 1.000 record/sec.
  • Limiti di lettura: 2 MB/sec o 5 chiamate GetRecords/sec.
Sviluppare applicazioni su AWS

Consumo dei record: inizio dello stream

 

inizio stream

  • I consumer indicano da dove iniziare a leggere nello stream.
  • TRIM_HORIZON: il record più vecchio nello shard.
  • LATEST: il prossimo record nuovo.
  • AT_SEQUENCE_NUMBER o AFTER_SEQUENCE_NUMBER: da un sequence number specifico.
  • AT_TIMESTAMP: a o dopo un timestamp specificato.
Sviluppare applicazioni su AWS

Consumo dei record: enhanced fan-out

  • Ogni consumer registrato ha 2 MB/sec dedicati per shard.
  • I record sono consegnati via push HTTP/2.
  • Usalo quando più consumer indipendenti richiedono bassa latenza.
  • E alto throughput sullo stesso stream.

 

pattern - enhanced fan-out

Sviluppare applicazioni su AWS

Mapping sorgente eventi Lambda

  • Lambda fa polling su Kinesis e invoca a batch.
  • BatchSize e MaximumBatchingWindow definiscono forma e dimensione dei batch.
  • ParallelizationFactor: invocazioni concorrenti per shard.
  • BisectBatchOnFunctionError: divide un batch in errore per isolare i record problematici.
  • ReportBatchItemFailures: evita retry di record già riusciti.
  • OnFailure: invia batch non processabili a una DLQ SQS o SNS.

 

lambda event sourcing

Sviluppare applicazioni su AWS

Data Firehose

  • Servizio gestito che consegna dati in streaming a destinazioni.
  • Target includono Amazon S3, Redshift e endpoint HTTP.
  • Supporta buffering, conversione formato e trasformazioni Lambda.
  • Latenza più alta, buffer di secondi/minuti.
  • Nessun replay o conservazione a lungo termine.
  • Ottimo per recapitare dati a una destinazione in modo semplice.

 

data firehose

Sviluppare applicazioni su AWS

Pattern: hot-cold

 

pattern - hot/cold

 

  • Il pattern hot-cold combina Data Streams e Data Firehose.
  • Data Streams abilita l'elaborazione real-time per i produttori.
  • Firehose archivia gli stessi record sorgente in S3.
Sviluppare applicazioni su AWS

Kinesis vs SQS

 

sqs vs kinesis

 

  • Usa Kinesis per stream ordinati, rigiocabili, con più consumer.
  • Kinesis è adatto a dati time-series e ordinati.
  • Usa SQS per distribuire task.
Sviluppare applicazioni su AWS

Scalare i dati in streaming

 

scaling

  • I data stream sono divisi in shard.
  • Le chiavi di partizione instradano i record a uno shard specifico.
  • La stessa chiave atterra sempre sullo stesso shard, preservando l'ordine.
  • Chiavi sbilanciate creano shard «caldi».
  • Mitiga con chiavi ad alta cardinalità.
  • Esegui reshard per scalare out/in.
Sviluppare applicazioni su AWS

Modalità di capacità: provisioned

  • Disponibili due modalità di capacità.
  • La modalità provisioned usa un numero fisso di shard, gestiti manualmente.
  • Offre costi e capacità prevedibili.
  • Richiede monitoraggio per evitare throttling e over-provisioning.

 

modalità di capacità

Sviluppare applicazioni su AWS

Modalità di capacità: on-demand

  • La modalità on-demand scala automaticamente in base al traffico.
  • Nessuna gestione manuale degli shard.
  • Semplifica le operazioni con carichi imprevedibili.
  • Costo più alto con traffico elevato e prolungato.

 

modalità di capacità

Sviluppare applicazioni su AWS

Kinesis Producer Library (KPL)

  • Le librerie AWS semplificano produttori e consumer scalabili.
  • La Kinesis Producer Library (KPL) fa batching e aggrega record.
  • Comprimi i record e ritenta in caso di errore.
  • Buffer asincrono per massimizzare il throughput del produttore.

 

librerie client

Sviluppare applicazioni su AWS

Kinesis Consumer Library (KCL)

 

librerie client

 

  • La Kinesis Consumer Library (KCL) gestisce il coordinamento degli shard.
  • Gestisce leasing dei worker e check-pointing.
  • Lo stato è salvato in una tabella DynamoDB.
  • I consumer usano quello stato per sapere da dove riprendere.
Sviluppare applicazioni su AWS

Gestione errori: idempotenza

 

gestione errori

 

  • La consegna Kinesis è at-least-once, quindi sono possibili duplicati.
  • Progetta i consumer per essere idempotenti.
  • Deduplica usando il sequence number o una chiave di business in DynamoDB.
  • I consumer possono rileggere record conservati da un sequence number o timestamp.
Sviluppare applicazioni su AWS

Gestione errori: throughput

 

gestione errori

 

  • ProvisionedThroughputExceededException segnala throttling.
  • Risolvi aumentando il numero di shard.
  • Migliora la distribuzione delle chiavi di partizione.
  • Ritenta con backoff.
  • Cambia modalità di capacità.
Sviluppare applicazioni su AWS
  • Traccia IncomingBytes e IncomingRecords in CloudWatch.
  • Monitora WriteProvisionedThroughputExceeded per il throttling.
  • Si ha lag quando il consumo è più lento della produzione.
  • Traccia GetRecords.IteratorAgeMilliseconds: età alta indica lag.

 

logging e monitoraggio

Sviluppare applicazioni su AWS

Sicurezza

  • Crittografia lato server a riposo con AWS KMS.
  • Crittografia in transito con TLS.
  • Le policy IAM controllano l'accesso.

 

sicurezza

Sviluppare applicazioni su AWS

Esercitiamoci!

Sviluppare applicazioni su AWS

Preparing Video For Download...