Amazon DVA-C02: Obsługa wiadomości, przesyłanie strumieniowe i architektury sterowane zdarzeniami (SNS, SQS, Kinesis, EventBridge, Step Functions) — Przewodnik do nauki

Część AWS Developer Associate DVA-C02 — Przewodnik do nauki. Ćwicz ze zweryfikowanymi odpowiedziami w centrum egzaminów Amazon, albo rozwiąż testy na czas na ExamRoll.io.

Wybór odpowiedniego prymitywu do przesyłania wiadomości i strumieniowania

Wybór pomiędzy SNS, SQS (standard vs FIFO), Kinesis, EventBridge i Step Functions zaczyna się od wzorca komunikacji: pub/sub, point-to-point, uporządkowane strumieniowanie, routing na szynie zdarzeń (event bus) czy orkiestracja przepływów pracy (workflow). SNS to wydawca pub/sub typu fan-out; użyj Publish (wywołanie SDK: Publish/PublishBatch) i subskrybuj punkty końcowe SQS, Lambda, HTTP/S lub punkty końcowe dla urządzeń mobilnych. SQS to trwały bufor point-to-point z semantyką ReceiveMessage/DeleteMessage; twórz kolejki za pomocą CreateQueue i atrybutów takich jak VisibilityTimeout, ReceiveMessageWaitTimeSeconds (long polling), MessageRetentionPeriod oraz polityk ponawiania (redrive policies) łączących z kolejką DLQ. Kolejki FIFO wymagają ustawienia FifoQueue=true i używają MessageGroupId oraz MessageDeduplicationId (lub ContentBasedDeduplication) do zachowania kolejności i deduplikacji. Kinesis Data Streams to uporządkowane strumieniowanie oparte na fragmentach (shardach); producenci wywołują PutRecord/PutRecords, a konsumenci używają GetShardIterator (TRIM_HORIZON, LATEST, AT_SEQUENCE_NUMBER), a następnie GetRecords. Kinesis Firehose zarządza dostarczaniem do S3/Redshift/OpenSearch i oferuje wskazówki buforowania (BufferingHints: SizeInMBs, IntervalInSeconds) oraz transformacje Lambda. EventBridge kieruje zdarzenia za pomocą PutEvents i filtrowania opartego na regułach, obsługuje rejestr schematów (schema registry) i szyny zdarzeń między kontami (cross-account buses). Step Functions orkiestruje złożone przepływy; użyj StartExecution (Standard) lub StartSyncExecution dla synchronicznych wzorców Express, z integracjami zadań (Task) takimi jak arn:aws:states:::lambda:invoke. Rozważ te kompromisy, gdy przepustowość, kolejność, gwarancje dostarczenia, retencja i potrzeby orkiestracji są ze sobą w konflikcie.

Wzorce SQS i SNS, deduplikacja i skalowanie konsumentów

Gdy potrzebujesz trwałej separacji komponentów (decoupling), SQS jest oczywistym wyborem; zaimplementuj SendMessage/SendMessageBatch dla producentów i używaj ReceiveMessage z WaitTimeSeconds, aby włączyć long polling i zredukować liczbę pustych odczytów. Dla ścisłego porządkowania i deduplikacji, utwórz kolejkę FIFO za pomocą CreateQueue (FifoQueue=true) i ustaw MessageGroupId dla uporządkowanych partycji; użyj MessageDeduplicationId lub włącz ContentBasedDeduplication, aby identyczne ładunki w oknie deduplikacji były pomijane. Kolejki standardowe mogą dostarczać duplikaty — dlatego konsumenci muszą być idempotentni, co można osiągnąć poprzez warunkowe zapisy do bazy danych (w DynamoDB PutItem z ConditionExpression attribute_not_exists(pk)) lub unikalne ograniczenia (unique constraints) i operacje upsert w transakcjach w RDS. Skonfiguruj polityki ponawiania (redrive policies), aby kierować niedostarczone wiadomości do kolejki DLQ po przekroczeniu maxReceiveCount; monitoruj ApproximateNumberOfMessages i ApproximateNumberOfMessagesNotVisible za pomocą GetQueueAttributes. Integracje z Lambda używają CreateEventSourceMapping dla SQS: ustaw BatchSize, MaximumBatchingWindowInSeconds i włącz FunctionResponseTypes = [“ReportBatchItemFailures”], aby używać semantyki częściowej odpowiedzi na paczkę (partial-batch-response) i unikać ponownego przetwarzania pomyślnie przetworzonych rekordów. Uważaj na semantykę Lambda z kolejkami FIFO: porządkowanie w ramach grupy wiadomości (message group) wymusza jednowątkowe przetwarzanie dla każdego MessageGroupId, ograniczając współbieżność w obrębie grupy; skaluj poprzez partycjonowanie na wiele identyfikatorów grup lub używając równoległych konsumentów z SNS do wielu kolejek. Pamiętaj również o visibility timeout: ustaw ChangeMessageVisibility, gdy przetwarzanie trwa dłużej, w przeciwnym razie ryzykujesz zduplikowane przetwarzanie.

