Mapări de surse de evenimente și procesarea fluxurilor

Aplicații serverless cu AWS Lambda

Claudio Canales

Senior DevOps Engineer

Locul mapărilor de surse de evenimente

  • Cozile și fluxurile sunt interogate în loturi.
  • Mapările de surse de evenimente gestionează interogarea și comportamentul de grupare.

Poziția mapării sursei de evenimente

Aplicații serverless cu AWS Lambda

Ce este o mapare de sursă de evenimente?

  • O resursă care conectează o sursă la funcția ta.
  • Lambda interoghează sursa pentru înregistrări.
  • Lambda invocă handlerul tău cu un lot.

Fluxul mapării sursei de evenimente

Aplicații serverless cu AWS Lambda

Push vs poll (de ce contează)

Modelul push

  • Sursele push trimit evenimente imediat.
  • Exemplu: un upload în Amazon S3.

Modelul poll

  • Sursele poll sunt verificate de Lambda.
  • Lambda construiește un lot și invocă handlerul.
  • Aceasta afectează latența și comportamentul la reîncercare.

Modelele push vs poll

Aplicații serverless cu AWS Lambda

Coadă vs flux

Coadă (Amazon SQS)

  • O coadă de mesaje gestionată, ca o căsuță de intrare.

Flux (DynamoDB Streams)

  • Un jurnal de modificări pentru un tabel DynamoDB.
  • Ambele ajung ca Records, dar semnificația diferă.

Comparație coadă vs flux

Aplicații serverless cu AWS Lambda

Structura evenimentului de lot (simplificată)

{
  "Records": [{
    "messageId": "abc-123",
    "body": "{\"order_id\": \"A-42\"}"
  }]
}
  • Majoritatea evenimentelor de mapare încep cu Records.
  • Fiecare înregistrare are un messageId și un șir body.
  • Parsează body în propriul tău payload.
Aplicații serverless cu AWS Lambda

Exemplu: procesarea fiecărei înregistrări

def lambda_handler(event, context):
    records = event.get("Records", [])
    for record in records:
        body = record.get("body", "")
        print("BODY:", body)
    return {"statusCode": 200}
  • Citește Records cu o listă implicită.
  • Parcurge fiecare înregistrare și citește body în siguranță.
  • Loghează ce ai nevoie, apoi returnează.
Aplicații serverless cu AWS Lambda

Exemplu: parsarea body-ului JSON

import json

def lambda_handler(event, context):
    record = event.get("Records", [])[0]
    payload = json.loads(record.get("body", "{}"))
    order_id = payload.get("order_id")
    print("ORDER_ID:", order_id)
    return {"statusCode": 200}
  • Importă json și citește body cu o valoare implicită sigură.
  • Parsează cu json.loads pentru a obține un dicționar.
  • Extrage câmpurile necesare; validează întotdeauna mai întâi.
Aplicații serverless cu AWS Lambda

Grupare: dimensiunea și fereastra lotului

  • Dimensiunea lotului: câte înregistrări per invocare.
  • Fereastra de lot: cât așteaptă Lambda să completeze un lot.
  • Loturi mai mari îmbunătățesc debitul, dar adaugă mai mult lucru per rulare.
  • Pentru SQS, urmărește ApproximateAgeOfOldestMessage pentru acumulare.

Diagrama compromisului dimensiunii lotului

Aplicații serverless cu AWS Lambda

Compromisurile dimensiunii lotului

  • Loturi mari sunt eficiente, dar pot crește latența.
  • Loturi mici reduc timpul până la prima procesare, dar cresc numărul de invocări.
  • Alege în funcție de volumul de lucru.

Grafic dimensiune lot vs latență

Aplicații serverless cu AWS Lambda

Scalare prin concurență

  • Mai multe loturi sunt procesate simultan.
  • Concurența mai mare crește debitul.
  • Însă crește și încărcarea sistemelor din aval.

Procesarea concurentă a loturilor

Aplicații serverless cu AWS Lambda

Protejează sistemele din aval

  • Prea multă concurență poate suprasolicita o bază de date sau un API.
  • Folosește limite pentru a controla încărcarea.
  • Preferă unități de lucru mici și sigure.

Protecția sistemelor din aval

Aplicații serverless cu AWS Lambda

Eșecuri parțiale de lot

  • Un lot poate conține înregistrări valide și invalide.
  • Fără eșecuri parțiale, o înregistrare defectă reîncarcă întregul lot.
  • Gestionarea eșecurilor parțiale permite reîncercarea doar a elementelor eșuate.

Fluxul eșecurilor parțiale de lot

Aplicații serverless cu AWS Lambda

Exemplu: raportarea eșecurilor parțiale (SQS)

def lambda_handler(event, context):
    failures = []
    for record in event.get("Records", []):
        try:
            process(record)
        except Exception:
            failures.append({"itemIdentifier": record["messageId"]})
    return {"batchItemFailures": failures}
  • Colectează messageId-urile înregistrărilor eșuate în batchItemFailures.
  • Lambda reîncearcă doar acele elemente.
  • Înregistrările procesate cu succes nu sunt retrimise.
Aplicații serverless cu AWS Lambda

Concluzii cheie

  • Mapările de surse de evenimente interoghează cozi și fluxuri.
  • Livrează Records în loturi către handlerul tău.
  • Dimensiunea și fereastra lotului echilibrează latența și debitul.
  • Eșecurile parțiale evită reprocesarea elementelor reușite.

Diagrama sumară a mapărilor

Aplicații serverless cu AWS Lambda

Hai să exersăm!

Aplicații serverless cu AWS Lambda

Preparing Video For Download...