Amazon DVA-C02: Mensajería, streaming y arquitecturas orientadas a eventos (SNS, SQS, Kinesis, EventBridge, Step Functions) — Guía de estudio

Forma parte de la AWS Developer Associate DVA-C02 — Guía de estudio. Practica con respuestas verificadas en el centro de exámenes de Amazon, o realiza tests cronometrados en ExamRoll.io.

Elegir la primitiva de mensajería y streaming adecuada

La elección entre SNS, SQS (estándar vs. FIFO), Kinesis, EventBridge y Step Functions comienza con el patrón de comunicación: pub/sub, punto a punto, streaming ordenado, enrutamiento de bus de eventos u orquestación de flujos de trabajo. SNS es un publicador pub/sub de tipo fan-out; utiliza Publish (llamada SDK: Publish/PublishBatch) y suscribe puntos de conexión de SQS, Lambda, HTTP/S o móviles. SQS es un búfer duradero punto a punto con semántica ReceiveMessage/DeleteMessage; crea colas con CreateQueue y atributos como VisibilityTimeout, ReceiveMessageWaitTimeSeconds (sondeo largo), MessageRetentionPeriod y políticas de reenvío (redrive policies) que enlazan a una DLQ. Las colas FIFO requieren FifoQueue=true y usan MessageGroupId más MessageDeduplicationId (o ContentBasedDeduplication) para el ordenamiento y la deduplicación. Kinesis Data Streams es un servicio de streaming ordenado y basado en shards; los productores llaman a PutRecord/PutRecords y los consumidores usan GetShardIterator (TRIM_HORIZON, LATEST, AT_SEQUENCE_NUMBER) y luego GetRecords. Kinesis Firehose gestiona la entrega a S3/Redshift/OpenSearch y ofrece BufferingHints (SizeInMBs, IntervalInSeconds) y transformaciones con Lambda. EventBridge enruta eventos con PutEvents y filtrado basado en reglas, admite el registro de esquemas (schema registry) y buses entre cuentas. Step Functions orquesta flujos complejos; StartExecution (Standard) o StartSyncExecution para patrones síncronos Express, con integraciones de Task como

undefined

. Considera estas compensaciones (tradeoffs) cuando el rendimiento (throughput), el ordenamiento, las garantías de entrega, la retención y las necesidades de orquestación entran en conflicto.

Patrones de SQS y SNS, deduplicación y escalado de consumidores

Cuando necesitas un desacoplamiento duradero, SQS es la opción preferida; implementa SendMessage/SendMessageBatch para los productores y usa ReceiveMessage con WaitTimeSeconds para habilitar el sondeo largo y reducir las recepciones vacías. Para un ordenamiento y deduplicación estrictos, crea una cola FIFO con CreateQueue (FifoQueue=true) y establece el MessageGroupId para particiones ordenadas; usa MessageDeduplicationId o habilita ContentBasedDeduplication para que las cargas útiles (payloads) idénticas dentro de la ventana de deduplicación sean suprimidas. Las colas estándar pueden entregar duplicados; por lo tanto, haz que los consumidores sean idempotentes mediante escrituras condicionales en la base de datos (PutItem en DynamoDB con ConditionExpression attribute_not_exists(pk)) o restricciones de unicidad (unique constraints) y operaciones de upsert en RDS dentro de transacciones. Configura políticas de reenvío (redrive policies) para enrutar mensajes fallidos a una DLQ después de un maxReceiveCount; monitorea ApproximateNumberOfMessages y ApproximateNumberOfMessagesNotVisible a través de GetQueueAttributes. Las integraciones con Lambda usan CreateEventSourceMapping para SQS: establece BatchSize, MaximumBatchingWindowInSeconds y habilita FunctionResponseTypes = [“ReportBatchItemFailures”] para usar la semántica de respuesta parcial de lote (partial-batch-response) y evitar el reprocesamiento de registros que ya tuvieron éxito. Ten cuidado con la semántica de Lambda para FIFO: el ordenamiento del grupo de mensajes impone un procesamiento de un solo hilo (single-threaded) por MessageGroupId, lo que limita la concurrencia por grupo; escala particionando en muchos ID de grupo o usando consumidores paralelos con SNS hacia múltiples colas. También ten en cuenta el tiempo de espera de visibilidad (visibility timeout): establece ChangeMessageVisibility cuando el procesamiento tarde más tiempo o corras el riesgo de un procesamiento duplicado.

