AWS SQS无法接收队列消息(Node.js)及配置咨询
问题分析与解决方案
为什么没收到预期的10条消息?
- 短轮询机制:默认短轮询会立即返回,哪怕队列有消息,也可能只返回部分甚至空结果;
- FIFO队列消息组阻塞:若FIFO队列所有消息用同一个
MessageGroupId,SQS会按顺序处理,同一时间仅一条消息可见,导致轮询只能拿到1条; - 未处理消息未删除:Worker拿到消息后未删除,消息会在可见性超时后重回队列,干扰后续轮询的消息数量;
- 定时器异步叠加:
setInterval不管前一次sendMessage是否完成都会触发,可能导致实际发送频率不稳定,甚至丢消息。
满足需求的配置与代码修改
队列配置(必须用FIFO队列)
- 启用内容去重:创建队列时开启「内容去重」,SQS会根据消息内容自动生成去重ID,避免重复消息;
- 设置可见性超时:建议设为10秒(比Worker轮询间隔5秒长),确保Worker处理消息期间,消息不会重新回到队列;
- 消息保留期:默认4天,满足「无消费者时存储消息」需求,如需更长可调整至最大14天。
Timer.js 修改(确保每秒发一条+适配FIFO)
import { SQSClient, SendMessageCommand } from "@aws-sdk/client-sqs"; const SQS = new SQSClient({ credentials: { accessKeyId: "__ACCESS_KEY_ID__", secretAccessKey: "__SECRET_ACCESS_KEY__", }, region: "__REGION__", }); const sendMessage = async () => { const messageBody = new Date().toISOString(); const input = { QueueUrl: "__QUEUE_URL__", MessageBody: messageBody, // 手动指定去重ID,与内容去重二选一(手动更可控) MessageDeduplicationId: messageBody, // 每个消息用不同组ID,让FIFO队列并行处理消息 MessageGroupId: `group-${Math.random().toString(36).slice(2, 10)}`, }; const command = new SendMessageCommand(input); await SQS.send(command); // 递归调用确保每秒发送一条,避免异步叠加 setTimeout(sendMessage, 1000); }; sendMessage();
Worker.js 修改(长轮询+消息删除)
import { SQSClient, ReceiveMessageCommand, DeleteMessageCommand } from "@aws-sdk/client-sqs"; const SQS = new SQSClient({ credentials: { accessKeyId: "__ACCESS_KEY_ID__", secretAccessKey: "__SECRET_ACCESS_KEY__", }, region: "__REGION__", }); const worker = async () => { console.log("polling for messages..."); const input = { QueueUrl: "__QUEUE_URL__", MaxNumberOfMessages: 10, // 启用长轮询,最长20秒,有消息立即返回,超时再返回空 WaitTimeSeconds: 20, // 可见性超时设为10秒,比轮询间隔长 VisibilityTimeout: 10, }; const command = new ReceiveMessageCommand(input); const response = await SQS.send(command); const { Messages } = response; if (!Messages) { console.log("no messages received"); return; } const messagesReceived = Messages.map((e) => ({ id: e.MessageId, receiptHandle: e.ReceiptHandle, body: e.Body, })); console.log("received messages:", messagesReceived.map(e => e.body)); // 逐个删除已处理消息 for (const msg of messagesReceived) { const deleteCommand = new DeleteMessageCommand({ QueueUrl: "__QUEUE_URL__", ReceiptHandle: msg.receiptHandle, }); await SQS.send(deleteCommand); } }; setInterval(worker, 5000);
关键说明
- 去重实现:FIFO队列的内容去重或手动指定
MessageDeduplicationId,确保同一内容的消息不会被重复存储; - 并行处理:不同的
MessageGroupId让FIFO队列突破顺序限制,同时返回多个消息,满足每次轮询拿10条的需求; - 长轮询:减少空轮询次数,提高消息接收的及时性和完整性;
- 消息删除:必须删除已处理的消息,否则消息会重回队列,导致重复消费和队列积压。
内容的提问来源于stack exchange,提问作者Sam Leurs
相关产品推荐
相关产品推荐

