Date în flux

Dezvoltarea aplicațiilor pe AWS

Ricardo Sueiras

Principal Technologist

Noțiuni de bază despre datele în flux

 

date în flux

  • Înregistrările sunt procesate continuu, pe măsură ce sosesc.
  • Producătorii scriu într-un flux.
  • Consumatorii citesc independent și în paralel.
  • Absoarbe vârfurile de trafic fără a aștepta joburi batch.
Dezvoltarea aplicațiilor pe AWS

Din ce este format un flux de date?

 

componente sistem de streaming

  • Producători: aplicații care emit înregistrări.
  • Fluxuri: buffere durabile și ordonate care decuplează producătorii de consumatori.
  • Consumatori: aplicații care citesc și procesează înregistrări.
  • Stocare și retenție: definește cât timp rămân disponibile înregistrările pentru reluare.
Dezvoltarea aplicațiilor pe AWS

Sharduri

 

shards

 

  • Shardurile sunt unități de debit și procesare paralelă.
  • Cheile de partiție distribuie înregistrările între sharduri.
  • Ordinea este garantată doar în cadrul unui shard.
  • Consumatorii folosesc check-pointing pentru a urmări progresul și a recupera după erori.
Dezvoltarea aplicațiilor pe AWS

Servicii AWS gestionate pentru date în flux

  • Amazon Kinesis Data Streams.
  • Amazon Data Firehose.
  • Blocuri de bază pentru aplicațiile cu date în flux.

 

servicii aws

Dezvoltarea aplicațiilor pe AWS

Kinesis Data Streams

  • Streaming cu latență redusă.
  • Retenție configurabilă: 24 de ore (implicit) până la 365 de zile.
  • Fluxuri durabile, cu posibilitate de reluare.
  • Mai mulți consumatori pot citi în paralel.
  • Scalarea se gestionează manual prin sharduri.

 

kinesis data streams

Dezvoltarea aplicațiilor pe AWS

Când să folosești Kinesis Data Streams

  • Folosește-l când ai nevoie de latență în timp real.
  • Când este necesară reluarea înregistrărilor.
  • Când există mai mulți consumatori independenți.
  • Când ai nevoie de logică personalizată de transformare.

 

kinesis data streams

Dezvoltarea aplicațiilor pe AWS

Scrierea înregistrărilor

 

biblioteci client

  • Cheile de partiție controlează plasarea în shard și ordinea.
  • PutRecords grupează până la 500 de înregistrări sau 5 MB per cerere.
  • Erorile sunt raportate per înregistrare.
  • FailedRecordCount ajută la identificarea înregistrărilor eșuate.
Dezvoltarea aplicațiilor pe AWS

Consumul înregistrărilor: classic

 

classic

 

  • Consumatorii interoghează Kinesis API la cerere.
  • Limite de scriere: 1 MB/sec sau 1.000 de înregistrări/sec.
  • Limite de citire: 2 MB/sec sau 5 apeluri GetRecords/sec.
Dezvoltarea aplicațiilor pe AWS

Consumul înregistrărilor: punctul de start

 

stream-start

  • Consumatorii specifică de unde să înceapă citirea din flux.
  • TRIM_HORIZON: cea mai veche înregistrare disponibilă în shard.
  • LATEST: următoarea înregistrare nouă.
  • AT_SEQUENCE_NUMBER sau AFTER_SEQUENCE_NUMBER: de la un număr de secvență specific.
  • AT_TIMESTAMP: la sau după un timestamp specificat.
Dezvoltarea aplicațiilor pe AWS

Consumul înregistrărilor: enhanced fan-out

  • Fiecare consumator înregistrat primește 2 MB/sec dedicat per shard.
  • Înregistrările sunt livrate prin push HTTP/2.
  • Folosește această opțiune când mai mulți consumatori independenți au nevoie de latență redusă.
  • Și de debit ridicat pe același flux.

 

pattern - enhanced fan-out

Dezvoltarea aplicațiilor pe AWS

Maparea surselor de evenimente Lambda

  • Lambda interoghează Kinesis și invocă funcții în loturi.
  • BatchSize și MaximumBatchingWindow controlează dimensiunea lotului.
  • ParallelizationFactor: invocări concurente per shard.
  • BisectBatchOnFunctionError: împarte un lot cu erori pentru a izola înregistrările problematice.
  • ReportBatchItemFailures: evită reluarea înregistrărilor deja procesate cu succes.
  • OnFailure: trimite loturile neprocesabile într-un DLQ SQS sau SNS.

 

lambda event sourcing

Dezvoltarea aplicațiilor pe AWS