Kinesis Data Streams y Firehose: ordenamiento, retención y manejo de contrapresión (back-pressure)

Kinesis Data Streams proporciona ordenamiento por shard y retención duradera para casos de uso de streaming. Los productores llaman a PutRecord o PutRecords (en lote) con una PartitionKey que se asigna a un shard; los consumidores llaman a GetShardIterator y GetRecords, y luego guardan los desplazamientos (offsets) como puntos de control (checkpoint) usando KCL (Kinesis Client Library) o una tabla de checkpoint personalizada en DynamoDB. La retención predeterminada es de 24 horas (ajustable a ventanas más largas según la configuración del stream, y con características de retención extendida donde estén disponibles); planifica el número de shards con UpdateShardCount para que coincida con el rendimiento de escritura y el paralelismo de lectura. El escalado de consumidores está restringido: un único mapeo de origen de eventos de Lambda (event source mapping) asigna un shard a una concurrencia de Lambda, así que aumenta los shards para aumentar la concurrencia de consumidores o habilita el fan-out mejorado (enhanced fan-out) para dar a cada consumidor su propio canal de 2 MB/s y escalado independiente usando la API SubscribeToShard (registro de consumidor). Usa PutRecords para un procesamiento por lotes eficiente; la contrapresión (back-pressure) surge cuando los consumidores se retrasan (monitorea GetRecords.IteratorAgeMilliseconds). Para manejar picos, almacena en búfer en Kinesis o antepón SQS, usa reintentos en el productor con retroceso exponencial (exponential backoff) y utiliza cifrado a nivel de stream con KMS para PII (información de identificación personal). Kinesis Data Firehose simplifica la entrega: configura BufferingHints (SizeInMBs, IntervalInSeconds), CompressionFormat y una transformación de datos con Lambda. Firehose maneja los reintentos/retrocesos hacia los destinos y puede escribir registros fallidos en un bucket de S3 de respaldo. Una trampa común es el aprovisionamiento insuficiente de shards: los consumidores se quedan sin datos (starve) y la latencia se dispara; mide y escala de forma proactiva.

EventBridge y Step Functions para enrutamiento y orquestación

EventBridge se destaca en el enrutamiento de eventos basado en esquemas y en las integraciones entre cuentas/socios de eventos usando PutEvents para inyectar eventos y PutRule/PutTargets para enrutar a SQS, Lambda, Kinesis, Step Functions o puntos de conexión HTTP. EventBridge utiliza patrones de eventos para filtrar y admite el archivado y la reproducción (replay) para reconstruir el estado. Utiliza colas de mensajes fallidos (dead-letter queues) para las reglas (Target con SqsParameters o DeadLetterConfig) y ten en cuenta que EventBridge proporciona reintentos con retroceso exponencial (exponential backoff) y luego envía a la DLQ en caso de fallo. Para la orquestación, elige Step Functions: máquinas de estado Standard para flujos de trabajo duraderos y de larga duración con historial de ejecución y reintentos/Catch integrados, y Express para flujos de trabajo de corta duración y alto rendimiento con un costo menor y ejecución de mejor esfuerzo (best-effort). Usa integraciones de tareas (Task) con integraciones de servicio (

undefined

o

undefined

) y patrones de devolución de llamada (callback) usando “waitForTaskToken” para implementar aprobaciones externas asíncronas. Implementa reintentos y Catch con retroceso exponencial y usa HeartbeatSeconds para tareas largas. Usa el estado Map para paralelizar colecciones grandes, pero vigila la concurrencia y las limitaciones (throttles) en los servicios de destino. Un error común es elegir Express para flujos de trabajo que requieren un historial duradero y de procesamiento único (exactly-once); elige Standard para la auditabilidad. Además, asegura la idempotencia en las tareas invocadas por Step Functions pasando un token de idempotencia y haciendo que los servicios de destino apliquen la unicidad en el momento de la escritura.

