AWS SQS FIFO队列消息丢失咨询:发送1000条仅收到69条
问题排查与解决方案
核心问题分析
AWS SQS本身是高可靠的消息队列服务,你的问题大概率出在代码逻辑或配置上,而非SQS服务本身。现有代码存在几个关键隐患:
未等待异步发送操作完成
sqs.send(sqsCommand)返回Promise,但循环中直接调用却未用await等待结果,也没有统一处理异步任务。这会导致:- Amplify执行环境可能在大量未完成的异步请求触发前就终止进程,直接中断消息发送
- 错误捕获不完整:
catch块抛出的错误无法被主线程感知,发送失败的消息不会被记录
MessageDeduplicationId 存在重复风险
用${Date.now()}-${x}生成去重ID时,若循环执行速度快于毫秒级(x递增速度超过时间戳变化),会出现多条消息的去重ID重复。FIFO队列会自动丢弃重复消息,直接导致部分消息丢失。缺少发送结果的确认日志
仅记录了发送启动日志,未记录发送成功的日志,无法确认实际成功发送的消息数量。
修复步骤
1. 修改代码,确保异步发送完成
将循环改为异步循环,等待每条消息发送Promise完成后再执行下一次,同时完善错误处理:
import { SQSClient, SendMessageCommand } from "@aws-sdk/client-sqs"; const sqs = new SQSClient({ region: env.SES_REGION, credentials: { accessKeyId: env.SES_ACCESS_KEY, secretAccessKey: env.SES_SECRET_KEY, }, }); const sqsMessageGroupId = `${Date.now()}`; // 改用异步循环确保每条消息发送完成 for (let x = 0; x < input.toAddress.length; x++) { const currentEmail = input.toAddress[x]!.emailAddress; console.log( `sending email ${x + 1} of ${input.toAddress.length}: ${currentEmail}, subject: ${input.subject}` ); const sqsCommand = new SendMessageCommand({ QueueUrl: env.SQS_QUEUE_URL, MessageBody: JSON.stringify({ to: [currentEmail], subject: input.subject, html: input.bodyHtml, text: input.bodyPlainText, }), MessageGroupId: sqsMessageGroupId, // 生成唯一去重ID,结合邮箱+时间戳+UUID避免重复 MessageDeduplicationId: `${currentEmail}-${Date.now()}-${crypto.randomUUID()}`, }); try { const result = await sqs.send(sqsCommand); console.log(`Email ${x + 1} sent successfully, MessageId: ${result.MessageId}`); } catch (err) { console.error(`Failed to send email ${x + 1} (${currentEmail}):`, err); // 根据业务需求选择是否终止流程,或继续发送剩余消息 // throw err; } }
2. 优化去重ID生成逻辑
结合唯一业务标识(如邮箱地址)+ 时间戳 + UUID生成MessageDeduplicationId,彻底避免重复,防止合法消息被误判为重复而丢弃。
3. 检查SQS队列配置
- 死信队列(DLQ):查看队列是否配置了死信队列,若消息发送时出现权限不足、消息大小超限等错误,可能被转发到死信队列,可在SQS控制台查看死信队列的消息数。
- 内容去重开关:确认FIFO队列的「基于内容的去重」是否开启,若开启,即使指定了
MessageDeduplicationId,SQS仍会根据消息内容判断重复,可能导致消息被丢弃。 - 消息保留期:默认保留期为4天,若配置了更短的保留期,消息可能被自动删除,但你的情况次日仍有69条消息,此可能性较低。
4. 补充监控与日志
- 查看CloudWatch中SQS的
NumberOfMessagesSent指标:该指标记录成功发送到队列的消息总数,可与循环次数对比,确认实际到达SQS的消息数量。 - 检查Amplify完整日志:排查是否存在未捕获的异步错误,如请求超时、权限不足等。
可靠性保障建议
SQS的可靠性基于正确使用:
- 必须等待发送请求完成,通过
await或Promise.all确保异步任务执行完毕 - 合理生成去重ID,避免合法消息被丢弃
- 配置死信队列,捕获发送或处理失败的消息
- 利用CloudWatch指标监控队列的消息流动情况
内容的提问来源于stack exchange,提问作者iceiceicy
相关产品推荐
相关产品推荐