Data Firehose

  • Serviciu gestionat care livrează date în flux către destinații.
  • Destinații: Amazon S3, Redshift și endpoint-uri HTTP.
  • Suportă bufferizare, conversie de format și transformări Lambda.
  • Latență mai mare, date bufferizate de la secunde la minute.
  • Fără reluare sau retenție pe termen lung.
  • Ideal pentru livrarea simplă a datelor către o destinație.

 

data firehose

Dezvoltarea aplicațiilor pe AWS

Model: hot-cold

 

pattern - hot/cold

 

  • Modelul hot-cold combină Data Streams și Data Firehose.
  • Data Streams permite procesarea în timp real pentru producători.
  • Firehose arhivează înregistrările din aceeași sursă în S3.
Dezvoltarea aplicațiilor pe AWS

Kinesis vs SQS

 

sqs vs kinesis

 

  • Folosește Kinesis pentru fluxuri ordonate, cu reluare și mai mulți consumatori.
  • Kinesis este potrivit pentru date de tip time-series și date ordonate.
  • Folosește SQS pentru distribuirea sarcinilor.
Dezvoltarea aplicațiilor pe AWS

Scalarea datelor în flux

 

scalare

  • Fluxurile de date sunt împărțite în sharduri.
  • Cheile de partiție direcționează înregistrările către un shard specific.
  • Aceeași cheie ajunge întotdeauna în același shard, păstrând ordinea.
  • Cheile de partiție neuniforme duc la sharduri supraîncărcate.
  • Folosește chei cu cardinalitate mare pentru a preveni acest lucru.
  • Aplică resharding pentru a scala în sus sau în jos.
Dezvoltarea aplicațiilor pe AWS

Moduri de capacitate: provizionat

  • Sunt disponibile două moduri de capacitate.
  • Modul provizionat folosește un număr fix de sharduri, gestionat manual.
  • Oferă cost și capacitate previzibile.
  • Necesită monitorizare pentru a evita throttling-ul și supraprovizionarea.

 

moduri de capacitate

Dezvoltarea aplicațiilor pe AWS

Moduri de capacitate: on-demand

  • Modul on-demand scalează automat în funcție de trafic.
  • Nu necesită gestionarea manuală a shardurilor.
  • Simplifică operațiunile pentru sarcini de lucru imprevizibile.
  • Cost mai ridicat la trafic intens și susținut.

 

moduri de capacitate

Dezvoltarea aplicațiilor pe AWS

Kinesis Producer Library (KPL)

  • Bibliotecile AWS simplifică construirea de producători și consumatori scalabili.
  • Kinesis Producer Library (KPL) grupează și agregă înregistrările.
  • Comprimă înregistrările și reia automat în caz de eroare.
  • Bufferizează asincron pentru a maximiza debitul producătorului.

 

biblioteci client

Dezvoltarea aplicațiilor pe AWS

Kinesis Consumer Library (KCL)

 

biblioteci client

 

  • Kinesis Consumer Library (KCL) gestionează coordonarea shardurilor.
  • Administrează leasingul workerilor și check-pointing-ul.
  • Starea este stocată într-un tabel DynamoDB.
  • Consumatorii folosesc acea stare pentru a determina de unde să reia.
Dezvoltarea aplicațiilor pe AWS

Gestionarea erorilor: idempotență

 

gestionarea erorilor

 

  • Livrarea Kinesis este de tip at-least-once, deci duplicatele sunt posibile.
  • Proiectează consumatorii să fie idempotenți.
  • Deduplicarea se face folosind numărul de secvență sau o cheie de business în DynamoDB.
  • Consumatorii pot reciti înregistrările reținute de la un număr de secvență sau timestamp.
Dezvoltarea aplicațiilor pe AWS

Gestionarea erorilor: erori de debit

 

gestionarea erorilor

 

  • ProvisionedThroughputExceededException semnalează throttling.
  • Rezolvă-l prin creșterea numărului de sharduri.
  • Îmbunătățește distribuția cheilor de partiție.
  • Reia cu backoff.
  • Schimbă modul de capacitate.
Dezvoltarea aplicațiilor pe AWS
  • Urmărește IncomingBytes și IncomingRecords în CloudWatch.
  • Monitorizează WriteProvisionedThroughputExceeded pentru throttling.
  • Întârzierea consumatorului apare când citirea este mai lentă decât scrierea.
  • Urmărește GetRecords.IteratorAgeMilliseconds: o valoare mare indică întârziere.

 

monitorizare

Dezvoltarea aplicațiilor pe AWS

Securitate

  • Criptare server-side la repaus folosind AWS KMS.
  • Criptare în tranzit folosind TLS.
  • Politicile IAM controlează accesul.

 

securitate

Dezvoltarea aplicațiilor pe AWS

Hai să exersăm!

Dezvoltarea aplicațiilor pe AWS

Preparing Video For Download...