Kinesis Data Streams i Firehose: porządkowanie, retencja i obsługa przeciwciśnienia (back-pressure)

Kinesis Data Streams zapewniają uporządkowanie w ramach fragmentu (per-shard) i trwałą retencję dla zastosowań strumieniowych. Producenci wywołują PutRecord lub PutRecords (wsadowo) z PartitionKey, który jest mapowany na fragment (shard); konsumenci wywołują GetShardIterator i GetRecords, a następnie zapisują punkty kontrolne (checkpoint) offsetów za pomocą KCL (Kinesis Client Library) lub niestandardowej tabeli DynamoDB do checkpointów. Domyślna retencja wynosi 24 godziny (można ją wydłużyć w konfiguracji strumienia, a tam gdzie dostępne, użyć funkcji rozszerzonej retencji); planuj liczbę fragmentów (shardów) za pomocą UpdateShardCount, aby dopasować ją do przepustowości zapisu i równoległości odczytu. Skalowanie konsumentów jest ograniczone: pojedyncze mapowanie źródła zdarzeń Lambda (event source mapping) mapuje jeden fragment (shard) na jedną jednostkę współbieżności Lambda, więc aby zwiększyć współbieżność konsumentów, zwiększ liczbę fragmentów lub włącz enhanced fan-out, aby dać każdemu konsumentowi własne pasmo 2 MB/s i niezależne skalowanie za pomocą API SubscribeToShard (rejestracja konsumenta). Używaj PutRecords do efektywnego przetwarzania wsadowego; przeciwciśnienie (back-pressure) pojawia się, gdy konsumenci mają opóźnienia (monitoruj GetRecords.IteratorAgeMilliseconds). Aby radzić sobie ze skokami obciążenia, buforuj w Kinesis lub użyj SQS jako bufora wejściowego, stosuj ponowienia z wykładniczym czasem oczekiwania (exponential backoff) po stronie producenta i używaj szyfrowania na poziomie strumienia za pomocą KMS dla danych PII. Kinesis Data Firehose upraszcza dostarczanie: skonfiguruj wskazówki buforowania (BufferingHints: SizeInMBs, IntervalInSeconds), CompressionFormat i transformację danych za pomocą Lambda. Firehose obsługuje ponowienia/backoff do miejsc docelowych i może zapisywać nieprzetworzone rekordy do zapasowego bucketa S3. Częstą pułapką jest niedostateczne alokowanie fragmentów (shardów): konsumenci „głodują”, a opóźnienia gwałtownie rosną; mierz i skaluj proaktywnie.

EventBridge i Step Functions do routingu i orkiestracji

EventBridge doskonale sprawdza się w routingu zdarzeń opartym na schematach oraz w integracjach międzykontowych i z partnerami zdarzeń, wykorzystując PutEvents do wstrzykiwania zdarzeń oraz PutRule/PutTargets do kierowania ich do SQS, Lambda, Kinesis, Step Functions lub punktów końcowych HTTP. EventBridge używa wzorców zdarzeń (event patterns) do filtrowania i wspiera archiwizację oraz odtwarzanie (replay) w celu odbudowy stanu. Używaj kolejek martwych listów (dead-letter queues) dla reguł (Target z SqsParameters lub DeadLetterConfig) i pamiętaj, że EventBridge zapewnia ponowne próby z wykładniczym czasem oczekiwania (exponential backoff), a w przypadku niepowodzenia kieruje zdarzenie do DLQ. Do orkiestracji wybierz Step Functions: maszyny stanów typu Standard dla długotrwałych, trwałych (durable) przepływów pracy z historią wykonań i wbudowanymi ponownymi próbami/Catch, oraz Express dla krótkotrwałych przepływów o wysokiej przepustowości, niższym koszcie i wykonaniu typu best-effort. Używaj integracji typu Task z integracjami usług (arn:aws:states:::lambda:invoke lub arn:aws:states:::aws-sdk:apigateway:invoke) oraz wzorców callback z użyciem "waitForTaskToken" do implementacji asynchronicznych zewnętrznych zatwierdzeń. Implementuj ponowne próby i Catch z wykładniczym czasem oczekiwania i używaj HeartbeatSeconds dla długich zadań. Używaj stanu Map do zrównoleglania przetwarzania dużych kolekcji, ale uważaj na współbieżność i throttling usług docelowych. Typową pułapką jest błędny wybór typu Express dla przepływów pracy, które wymagają trwałej historii z semantyką exactly-once — wybierz Standard dla zapewnienia audytowalności. Ponadto, zapewnij idempotencję w zadaniach wywoływanych przez Step Functions, przekazując token idempotencji i wymuszając unikalność po stronie usług docelowych podczas zapisu.

