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.
- SNS: wysoka przepustowość fan-out, brak gwarancji kolejności, Publish/Subscribe, użyj MessageAttributes do routingu.
- SQS Standard: co najmniej jednokrotne dostarczenie (at-least-once), kolejność „na miarę możliwości” (best-effort), long polling, tańsze do separacji komponentów (decoupling).
- SQS FIFO: dokładnie jednokrotne przetwarzanie (exactly-once) i zachowanie kolejności w ramach MessageGroupId, używaj do ścisłego porządkowania i deduplikacji.
- Kinesis Data Streams: uporządkowane w ramach fragmentu (per-shard), strumieniowanie o wysokiej przepustowości, PutRecord/PutRecords, wymagane skalowanie fragmentów (shardów).
- Kinesis Firehose: zarządzane dostarczanie i buforowanie, obsługuje szyfrowanie po stronie serwera i transformacje Lambda.
- EventBridge: szyna zdarzeń (event bus) z regułami routingu, rejestr schematów, archiwizacja i odtwarzanie (archive & replay), API PutEvents.
- Step Functions: stanowa orkiestracja, ponowienia/Catch, kompromisy między Standard a Express w zakresie trwałości i przepustowości.
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:
- 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 (CreateQueuezFifoQueue=true), kluczując komunikaty poPlayerIDjakoMessageGroupIdi używającMessageDeduplicationIdopartego naETagz S3. - Skonfiguruj konsumenta Lambda z mapowaniem źródła zdarzeń SQS (
CreateEventSourceMapping), używającBatchSize=1,FunctionResponseTypes=["ReportBatchItemFailures"]i ustawVisibilityTimeoutna wartość większą niż maksymalny czas przetwarzania; używajChangeMessageVisibilitypodczas wywoływania zewnętrznych API do ML. - Funkcja Lambda wykonuje idempotentne zapisy do bazy danych Aurora, używając deterministycznego klucza idempotencji (
INSERT ... ON CONFLICT DO NOTHINGlub 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. - 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; monitorujIteratorAgeMillisecondsi używajUpdateShardCountdo 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 →