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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 14:12:04