Consumidores de AWS SQS
Amazon SQS proporciona colas gestionadas sin necesidad de operar Redis. Los workers de Node usan AWS SDK v3 con sondeo largo, tiempos de espera de visibilidad y colas de mensajes fallidos (DLQ).
Busca en todas las páginas de la documentación
Amazon SQS proporciona colas gestionadas sin necesidad de operar Redis. Los workers de Node usan AWS SDK v3 con sondeo largo, tiempos de espera de visibilidad y colas de mensajes fallidos (DLQ).
Tarjeta de receta de referencia rápida: lista para copiar y pegar.
import {
SQSClient,
ReceiveMessageCommand,
DeleteMessageCommand,
} from "@aws-sdk/client-sqs";
const sqs = new SQSClient({});
const queueUrl = process.env.SQS_QUEUE_URL!;
const res = await sqs.send(
new ReceiveMessageCommand({
QueueUrl: queueUrl,
MaxNumberOfMessages: 10,
WaitTimeSeconds: 20,
VisibilityTimeout: 60,
})
);
for (const msg of res.Messages ?? []) {
await processBody(msg.Body!);
await sqs.send(
new DeleteMessageCommand({ QueueUrl: queueUrl, ReceiptHandle: msg.ReceiptHandle! })
);
}Cuándo usarlo:
// src/sqs/poll.ts
import {
SQSClient,
ReceiveMessageCommand,
DeleteMessageCommand,
ChangeMessageVisibilityCommand,
} from "@aws-sdk/client-sqs";
const sqs = new SQSClient({ region: process.env.AWS_REGION });
const queueUrl = process.env.SQS_QUEUE_URL!;
async function handleMessage(body: string) {
const payload = JSON.parse(body) as { type: string; orderId: string };
if (payload.type === "fulfill") {
await fulfillOrder(payload.orderId);
}
}
export async function pollForever(signal: AbortSignal) {
while (!signal.aborted) {
const res = await sqs.send(
new ReceiveMessageCommand({
QueueUrl: queueUrl,
MaxNumberOfMessages: 10,
WaitTimeSeconds: 20,
VisibilityTimeout: 120,
MessageAttributeNames: ["All"],
})
);
for (const msg of res.Messages ?? []) {
try {
await handleMessage(msg.Body ?? "{}");
await sqs.send(
new DeleteMessageCommand({
QueueUrl: queueUrl,
ReceiptHandle: msg.ReceiptHandle!,
})
);
} catch (err) {
console.error({ msg: "sqs_handler_error", messageId: msg.MessageId, err });
// El mensaje vuelve a la cola después del tiempo de espera de visibilidad
}
}
}
}
// src/worker-main.ts
const ac = new AbortController();
process.on("SIGTERM", () => ac.abort());
await pollForever(ac.signal);Lo que esto demuestra:
WaitTimeSeconds: 20 reduce el costo y el uso de CPUvisibility_timeout >= p99_processing_time * 1.5
ChangeMessageVisibility para trabajos de duración variable{
"RedrivePolicy": {
"deadLetterTargetArn": "arn:aws:sqs:...:dlq",
"maxReceiveCount": 5
}
}| Tipo | Orden | Rendimiento | Deduplicación |
|---|---|---|---|
| Estándar | Mejor esfuerzo | Muy alto | Idempotencia a nivel de aplicación |
| FIFO | Por grupo de mensajes | 300 TPS/grupo | Deduplicación de contenido opcional |
MessageGroupId// El cuerpo del mensaje puede ser JSON envuelto de SNS - desenvolver el sobre
const outer = JSON.parse(body);
const inner = outer.Message ? JSON.parse(outer.Message) : outer;maxReceiveCount + alarma de DLQ.Message.| Alternativa | Usar cuándo | No usar cuándo |
|---|---|---|
| BullMQ | Redis ya en ejecución, API de trabajos rica | Quieres cero infraestructura de cola |
| Kinesis | Análisis de streaming | Cola de tareas simple |
| EventBridge | Reglas de enrutamiento de eventos | Solo cola de worker punto a punto |
| Google Pub/Sub | Stack de GCP | Organización solo de AWS |
sqs-consumer envuelve el sondeo con eventos. Está bien para workers de ECS. Comprende la semántica de visibilidad de cualquier manera.
Escala los consumidores en la métrica ApproximateNumberOfMessagesVisible. Evita bucles cerrados duplicados sin sondeo largo.
No. La cola estándar es al menos una vez. Diseña manejadores idempotentes.
Usa el patrón de cliente extendido de S3: SQS lleva un puntero de S3 cuando >256KB.
sqs:ReceiveMessage, DeleteMessage, ChangeMessageVisibility solo en el ARN de la cola.
LocalStack o ElasticMQ para pruebas de integración. Docker de ElasticMQ para laptop.
SendMessageCommand desde la API con un pequeño cuerpo JSON. Mismas reglas de idempotencia que el productor de BullMQ.
4 días por defecto. Aumenta a 14 para ventanas de reejecución durante interrupciones.
El procesamiento por lotes y múltiples IDs de grupo de mensajes escalan flujos ordenados paralelos.
Conteo de solicitudes. El sondeo largo reduce las recepciones vacías. Recepción por lotes de hasta 10 mensajes.
Versiones de stack: Esta página fue escrita para Node.js 24.18.0 (LTS activa), npm 10+, TypeScript 5.6+, Express 5, Fastify 5 y NestJS 11.
Revisado por Chris St. John·Última actualización: 19 jul 2026