SNS消息已推送至SQS队列但订阅未消费问题咨询
问题解答:SNS发布后为何SQS消息未被消费?是否必须单独处理SQS消费?
核心结论
是的,你必须单独实现SQS的消息消费逻辑。SNS和SQS是分工协作的:
- SNS负责发布消息并转发到订阅的端点(比如你的SQS队列)
- SQS负责存储消息并等待消费端主动拉取/触发处理
你的代码目前只完成了「绑定SQS到SNS主题」的订阅关系创建,没有实现从SQS拉取消息的消费逻辑。
你的代码问题分析
你当前的subscribe函数只是调用了AWS SNS的subscribe API,这一步的作用是建立SNS主题到SQS队列的转发关系,只会在服务器启动时执行一次,返回的是订阅ARN(就是你看到的日志内容)。这个操作本身不负责消费任何消息——消息被SNS推送到SQS后,会一直存在队列里,直到你主动去拉取处理。
修正方案:添加SQS消费逻辑
你需要在代码中新增SQS客户端,实现消息拉取和处理的逻辑。以下是修改后的插件示例:
import { SNS, SQS } from 'aws-sdk'; import fp from 'fastify-plugin'; export const awsSNSPlugin = fp(async server => { // 初始化SNS和SQS客户端 const snsService = new SNS({ apiVersion: '2010-03-31' }); const sqsService = new SQS({ apiVersion: '2012-11-05' }); // SNS发布逻辑(保留原代码) const publish: AWSPubSub['publish'] = async params => { try { const response = await snsService.publish(params).promise(); server.log.info(`Message ${params.Message} sent to the topic ${params.TopicArn}`); server.log.info(response); } catch (err: any) { server.log.error(err, err.stack); } }; // 订阅SNS到SQS,并启动SQS消费 const subscribe: AWSPubSub['subscribe'] = async (params, messageHandler) => { try { // 1. 创建SNS到SQS的订阅关系 const subscribeResponse = await snsService.subscribe(params).promise(); server.log.info(`订阅关系创建成功: ${JSON.stringify(subscribeResponse)}`); // 2. 启动SQS消息消费轮询 const queueUrl = params.Endpoint; // 假设你的Endpoint是SQS的队列URL const pollMessages = async () => { try { // 拉取SQS消息(最多10条,等待时间20秒长轮询) const receiveResponse = await sqsService.receiveMessage({ QueueUrl: queueUrl, MaxNumberOfMessages: 10, WaitTimeSeconds: 20, MessageAttributeNames: ['All'] }).promise(); if (receiveResponse.Messages) { for (const msg of receiveResponse.Messages) { // 调用传入的消息处理回调 if (messageHandler) { await messageHandler(msg.Body); } // 处理完成后删除消息,避免重复消费 await sqsService.deleteMessage({ QueueUrl: queueUrl, ReceiptHandle: msg.ReceiptHandle! }).promise(); } } } catch (err: any) { server.log.error('SQS消费出错:', err); } finally { // 继续轮询 setTimeout(pollMessages, 0); } }; // 启动轮询 pollMessages(); } catch (err: any) { server.log.error(err, err.stack); } }; server.decorate('pubSub', { publish, subscribe, }); });
使用说明
调用subscribe时,传入你的SQS队列URL作为Endpoint,并传入消息处理回调:
server.pubSub.subscribe( { TopicArn: '你的SNS主题ARN', Protocol: 'sqs', Endpoint: '你的SQS队列URL' }, (message) => { server.log.info(`收到消息: ${message}`); // 这里写你的消息处理逻辑 } );
补充说明
SNS+SQS的组合是AWS推荐的异步消息架构,优势在于:
- 解耦发布者和消费者:发布者不用关心谁在消费,只需发送到SNS
- 消息持久化:SQS会保存消息直到被消费,避免消息丢失
- 支持多消费者:多个服务可以同时订阅同一个SNS主题,各自用SQS接收消息
这种模式并不奇怪,是分布式系统中常见的消息解耦方案。
内容的提问来源于stack exchange,提问作者Mike Yim
相关产品推荐
相关产品推荐

