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

基于.NET Core与AWS SNS、SQS实现微服务发布订阅系统咨询

.NET Core 集成 AWS SNS + SQS 实现发布订阅方案

前置依赖安装

首先安装对应AWS .NET SDK的NuGet包:

  • AWSSDK.SimpleNotificationService:SNS操作核心包
  • AWSSDK.SQS:SQS操作核心包
  • AWSSDK.Extensions.NETCore.Setup:.NET Core DI容器集成AWS服务的扩展包

AWS凭证配置可以写在appsettings.json,本地开发也可以用AWS CLI生成的本地凭证文件,生产环境直接使用服务对应的IAM角色即可,禁止硬编码AK/SK。

// appsettings.json配置示例
"AWS": {
  "Region": "替换为你的AWS区域,如cn-north-1",
  "Profile": "本地开发的凭证profile名称,生产环境可删除"
}

服务注册

在Program.cs中注入SNS、SQS客户端到DI容器:

// 读取AWS配置
builder.Services.AddDefaultAWSOptions(builder.Configuration.GetAWSOptions());
// 注入SNS、SQS客户端
builder.Services.AddAWSService<IAmazonSimpleNotificationService>();
builder.Services.AddAWSService<IAmazonSQS>();

发布端(SNS主题消息推送)实现

你可以提前在AWS控制台创建SNS主题,也可以用代码动态创建,发布消息的示例封装如下:

public class SnsMessagePublisher
{
    private readonly IAmazonSimpleNotificationService _snsClient;
    // 替换为你的SNS主题ARN
    private const string SnsTopicArn = "arn:aws-cn:sns:cn-north-1:123456789012:your-topic-name";

    public SnsMessagePublisher(IAmazonSimpleNotificationService snsClient)
    {
        _snsClient = snsClient;
    }

    /// <summary>
    /// 发布自定义类型消息到SNS主题
    /// </summary>
    public async Task PublishAsync<T>(T messageObj)
    {
        var messageContent = System.Text.Json.JsonSerializer.Serialize(messageObj);
        var publishRequest = new PublishRequest
        {
            TopicArn = SnsTopicArn,
            Message = messageContent,
            // 可添加消息属性,用于后续SNS订阅端过滤消息
            MessageAttributes = new Dictionary<string, MessageAttributeValue>
            {
                {
                    "MessageType", 
                    new MessageAttributeValue { StringValue = typeof(T).Name, DataType = "String" }
                }
            }
        };
        await _snsClient.PublishAsync(publishRequest);
    }
}

订阅端(SNS+SQS搭配消费)实现

微服务场景下推荐用SQS作为SNS的订阅端点,相比HTTP推送更稳定,可避免消息丢失、支持削峰填谷。

前提配置

  1. AWS控制台创建SQS队列,给对应SNS主题授予发送消息到该队列的权限
  2. 将SQS队列绑定订阅到SNS主题,可按需配置消息过滤规则,只接收指定类型的消息

消息消费封装示例

如果是常驻消费场景,可封装为后台服务:

// 消费逻辑封装
public class SqsMessageConsumer
{
    private readonly IAmazonSQS _sqsClient;
    // 替换为你的SQS队列URL
    private const string SqsQueueUrl = "https://sqs.cn-north-1.amazonaws.com.cn/123456789012/your-queue-name";

    public SqsMessageConsumer(IAmazonSQS sqsClient)
    {
        _sqsClient = sqsClient;
    }

    public async Task StartConsumeAsync(CancellationToken stoppingToken)
    {
        while (!stoppingToken.IsCancellationRequested)
        {
            var receiveRequest = new ReceiveMessageRequest
            {
                QueueUrl = SqsQueueUrl,
                MaxNumberOfMessages = 10, // 单次最多拉取10条
                WaitTimeSeconds = 20, // 开启长轮询,减少空请求开销
                MessageAttributeNames = new List<string> { "All" }
            };

            var receiveResponse = await _sqsClient.ReceiveMessageAsync(receiveRequest, stoppingToken);
            foreach (var message in receiveResponse.Messages)
            {
                try
                {
                    // 此处编写你的业务处理逻辑
                    Console.WriteLine($"收到消息:{message.Body}");

                    // 处理成功后删除消息,避免重复消费
                    await _sqsClient.DeleteMessageAsync(SqsQueueUrl, message.ReceiptHandle, stoppingToken);
                }
                catch (Exception ex)
                {
                    // 处理失败的消息会在可见性超时后重新回到队列,建议配置死信队列存储多次消费失败的消息
                    Console.WriteLine($"消息处理失败:{ex.Message}");
                }
            }
        }
    }
}

// 后台消费服务注册
public class SqsConsumeBackgroundService : BackgroundService
{
    private readonly SqsMessageConsumer _consumer;

    public SqsConsumeBackgroundService(SqsMessageConsumer consumer)
    {
        _consumer = consumer;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        await _consumer.StartConsumeAsync(stoppingToken);
    }
}

// Program.cs中注册后台服务
builder.Services.AddHostedService<SqsConsumeBackgroundService>();

注意事项

  • 幂等处理:SQS默认是至少一次投递,业务逻辑必须做幂等校验,避免重复消费导致数据异常
  • 权限控制:生产环境不要使用AK/SK硬编码,直接为ECS、EKS等服务对应的IAM角色授予SNS发布、SQS消费的最小权限即可
  • 消息过滤:如果同一个SNS主题需要推送多种业务消息,可通过消息属性配置订阅过滤规则,无需消费全量消息

内容的提问来源于stack exchange,提问作者heisenberg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 21:15:03