基于.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推送更稳定,可避免消息丢失、支持削峰填谷。
前提配置
- AWS控制台创建SQS队列,给对应SNS主题授予发送消息到该队列的权限
- 将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
相关产品推荐
相关产品推荐

