Mapowania źródeł zdarzeń i przetwarzanie strumieni

Serverless Applications with AWS Lambda

Claudio Canales

Senior DevOps Engineer

Gdzie mapowania źródeł zdarzeń pasują do układanki

  • Kolejki i strumienie są zazwyczaj odpytywane partiami.
  • Mapowania źródeł zdarzeń zarządzają odpytywaniem i tworzeniem partii.

Pozycja mapowania źródła zdarzeń

Serverless Applications with AWS Lambda

Czym jest mapowanie źródła zdarzeń?

  • Zasób łączący źródło z twoją funkcją.
  • Lambda odpytuje źródło w poszukiwaniu rekordów.
  • Lambda wywołuje handler z partią rekordów.

Przepływ mapowania źródła zdarzeń

Serverless Applications with AWS Lambda

Push vs. poll — dlaczego to ważne

Model push

  • Źródła push wysyłają zdarzenia natychmiast.
  • Przykład: przesłanie pliku do Amazon S3.

Model poll

  • Źródła poll są odpytywane przez Lambda.
  • Lambda buduje partię i wywołuje handler.
  • To wpływa na opóźnienia i sposób ponawiania prób.

Modele push i poll

Serverless Applications with AWS Lambda

Kolejka vs. strumień

Kolejka (Amazon SQS)

  • Zarządzana kolejka wiadomości, jak skrzynka odbiorcza.

Strumień (DynamoDB Streams)

  • Dziennik zmian tabeli DynamoDB.
  • Oba dostarczają rekordy (Records), ale mają inne znaczenie.

Porównanie kolejki i strumienia

Serverless Applications with AWS Lambda

Struktura zdarzenia partii (uproszczona)

{
  "Records": [{
    "messageId": "abc-123",
    "body": "{\"order_id\": \"A-42\"}"
  }]
}
  • Większość zdarzeń z mapowań zaczyna się od Records.
  • Każdy rekord ma messageId oraz ciąg body.
  • Sparsuj body do własnego ładunku.
Serverless Applications with AWS Lambda

Omówienie: przetwarzanie każdego rekordu

def lambda_handler(event, context):
    records = event.get("Records", [])
    for record in records:
        body = record.get("body", "")
        print("BODY:", body)
    return {"statusCode": 200}
  • Odczytaj Records z domyślną pustą listą.
  • Iteruj po rekordach i bezpiecznie odczytuj body.
  • Zaloguj potrzebne dane, a następnie zwróć odpowiedź.
Serverless Applications with AWS Lambda

Omówienie: parsowanie ciała 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}
  • Zaimportuj json i odczytaj body z bezpieczną wartością domyślną.
  • Użyj json.loads, aby uzyskać słownik.
  • Wyciągnij potrzebne pola; zawsze najpierw je zweryfikuj.
Serverless Applications with AWS Lambda

Partie: rozmiar i okno czasowe

  • Rozmiar partii: liczba rekordów na wywołanie.
  • Okno partii: jak długo Lambda czeka na wypełnienie partii.
  • Większe partie zwiększają przepustowość, ale wydłużają pracę na wywołanie.
  • W SQS monitoruj ApproximateAgeOfOldestMessage pod kątem zaległości.

Diagram kompromisów rozmiaru partii

Serverless Applications with AWS Lambda

Kompromisy rozmiaru partii

  • Duże partie są wydajne, ale mogą zwiększać opóźnienia.
  • Małe partie skracają czas do pierwszego przetworzenia, lecz mnożą wywołania.
  • Dobierz rozmiar do swojego obciążenia.

Wykres rozmiaru partii vs. opóźnienie

Serverless Applications with AWS Lambda

Skalowanie dzięki współbieżności

  • Wiele partii jest przetwarzanych równocześnie.
  • Wyższy poziom współbieżności zwiększa przepustowość.
  • Jednocześnie jednak rośnie obciążenie systemów downstream.

Współbieżne przetwarzanie partii

Serverless Applications with AWS Lambda

Chroń systemy downstream

  • Zbyt duża współbieżność może przeciążyć bazę danych lub API.
  • Używaj limitów, aby kontrolować obciążenie.
  • Wybieraj małe, bezpieczne jednostki pracy.

Ochrona systemów downstream

Serverless Applications with AWS Lambda

Częściowe błędy partii

  • Partia może zawierać dobre i złe rekordy.
  • Bez obsługi częściowych błędów jeden zły rekord powoduje ponowienie całej partii.
  • Obsługa częściowych błędów pozwala ponawiać tylko nieudane elementy.

Przepływ częściowego błędu partii

Serverless Applications with AWS Lambda

Omówienie: raportowanie częściowych błędów (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}
  • Zbieraj messageId nieudanych rekordów w batchItemFailures.
  • Lambda ponowi tylko te elementy.
  • Pomyślnie przetworzone rekordy nie są wysyłane ponownie.
Serverless Applications with AWS Lambda

Najważniejsze wnioski

  • Mapowania źródeł zdarzeń odpytują kolejki i strumienie.
  • Dostarczają rekordy w partiach do handlera.
  • Rozmiar i okno partii to kompromis między opóźnieniem a przepustowością.
  • Częściowe błędy zapobiegają ponownemu przetwarzaniu udanych rekordów.

Diagram podsumowania mapowań

Serverless Applications with AWS Lambda

Czas na praktykę!

Serverless Applications with AWS Lambda

Preparing Video For Download...