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

AWS SNS FIFO主题消息无法投递至订阅的SQS FIFO队列问题排查

问题:SNS FIFO主题消息无法投递到订阅的SQS FIFO队列

我正尝试结合使用AWS SNS与SQS队列,期望将消息发布至AWS SNS FIFO主题后,消息能投递至订阅该主题的SQS FIFO队列,后续再处理队列中的消息。我通过以下代码创建了主题、队列及订阅:

import { SNSClient, SubscribeCommand, CreateTopicCommand } from "@aws-sdk/client-sns";
import {
    CreateQueueCommand,
    SQSClient,
    GetQueueAttributesCommand,
} from "@aws-sdk/client-sqs";

const qName = `OrderQ.fifo`;
const sqsClient = new SQSClient({});
const topicName = `OrderTopic.fifo`;
const snsClient = new SNSClient({});

const createTopic = async () => {
    try {
        const command = new CreateTopicCommand({
            Name: topicName,
            Attributes: { // for FIFO topic
                FifoTopic: "true",
                ContentBasedDeduplication: "true"
            }
        });

        const result = await snsClient.send(command);

        return result;
    } catch (e) {
        throw e;
    }
}

const createQ = async () => {
    try {
        const command = new CreateQueueCommand({
            QueueName: qName,
            Attributes: {
                "FifoQueue": "true"
            }
        });

        const result = await sqsClient.send(command);

        return result;
    } catch (e) {
        throw e;
    }
}

/**
 * @param {string} qUrl
 * @param {string} topicArn
 */
const subscribeQToTopic = async (qUrl, topicArn) => {
    try {
        const command = new SubscribeCommand({
            Protocol: "sqs",
            TopicArn: topicArn,
            Endpoint: qUrl,
        });

        const result = await snsClient.send(command);

        return result;
    } catch (e) {
        throw e;
    }
}

const createTopicAndSubscription = async () => {
    try {
        const qResult = await createQ();

        const qAttrsResult = await sqsClient.send(new GetQueueAttributesCommand({
            QueueUrl: qResult.QueueUrl,
            AttributeNames: ["QueueArn"]
        }));

        const topicResult = await createTopic();

        const result = await subscribeQToTopic(qAttrsResult.Attributes.QueueArn, topicResult.TopicArn);

        return result;
    } catch (e) {
        throw e;
    }
};

createTopicAndSubscription().then((result) => {
    console.log(result);
}).catch(e => {
    console.error(e);
})

执行脚本后,主题、队列及订阅均成功创建。随后我通过主题ARN发布消息,代码如下:

import {PublishCommand, SNSClient} from "@aws-sdk/client-sns";
    
const topicArn = `topc_arn`;

const snsClient = new SNSClient({
    region: `eu-west-2`,
});

const publishMessage = async () => {
    try {
        const result = await snsClient.send(new PublishCommand({
            Message: "this is the test message", // MESSAGE_TEXT
            TopicArn: topicArn, //TOPIC_ARN
            MessageGroupId: `order_group`,
        }));

        console.log(result);
    } catch (e) {
        console.error(e);
    }
}

publishMessage();

执行该代码未抛出错误,返回了如下成功响应:

{
  '$metadata': {
    httpStatusCode: 200,
    requestId: 'ed23d6fb-185e-5487-933a-8bf46b15f727',
    extendedRequestId: undefined,
    cfId: undefined,
    attempts: 1,
    totalRetryDelay: 0
  },
  MessageId: 'a5b68fe0-e944-5b6f-bec2-d2ef9a0fe8cf',
  SequenceNumber: '10000000000000008000'
}

但消息并未出现在队列中。请问我的代码存在什么问题,该如何修复?


问题原因及修复方案

核心问题分析

  1. 客户端区域不匹配:创建主题和队列的代码中,SQSClient和SNSClient未指定区域,可能默认使用了与发布消息时不同的区域(发布代码指定了eu-west-2)。SNS与SQS是区域性服务,跨区域无法完成消息投递。
  2. SQS队列缺少SNS投递权限:即使订阅成功,SNS默认没有向SQS发送消息的权限,必须为队列添加明确的权限策略,允许指定SNS主题执行sqs:SendMessage动作。
  3. 订阅函数参数名混淆:代码中subscribeQToTopic函数的参数名是qUrl,但实际传递的是队列ARN,虽然当前逻辑未出错,但容易造成后续维护混淆。

具体修复步骤

1. 统一客户端区域

在创建资源的代码中,为SQSClient和SNSClient指定与发布代码一致的区域:

const sqsClient = new SQSClient({ region: "eu-west-2" });
const snsClient = new SNSClient({ region: "eu-west-2" });

2. 为SQS队列添加权限策略

修改createQ函数,在创建队列后添加允许SNS投递消息的权限策略,同时调整创建顺序(先创建主题再创建队列,以便获取主题ARN生成策略):

// 新增导入SetQueueAttributesCommand
import {
    CreateQueueCommand,
    SQSClient,
    GetQueueAttributesCommand,
    SetQueueAttributesCommand,
} from "@aws-sdk/client-sqs";

// 修改createQ函数,接收topicArn参数
const createQ = async (topicArn) => {
    try {
        const command = new CreateQueueCommand({
            QueueName: qName,
            Attributes: {
                "FifoQueue": "true"
            }
        });

        const result = await sqsClient.send(command);

        // 获取队列ARN
        const qAttrsResult = await sqsClient.send(new GetQueueAttributesCommand({
            QueueUrl: result.QueueUrl,
            AttributeNames: ["QueueArn"]
        }));
        const queueArn = qAttrsResult.Attributes.QueueArn;

        // 构建并设置权限策略
        const policy = JSON.stringify({
            Version: "2012-10-17",
            Statement: [
                {
                    Effect: "Allow",
                    Principal: { Service: "sns.amazonaws.com" },
                    Action: "sqs:SendMessage",
                    Resource: queueArn,
                    Condition: {
                        ArnEquals: { "aws:SourceArn": topicArn }
                    }
                }
            ]
        });

        await sqsClient.send(new SetQueueAttributesCommand({
            QueueUrl: result.QueueUrl,
            Attributes: { Policy: policy }
        }));

        return { ...result, QueueArn: queueArn };
    } catch (e) {
        throw e;
    }
}

// 调整createTopicAndSubscription的执行顺序
const createTopicAndSubscription = async () => {
    try {
        const topicResult = await createTopic();
        const qResult = await createQ(topicResult.TopicArn);
        const result = await subscribeQToTopic(qResult.QueueArn, topicResult.TopicArn);
        return result;
    } catch (e) {
        throw e;
    }
};

3. 修正订阅函数参数名(可选但推荐)

修改subscribeQToTopic的参数名,明确传递的是队列ARN而非URL:

/**
 * @param {string} queueArn
 * @param {string} topicArn
 */
const subscribeQToTopic = async (queueArn, topicArn) => {
    try {
        const command = new SubscribeCommand({
            Protocol: "sqs",
            TopicArn: topicArn,
            Endpoint: queueArn,
        });

        const result = await snsClient.send(command);
        return result;
    } catch (e) {
        throw e;
    }
}

验证修复

  1. 删除之前创建的主题、队列和订阅(避免残留资源干扰)。
  2. 重新运行创建资源的脚本,确保所有资源在eu-west-2区域创建完成。
  3. 再次执行消息发布代码,检查SQS队列是否成功接收消息。

内容的提问来源于stack exchange,提问作者Wai Yan Hein

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 09:15:00