You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

AWS SQS无法接收队列消息(Node.js)及配置咨询

问题分析与解决方案

为什么没收到预期的10条消息?

  1. 短轮询机制:默认短轮询会立即返回,哪怕队列有消息,也可能只返回部分甚至空结果;
  2. FIFO队列消息组阻塞:若FIFO队列所有消息用同一个MessageGroupId,SQS会按顺序处理,同一时间仅一条消息可见,导致轮询只能拿到1条;
  3. 未处理消息未删除:Worker拿到消息后未删除,消息会在可见性超时后重回队列,干扰后续轮询的消息数量;
  4. 定时器异步叠加: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 07:37:44