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.
- SNS: fan-out de alto rendimiento, sin ordenamiento, Publish/Subscribe, usa MessageAttributes para enrutamiento.
- SQS Standard: entrega al menos una vez (at-least-once), ordenamiento de mejor esfuerzo (best-effort), sondeo largo (long polling), más económico para desacoplamiento.
- SQS FIFO: ordenamiento exactamente una vez (exactly-once) dentro de un MessageGroupId, úsalo para ordenamiento estricto y deduplicación.
- Kinesis Data Streams: ordenado por shard, streaming de alto rendimiento, PutRecord/PutRecords, requiere escalado de shards.
- Kinesis Firehose: entrega y almacenamiento en búfer gestionados, admite cifrado del lado del servidor y transformaciones con Lambda.
- EventBridge: bus de eventos con reglas de enrutamiento, registro de esquemas, archivo y reproducción (archive & replay), API PutEvents.
- Step Functions: orquestación con estado (stateful), reintentos/Catch, compensaciones entre Standard y Express en cuanto a durabilidad y rendimiento.
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:
- 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.
- 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.
- 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 →