Problema práctico: Escenario de caso de uso

Escenario: StreamlyGames opera un backend de juegos global en un entorno de AWS de múltiples cuentas. Los jugadores suben clips de juego de 10 MB a S3; una canalización de procesamiento debe transcodificar los videos, ejecutar análisis de ML y escribir los resultados en Aurora Serverless con un procesamiento ordenado, sin duplicados y con consumidores escalables.

Desafío: Asegurar que cada archivo subido active un procesamiento único (exactly-once) y en orden por cada jugador, gestionar los picos de carga sin perder eventos y escalar los consumidores para la inferencia de ML mientras se evitan escrituras duplicadas en la base de datos.

Enfoque recomendado:

  1. Crear una notificación de eventos de S3 para publicar eventos de creación de objetos (object-created) en un bus personalizado de EventBridge (PutEvents) y también en una cola SQS FIFO (CreateQueue con FifoQueue=true) usando el PlayerID como MessageGroupId y un MessageDeduplicationId basado en el ETag de S3.
  2. Configurar un consumidor Lambda con un mapeo de origen de eventos de SQS (CreateEventSourceMapping) usando BatchSize=1, FunctionResponseTypes=[“ReportBatchItemFailures”], y establecer un VisibilityTimeout mayor que el tiempo máximo de procesamiento; usar ChangeMessageVisibility al llamar a API de ML de terceros.
  3. La función Lambda realiza escrituras idempotentes en la base de datos Aurora usando una clave de idempotencia determinista (

undefined

o una restricción de unicidad) y guarda puntos de control del progreso; para llamadas de ML largas, usar Step Functions de forma asíncrona con tokens de tarea (

undefined

) o Step Functions Express para un alto rendimiento. 4. Para escalar la inferencia, escribir eventos intermedios en Kinesis Data Streams por shard y por región para consumidores de alto rendimiento y habilitar consumidores con distribución ramificada mejorada (enhanced fan-out) (SubscribeToShard) para flotas de trabajadores de ML dedicadas; monitorear IteratorAgeMilliseconds y usar UpdateShardCount para escalar.

Justificación: Usar SQS FIFO para garantizar el orden por jugador y la deduplicación en la ingesta, escrituras idempotentes en la base de datos para asegurar una semántica de procesamiento único (exactly-once), y Kinesis con distribución ramificada mejorada (enhanced fan-out) o Step Functions para manejar el procesamiento de ML de alto rendimiento y con ráfagas de carga, manteniendo los consumidores desacoplados y escalables.


Bases de datos y almacenamiento en caché (RDS · Todos los dominios

Practica estas preguntas → · Práctica cronometrada en 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.

Aprueba tu examen →

Explorar Amazon →

Related guides

Acceso todo en uno

Una suscripción. Todos los exámenes.

Cada plan desbloquea la búsqueda ilimitada de respuestas, pruebas de práctica, explicaciones de AI y la biblioteca completa de recursos, en más de 20 idiomas.

Mensual
24.87
Just €0.83/day
Todo incluido:
  • Búsqueda ilimitada de respuestas
  • Pruebas de práctica ilimitadas
  • Explicaciones con tecnología AI
  • Biblioteca completa de recursos
  • Más de 20 idiomas
  • Actualizaciones semanales de contenido
  • Recompensas y referencias
  • Soporte prioritario
Iniciar prueba gratuita

No se requiere tarjeta de crédito*

Mejor valor
12 meses
179.87
Just €0.49/daySave 40%
Todo incluido:
  • Búsqueda ilimitada de respuestas
  • Pruebas de práctica ilimitadas
  • Explicaciones con tecnología AI
  • Biblioteca completa de recursos
  • Más de 20 idiomas
  • Actualizaciones semanales de contenido
  • Recompensas y referencias
  • Soporte prioritario
Iniciar prueba gratuita

No se requiere tarjeta de crédito*

✓ Plan gratuito incluido · ✓ Cancela en cualquier momento · ✓ Todos los planes desbloquean el producto completo