Problem praktyczny: Scenariusz użycia

Scenariusz: StreamlyGames zarządza globalnym backendem dla gier w wielokontowym środowisku AWS. Gracze przesyłają 10 MB klipy z rozgrywki do S3; potok przetwarzania musi transkodować wideo, przeprowadzać analizę ML i zapisywać wyniki do Aurora Serverless, zapewniając uporządkowane, zdeduplikowane przetwarzanie i skalowalnych konsumentów.

Wyzwanie: Zapewnij, aby każdy przesłany plik wyzwalał przetwarzanie dokładnie raz (exactly-once) w kolejności dla danego gracza, obsługuj skoki obciążenia bez utraty zdarzeń i skaluj konsumentów do inferencji ML, jednocześnie zapobiegając zduplikowanym zapisom do bazy danych.

Zalecane podejście:

  1. Utwórz powiadomienie o zdarzeniach S3, aby publikować zdarzenia utworzenia obiektu (object-created) na niestandardową magistralę EventBridge (custom bus) za pomocą PutEvents, a także do kolejki SQS FIFO (CreateQueue z FifoQueue=true), kluczując komunikaty po PlayerID jako MessageGroupId i używając MessageDeduplicationId opartego na ETag z S3.
  2. Skonfiguruj konsumenta Lambda z mapowaniem źródła zdarzeń SQS (CreateEventSourceMapping), używając BatchSize=1, FunctionResponseTypes=["ReportBatchItemFailures"] i ustaw VisibilityTimeout na wartość większą niż maksymalny czas przetwarzania; używaj ChangeMessageVisibility podczas wywoływania zewnętrznych API do ML.
  3. Funkcja Lambda wykonuje idempotentne zapisy do bazy danych Aurora, używając deterministycznego klucza idempotencji (INSERT ... ON CONFLICT DO NOTHING lub ograniczenie unikalności) i zapisuje punkty kontrolne postępu; dla długich wywołań ML użyj asynchronicznych Step Functions z tokenami zadań (arn:aws:states:::lambda:invoke.waitForTaskToken) lub Step Functions Express dla wysokiej przepustowości.
  4. Aby skalować inferencję, zapisuj zdarzenia pośrednie do Kinesis Data Streams per shard per region dla konsumentów o wysokiej przepustowości i włącz konsumentów z rozszerzonym rozproszeniem (enhanced fan-out) (SubscribeToShard) dla dedykowanych flot workerów ML; monitoruj IteratorAgeMilliseconds i używaj UpdateShardCount do skalowania.

Uzasadnienie: Użycie SQS FIFO gwarantuje kolejność per gracz i deduplikację na etapie pozyskiwania danych (ingestion), idempotentne zapisy do bazy danych wymuszają semantykę exactly-once, a Kinesis z enhanced fan-out lub Step Functions obsługują nagłe, wysokie obciążenia przetwarzania ML, utrzymując konsumentów oddzielonymi (decoupled) i skalowalnymi.


Bazy danych i buforowanie (RDS · Wszystkie domeny

Przećwicz te pytania → · Testy na czas na ExamRoll.io →

Pass the whole exam — not just this question

You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.

Zdaj egzamin →

Przeglądaj Amazon →

Related guides

Dostęp all-in-one

Jedna subskrypcja. Każdy egzamin.

Każdy plan odblokowuje nieograniczone wyszukiwanie odpowiedzi, testy praktyczne, wyjaśnienia AI i pełną bibliotekę zasobów — w ponad 20 językach.

Miesięczny
24.87
Just €0.83/day
Wszystko w cenie:
  • Nieograniczone wyszukiwanie odpowiedzi
  • Nieograniczone testy praktyczne
  • Wyjaśnienia wspomagane AI
  • Pełna biblioteka zasobów
  • Ponad 20 języków
  • Cotygodniowe aktualizacje treści
  • Nagrody i polecenia
  • Priorytetowe wsparcie
Rozpocznij bezpłatny okres próbny

Karta kredytowa nie jest wymagana*

Najlepsza wartość
12 miesięcy
179.87
Just €0.49/daySave 40%
Wszystko w cenie:
  • Nieograniczone wyszukiwanie odpowiedzi
  • Nieograniczone testy praktyczne
  • Wyjaśnienia wspomagane AI
  • Pełna biblioteka zasobów
  • Ponad 20 języków
  • Cotygodniowe aktualizacje treści
  • Nagrody i polecenia
  • Priorytetowe wsparcie
Rozpocznij bezpłatny okres próbny

Karta kredytowa nie jest wymagana*

✓ Plan darmowy w zestawie · ✓ Anuluj w dowolnym momencie · ✓ Wszystkie plany odblokowują pełny produkt