Amazon DVA-C02: Messaging, Streaming und ereignisgesteuerte Architekturen (SNS, SQS, Kinesis, EventBridge, Step Functions) — Lernleitfaden
Teil des AWS Developer Associate DVA-C02 — Lernleitfaden. Üben Sie mit verifizierten Antworten im Amazon-Prüfungscenter, oder absolvieren Sie zeitlich begrenzte Übungstests auf ExamRoll.io.
Auswahl des richtigen Messaging- und Streaming-Grundbausteins
Die Wahl zwischen SNS, SQS (Standard vs. FIFO), Kinesis, EventBridge und Step Functions beginnt mit dem Kommunikationsmuster: Pub/Sub, Point-to-Point, geordnetes Streaming, Event-Bus-Routing oder Workflow-Orchestrierung. SNS ist ein Fan-Out-Pub/Sub-Publisher; verwenden Sie Publish (SDK-Aufruf: Publish/PublishBatch) und abonnieren Sie SQS-Endpunkte, Lambda, HTTP/S oder mobile Endpunkte. SQS ist ein langlebiger Point-to-Point-Puffer mit ReceiveMessage/DeleteMessage-Semantik; erstellen Sie Warteschlangen mit CreateQueue und Attributen wie VisibilityTimeout, ReceiveMessageWaitTimeSeconds (Long Polling), MessageRetentionPeriod und Redrive-Richtlinien, die eine DLQ verknüpfen. FIFO-Warteschlangen erfordern FifoQueue=true und verwenden MessageGroupId sowie MessageDeduplicationId (oder ContentBasedDeduplication) für Reihenfolge und Deduplizierung. Kinesis Data Streams ist geordnetes, Shard-basiertes Streaming; Producer rufen PutRecord/PutRecords auf und Consumer verwenden GetShardIterator (TRIM_HORIZON, LATEST, AT_SEQUENCE_NUMBER) und dann GetRecords. Kinesis Firehose verwaltet die Zustellung an S3/Redshift/OpenSearch und bietet BufferingHints (SizeInMBs, IntervalInSeconds) sowie Lambda-Transformationen. EventBridge leitet Ereignisse mit PutEvents und regelbasiertem Filtern weiter, unterstützt eine Schema-Registry und kontoübergreifende Busse. Step Functions orchestrieren komplexe Abläufe; StartExecution (Standard) oder StartSyncExecution für synchrone Express-Muster, mit Task-Integrationen wie arn:aws:states:::lambda:invoke. Berücksichtigen Sie diese Kompromisse, wenn Durchsatz, Reihenfolge, Zustellgarantien, Aufbewahrung und Orchestrierungsanforderungen im Widerspruch zueinander stehen.
- SNS: Fan-Out mit hohem Durchsatz, keine Reihenfolge, Publish/Subscribe, MessageAttributes für das Routing verwenden.
- SQS Standard: Mindestens einmalige Zustellung (At-least-once), Best-Effort-Reihenfolge, Long Polling, günstiger zur Entkopplung.
- SQS FIFO: Genau einmalige Verarbeitung (Exactly-once) innerhalb einer MessageGroupId, für strikte Reihenfolge und Deduplizierung verwenden.
- Kinesis Data Streams: Geordnet pro Shard, Streaming mit hohem Durchsatz, PutRecord/PutRecords, Skalierung der Shards erforderlich.
- Kinesis Firehose: Verwaltete Zustellung und Pufferung, unterstützt serverseitige Verschlüsselung und Lambda-Transformationen.
- EventBridge: Event-Bus mit Routing-Regeln, Schema-Registry, Archivierung & Wiederholung, PutEvents-API.
- Step Functions: Zustandsbehaftete Orchestrierung, Wiederholungsversuche/Catch, Kompromisse zwischen Standard und Express bei Langlebigkeit und Durchsatz.
SQS- und SNS-Muster, Deduplizierung und Skalierung von Consumern
Wenn Sie eine langlebige Entkopplung benötigen, ist SQS die erste Wahl; implementieren Sie SendMessage/SendMessageBatch für Producer und verwenden Sie ReceiveMessage mit WaitTimeSeconds, um Long Polling zu aktivieren und leere Empfangsvorgänge (Empty Receives) zu reduzieren. Für eine strikte Reihenfolge und Deduplizierung erstellen Sie eine FIFO-Warteschlange mit CreateQueue (FifoQueue=true) und setzen Sie eine MessageGroupId für geordnete Partitionen; verwenden Sie eine MessageDeduplicationId oder aktivieren Sie ContentBasedDeduplication, damit identische Payloads innerhalb des Deduplizierungsfensters unterdrückt werden. Standard-Warteschlangen können Duplikate zustellen – daher sollten Consumer idempotent gestaltet werden, z. B. durch bedingte Schreibvorgänge in der Datenbank (DynamoDB PutItem mit ConditionExpression attribute_not_exists(pk)) oder durch Unique Constraints und Upserts innerhalb von Transaktionen in RDS. Konfigurieren Sie Redrive-Richtlinien, um fehlerhafte Nachrichten nach Erreichen von maxReceiveCount an eine DLQ weiterzuleiten; überwachen Sie ApproximateNumberOfMessages und ApproximateNumberOfMessagesNotVisible über GetQueueAttributes. Lambda-Integrationen verwenden CreateEventSourceMapping für SQS: Setzen Sie BatchSize, MaximumBatchingWindowInSeconds und aktivieren Sie FunctionResponseTypes = [“ReportBatchItemFailures”], um die Semantik für Teil-Batch-Antworten (Partial-Batch-Response) zu nutzen und die erneute Verarbeitung bereits erfolgreicher Datensätze zu vermeiden. Beachten Sie die Semantik von FIFO-Lambda: Die Reihenfolge der Nachrichtengruppen erzwingt eine Single-Thread-Verarbeitung pro MessageGroupId, was die gleichzeitige Verarbeitung pro Gruppe einschränkt; skalieren Sie durch Partitionierung in viele Gruppen-IDs oder durch den Einsatz paralleler Consumer mit SNS an mehrere Warteschlangen. Achten Sie auch auf das Visibility Timeout: Setzen Sie ChangeMessageVisibility, wenn die Verarbeitung länger dauert, sonst riskieren Sie eine doppelte Verarbeitung.
Kinesis Data Streams und Firehose: Reihenfolge, Aufbewahrung und Umgang mit Gegendruck
Kinesis Data Streams bieten Reihenfolgenerhaltung pro Shard und langlebige Aufbewahrung für Streaming-Anwendungsfälle. Producer rufen PutRecord oder PutRecords (Batch) mit einem PartitionKey auf, der einem Shard zugeordnet wird; Consumer rufen GetShardIterator und GetRecords auf und setzen dann Checkpoint-Offsets mit der KCL (Kinesis Client Library) oder einer benutzerdefinierten DynamoDB-Checkpoint-Tabelle. Die standardmäßige Aufbewahrungsfrist beträgt 24 Stunden (kann pro Stream-Konfiguration auf längere Zeiträume angepasst werden, und erweiterte Aufbewahrungsfunktionen, wo verfügbar); planen Sie die Anzahl der Shards mit UpdateShardCount, um sie an den Schreibdurchsatz und die Parallelität der Lesevorgänge anzupassen. Die Skalierung der Consumer ist eingeschränkt: Ein einzelnes Lambda Event Source Mapping ordnet einen Shard einer Lambda-Concurrency zu. Erhöhen Sie also die Anzahl der Shards, um die Consumer-Concurrency zu erhöhen, oder aktivieren Sie Enhanced Fan-Out, um jedem Consumer seine eigene 2 MB/s-Pipe und unabhängige Skalierung mithilfe der SubscribeToShard-API (Consumer-Registrierung) zu geben. Verwenden Sie PutRecords für effizientes Batching; Gegendruck (Back-Pressure) entsteht, wenn die Consumer hinterherhinken (überwachen Sie GetRecords.IteratorAgeMilliseconds). Um Lastspitzen zu bewältigen, puffern Sie auf Kinesis oder schalten Sie SQS vor, verwenden Sie Wiederholungsversuche des Producers mit exponentiellem Backoff und nutzen Sie Verschlüsselung auf Stream-Ebene mit KMS für personenbezogene Daten (PII). Kinesis Data Firehose vereinfacht die Zustellung: Konfigurieren Sie BufferingHints (SizeInMBs, IntervalInSeconds), CompressionFormat und eine Lambda-Datentransformation. Firehose kümmert sich um Wiederholungsversuche/Backoff zu den Zielen und kann fehlgeschlagene Datensätze in einen Backup-S3-Bucket schreiben. Eine häufige Falle ist die Unterprovisionierung von Shards: Consumer hungern aus und die Latenz steigt sprunghaft an; messen und skalieren Sie proaktiv.
EventBridge und Step Functions für Routing und Orchestrierung
EventBridge eignet sich hervorragend für schema-gesteuertes Event-Routing und kontoübergreifende/Event-Partner-Integrationen, indem es PutEvents zum Einspeisen von Events und PutRule/PutTargets zum Routen an SQS, Lambda, Kinesis, Step Functions oder HTTP-Endpunkte verwendet. EventBridge verwendet Event Patterns zum Filtern und unterstützt Archivierung und Replay, um den Zustand wiederherzustellen. Verwenden Sie Dead-Letter-Queues für Regeln (Target mit SqsParameters oder DeadLetterConfig) und beachten Sie, dass EventBridge bei einem Fehler Wiederholungsversuche mit exponentiellem Backoff und anschließend eine Weiterleitung an die DLQ bietet. Für die Orchestrierung wählen Sie Step Functions: Standard State Machines für langlebige, dauerhafte Workflows mit Ausführungsverlauf und integrierten Retries/Catch, und Express für kurzlebige Workflows mit hohem Durchsatz, geringeren Kosten und Best-Effort-Ausführung. Nutzen Sie Task-Integrationen mit Service-Integrationen (arn:aws:states:::lambda:invoke oder arn:aws:states:::aws-sdk:apigateway:invoke) und Callback Patterns mit „waitForTaskToken“, um asynchrone externe Genehmigungen zu implementieren. Implementieren Sie Retries und Catch mit exponentiellem Backoff und verwenden Sie HeartbeatSeconds für lang andauernde Tasks. Verwenden Sie den Map-Zustand, um große Sammlungen zu parallelisieren, aber achten Sie auf Gleichzeitigkeit (Concurrency) und nachgelagerte Drosselungen (Throttles). Ein typischer Fallstrick ist die falsche Wahl von Express für Workflows, die einen dauerhaften, genau einmaligen Verlauf (exactly-once durable history) erfordern – wählen Sie Standard für die Auditierbarkeit. Stellen Sie außerdem die Idempotenz bei von Step Functions aufgerufenen Tasks sicher, indem Sie ein Idempotency Token übergeben und die Zieldienste die Eindeutigkeit beim Schreiben erzwingen lassen.
Praktisches Problem: Anwendungsfallszenario
Szenario: StreamlyGames betreibt ein globales Gaming-Backend in einer AWS-Umgebung mit mehreren Konten. Spieler laden 10-MB-Gameplay-Clips auf S3 hoch; eine Verarbeitungspipeline muss Videos transkodieren, ML-Analysen durchführen und die Ergebnisse mit geordneter, deduplizierter Verarbeitung und skalierbaren Consumern in Aurora Serverless schreiben.
Herausforderung: Sicherstellen, dass jede hochgeladene Datei genau eine Verarbeitung pro Spieler in der richtigen Reihenfolge auslöst, Lastspitzen ohne Event-Verlust bewältigt werden und die Consumer für die ML-Inferenz skaliert werden, während doppelte DB-Schreibvorgänge verhindert werden.
Empfohlener Ansatz:
- Erstellen Sie eine S3-Event-Benachrichtigung, um „Object-Created“-Events an einen benutzerdefinierten EventBridge-Bus (
PutEvents) und auch an eine SQS-FIFO-Warteschlange (CreateQueuemitFifoQueue=true) zu veröffentlichen. Verwenden Sie diePlayerIDalsMessageGroupIdund eine auf dem S3-ETag basierendeMessageDeduplicationId. - Konfigurieren Sie einen Lambda-Consumer mit einem SQS Event Source Mapping (
CreateEventSourceMapping) unter Verwendung vonBatchSize=1,FunctionResponseTypes=["ReportBatchItemFailures"]und setzen SieVisibilityTimeout> maximale Verarbeitungszeit; verwenden SieChangeMessageVisibility, wenn Sie ML-APIs von Drittanbietern aufrufen. - Die Lambda-Funktion führt idempotente DB-Schreibvorgänge in Aurora durch, indem sie einen deterministischen Idempotenzschlüssel (
INSERT ... ON CONFLICT DO NOTHINGoder eine Unique Constraint) verwendet und den Fortschritt speichert (Checkpoints); für lange ML-Aufrufe verwenden Sie asynchrone Step Functions mit Task-Token (arn:aws:states:::lambda:invoke.waitForTaskToken) oder Step Functions Express für hohen Durchsatz. - Um die Inferenz zu skalieren, schreiben Sie Zwischenereignisse in Kinesis Data Streams pro Shard pro Region für Hochdurchsatz-Consumer und aktivieren Sie Enhanced-Fan-Out-Consumer (
SubscribeToShard) für dedizierte ML-Worker-Flotten; überwachen SieIteratorAgeMillisecondsund verwenden SieUpdateShardCountzur Skalierung.
Begründung: Verwenden Sie SQS FIFO, um die Reihenfolge pro Spieler und die Deduplizierung bei der Erfassung zu garantieren, idempotente DB-Schreibvorgänge, um eine Exactly-Once-Semantik zu erzwingen, und Kinesis plus Enhanced Fan-Out oder Step Functions, um stoßweise, hochvolumige ML-Verarbeitung zu bewältigen, während die Consumer entkoppelt und skalierbar bleiben.
← Datenbanken und Caching (RDS · Alle Domänen
Diese Fragen üben → · Zeitlich begrenzte Übung auf 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.
Bestehe deine Prüfung →