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' }
但消息并未出现在队列中。请问我的代码存在什么问题,该如何修复?
问题原因及修复方案
核心问题分析
- 客户端区域不匹配:创建主题和队列的代码中,
SQSClient和SNSClient未指定区域,可能默认使用了与发布消息时不同的区域(发布代码指定了eu-west-2)。SNS与SQS是区域性服务,跨区域无法完成消息投递。 - SQS队列缺少SNS投递权限:即使订阅成功,SNS默认没有向SQS发送消息的权限,必须为队列添加明确的权限策略,允许指定SNS主题执行
sqs:SendMessage动作。 - 订阅函数参数名混淆:代码中
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; } }
验证修复
- 删除之前创建的主题、队列和订阅(避免残留资源干扰)。
- 重新运行创建资源的脚本,确保所有资源在
eu-west-2区域创建完成。 - 再次执行消息发布代码,检查SQS队列是否成功接收消息。
内容的提问来源于stack exchange,提问作者Wai Yan Hein
相关产品推荐
相关产品推荐

