- 示例工程
- 教程
- 后端
【免费下载链接】aws-doc-sdk-examples
Welcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below.
导读
本文以 dotnetv3/SQS/README.md 为核心,系统讲解如何使用 AWS SDK for .NET 3.x 开发 Amazon Simple Queue Service(SQS)应用:涵盖ListQueues入门示例、创建/发送/接收/删除队列等全部单操作代码,以及"SNS 主题 + SQS 队列"消息发布场景和 AWS Message Processing Framework for .NET 框架化开发两大完整场景。读完本文,你将掌握 SQS 在 .NET 中的全部核心 API 调用方式、队列策略与 FIFO 配置要点,并能直接运行仓库中的可执行示例与测试。
Amazon SQS 是一项全托管的消息队列服务,用于解耦和扩展微服务、分布式系统与无服务器应用。仓库中的示例覆盖了从"Hello World"级别的入门调用到跨服务消息工作流的完整路径,是学习 .NET 集成 SQS 最直接的参考实现。
⚠️ 重要提示
- 运行这些代码可能产生 AWS 账户费用,详见 AWS 定价与免费套餐说明。
- 运行测试同样可能产生费用。
- 建议按最小权限原则(Grant least privilege)授权,仅授予任务所需的最低权限。
- 本仓库代码未经所有 AWS 区域验证,部分服务仅在特定区域可用。
一、示例概览与代码布局
dotnetv3目录下的 SQS 示例分为两类:单操作代码片段(Single actions)与完整场景(Scenarios),此外还有一个用于快速上手的入门示例(Get started)。
| 类别 | 内容 | 位置 |
|---|---|---|
| 入门(Hello SQS) | 列出队列ListQueues | HelloSQS.cs |
| 单操作 | CreateQueue、GetQueueAttributes、SetQueueAttributes、ReceiveMessage、DeleteMessageBatch、DeleteQueue | SQSWrapper.cs |
| 单操作 | SendMessage、CreateQueue(独立示例) | CreateSendExample.cs |
| 单操作 | GetQueueUrl | GetQueueUrl.cs |
| 单操作 | ReceiveMessage+DeleteMessage | ReceiveDeleteExample.cs |
| 场景 | SNS 主题发布消息到 SQS 队列 | TopicsAndQueues.cs |
| 场景 | AWS Message Processing Framework for .NET | MessageProcessingFramework |
上述代码同时被 TopicsAndQueues 跨服务示例 复用,其中SQSWrapper.cs是 SQS 操作的核心封装类,单操作代码片段即取自于此。
二、环境准备(Prerequisites)
运行本目录示例前,需要先完成 dotnetv3/README.md 中列出的前置条件,核心包括:
- 安装合适的 .NET SDK(多数示例使用 .NET 6,部分要求 .NET 5)。
- 引入 AWS SDK for .NET 及对应服务包(SQS 需要
Amazon.SQS命名空间)。 - 配置 AWS 凭证:可通过本地 credentials 文件配置,或设置
AWS_ACCESS_KEY_ID与AWS_SECRET_ACCESS_KEY环境变量。
如果使用 AWS IAM Identity Center 进行身份认证,还需为项目额外添加AWSSDK.SSO和AWSSDK.SSOOIDCNuGet 包(见 MessageProcessingFramework 说明)。
三、入门示例:Hello Amazon SQS(ListQueues)
仓库提供了一个最小入门示例,演示如何使用异步 API 列出当前账户下的队列:
var sqsClient = new AmazonSQSClient(); Console.WriteLine("Hello Amazon SQS! Following are some of your queues:"); // Let's get the first five queues. var response = await sqsClient.ListQueuesAsync( new ListQueuesRequest() { MaxResults = 5 }); foreach (var queue in response.QueueUrls) { Console.WriteLine($"\tQueue Url: {queue}"); }关键点(对应 HelloSQS.cs):
- 直接构造
AmazonSQSClient(不传区域时使用默认凭证链的默认区域)。 - 使用
await调用异步方法ListQueuesAsync。 - 通过
MaxResults = 5限制返回队列数量,响应中的QueueUrls保存各队列的 URL。
这是验证 SDK 安装、凭证配置与网络连通性最快的方式。
四、单操作详解(Single Actions)
以下示例展示如何调用 SQS 的单个服务函数。所有代码均可独立运行,路径与行号见文末代码地图。
4.1 创建队列(CreateQueue)
SQSWrapper.cs 中的CreateQueueWithName展示了标准队列与 FIFO 队列的创建差异:
public async Task<string> CreateQueueWithName(string queueName, bool useFifoQueue) { int maxMessage = 256 * 1024; var queueAttributes = new Dictionary<string, string> { { QueueAttributeName.MaximumMessageSize, maxMessage.ToString() } }; var createQueueRequest = new CreateQueueRequest() { QueueName = queueName, Attributes = queueAttributes }; if (useFifoQueue) { // Update the name if it is not correct for a FIFO queue. if (!queueName.EndsWith(".fifo")) { createQueueRequest.QueueName = queueName + ".fifo"; } // Add an attribute for a FIFO queue. createQueueRequest.Attributes.Add( QueueAttributeName.FifoQueue, "true"); } var createResponse = await _amazonSQSClient.CreateQueueAsync( new CreateQueueRequest() { QueueName = queueName }); return createResponse.QueueUrl; }要点:
- 最大消息大小:示例将
MaximumMessageSize设为256 * 1024(256 KB),这也是 SQS 支持的上限。 - FIFO 队列:启用
QueueAttributeName.FifoQueue = "true"属性,且队列名必须以.fifo结尾,否则调用会被拒绝——代码中通过EndsWith(".fifo")检查并自动补全后缀。 - 创建成功返回队列 URL(
QueueUrl),后续的收发消息、设置策略等操作都以该 URL 为定位依据。
独立的 CreateSendExample.cs 则展示了另一种配置方式,通过DelaySeconds(延迟投递 60 秒)和MessageRetentionPeriod(消息保留 86400 秒,即 1 天)设置队列属性:
var request = new CreateQueueRequest { QueueName = queueName, Attributes = new Dictionary<string, string> { { "DelaySeconds", "60" }, { "MessageRetentionPeriod", "86400" }, }, }; var response = await client.CreateQueueAsync(request);4.2 获取队列属性(GetQueueAttributes)
创建队列后,通常需要获取其 ARN(Amazon Resource Name)用于配置策略或订阅。参见 SQSWrapper.cs:
public async Task<string> GetQueueArnByUrl(string queueUrl) { var getAttributesRequest = new GetQueueAttributesRequest() { QueueUrl = queueUrl, AttributeNames = new List<string>() { QueueAttributeName.QueueArn } }; var getAttributesResponse = await _amazonSQSClient.GetQueueAttributesAsync(getAttributesRequest); return getAttributesResponse.QueueARN; }通过AttributeNames指定要查询的属性(此处为QueueArn),响应中直接返回QueueARN。
4.3 设置队列策略(SetQueueAttributes)
要让 SNS 主题能够向队列投递消息,必须为队列设置允许sqs:SendMessage的访问策略。参见 SQSWrapper.cs:
var queuePolicy = "{" + "\"Version\": \"2012-10-17\"," + "\"Statement\": [{" + "\"Effect\": \"Allow\"," + "\"Principal\": {" + $"\"Service\": " + "\"sns.amazonaws.com\"" + "}," + "\"Action\": \"sqs:SendMessage\"," + $"\"Resource\": \"{queueArn}\"," + "\"Condition\": {" + "\"ArnEquals\": {" + $"\"aws:SourceArn\": \"{topicArn}\"" + "}" + "}" + "}]" + "}"; var attributesResponse = await _amazonSQSClient.SetQueueAttributesAsync( new SetQueueAttributesRequest() { QueueUrl = queueUrl, Attributes = new Dictionary<string, string>() { { "Policy", queuePolicy } } }); return attributesResponse.HttpStatusCode == HttpStatusCode.OK;策略要点:
Principal指定服务主体sns.amazonaws.com,仅允许 SNS 服务访问。Action限制为sqs:SendMessage,遵循最小权限原则。Condition使用ArnEquals的aws:SourceArn约束来源主题 ARN,防止其他主题滥用该队列——这是防御跨账户/跨资源混淆代理攻击(confused deputy)的标准写法。
4.4 发送消息(SendMessage)
CreateSendExample.cs 展示了带消息属性的发送方式:
var sendMessageRequest = new SendMessageRequest { DelaySeconds = 10, MessageAttributes = messageAttributes, MessageBody = messageBody, QueueUrl = queueUrl, }; var response = await client.SendMessageAsync(sendMessageRequest); Console.WriteLine($"Sent a message with id : {response.MessageId}");其中messageAttributes通过MessageAttributeValue构造,支持String与Number等数据类型:
Dictionary<string, MessageAttributeValue> messageAttributes = new Dictionary<string, MessageAttributeValue> { { "Title", new MessageAttributeValue { DataType = "String", StringValue = "The Whistler" } }, { "Author", new MessageAttributeValue { DataType = "String", StringValue = "John Grisham" } }, { "WeeksOn", new MessageAttributeValue { DataType = "Number", StringValue = "6" } }, };消息属性(Message Attributes)让消息附带结构化元数据,接收方可用于过滤或业务判断。DelaySeconds = 10表示消息在队列中延迟 10 秒后才可见。
4.5 获取队列 URL(GetQueueUrl)
GetQueueUrl.cs 演示如何按名称查询队列 URL,并处理队列不存在的异常:
var client = new AmazonSQSClient(); string queueName = "New-Example-Queue"; try { var response = await client.GetQueueUrlAsync(queueName); if (response.HttpStatusCode == System.Net.HttpStatusCode.OK) { Console.WriteLine($"The URL for {queueName} is: {response.QueueUrl}"); } } catch (QueueDoesNotExistException ex) { Console.WriteLine(ex.Message); Console.WriteLine($"The queue {queueName} was not found."); }注意代码注释中的提醒:如果队列与默认用户的 AWS 区域不同,需要在客户端构造函数中显式指定区域。
4.6 接收并删除消息(ReceiveMessage + DeleteMessage)
ReceiveDeleteExample.cs 展示"先接收、后删除"的完整消费流程:
var receiveMessageRequest = new ReceiveMessageRequest { AttributeNames = { "SentTimestamp" }, MaxNumberOfMessages = 1, MessageAttributeNames = { "All" }, QueueUrl = queueUrl, VisibilityTimeout = 0, WaitTimeSeconds = 0, }; var receiveMessageResponse = await client.ReceiveMessageAsync(receiveMessageRequest); var deleteMessageRequest = new DeleteMessageRequest { QueueUrl = queueUrl, ReceiptHandle = receiveMessageResponse.Messages[0].ReceiptHandle, }; await client.DeleteMessageAsync(deleteMessageRequest);要点:
- 可见性超时(VisibilityTimeout):接收后消息进入"不可见"状态,防止多个消费者重复处理;处理完成必须调用
DeleteMessage并携带ReceiptHandle才能删除。 MessageAttributeNames = { "All" }表示返回全部消息属性,AttributeNames = { "SentTimestamp" }表示返回发送时间戳系统属性。
4.7 批量删除消息(DeleteMessageBatch)与删除队列(DeleteQueue)
SQSWrapper.cs 提供批量删除与队列清理:
// 批量删除:一次请求携带多个 ReceiptHandle var deleteRequest = new DeleteMessageBatchRequest() { QueueUrl = queueUrl, Entries = new List<DeleteMessageBatchRequestEntry>() }; foreach (var message in messages) { deleteRequest.Entries.Add(new DeleteMessageBatchRequestEntry() { ReceiptHandle = message.ReceiptHandle, Id = message.MessageId }); } var deleteResponse = await _amazonSQSClient.DeleteMessageBatchAsync(deleteRequest); return deleteResponse.Failed.Any(); // 删除队列 var deleteResponse = await _amazonSQSClient.DeleteQueueAsync( new DeleteQueueRequest() { QueueUrl = queueUrl }); return deleteResponse.HttpStatusCode == HttpStatusCode.OK;DeleteMessageBatch每个条目需提供ReceiptHandle与唯一Id,返回结果中Failed列表可用于重试失败条目。- 批量操作可显著降低 API 调用次数与成本,是清理积压消息的首选方式。
4.8 接收消息时的长轮询(ReceiveMessage)
场景代码中的 ReceiveMessagesByUrl 展示了长轮询(Long Polling)配置:
// Setting WaitTimeSeconds to non-zero enables long polling. var messageResponse = await _amazonSQSClient.ReceiveMessageAsync( new ReceiveMessageRequest() { QueueUrl = queueUrl, MaxNumberOfMessages = maxMessages, WaitTimeSeconds = 1 }); return messageResponse.Messages;WaitTimeSeconds设为非零值即启用长轮询:当队列为空时,请求会挂起等待消息到达,从而减少空轮询造成的请求次数与成本。MaxNumberOfMessages单次最多可请求 10 条。
五、场景一:向队列发布消息(Publish messages to queues)
这是 README 中列出的核心场景,代码位于 TopicsAndQueues.cs,完整演示 SNS 与 SQS 的集成消息流。该场景可完成以下任务:
- 创建主题(Topic),可选 FIFO 或非 FIFO。
- 创建若干 SQS 队列并订阅到主题,可为订阅添加过滤器(Filter)。
- 向主题发布消息。
- 轮询各队列,验证消息是否被正确投递。
5.1 场景执行流程
从源码注释与RunScenario方法可还原出完整调用链:
1. CreateTopic —— 创建 SNS 主题(FIFO 或非 FIFO) 2. CreateQueue —— 创建 SQS 队列 3. GetQueueAttributes —— 获取队列 ARN 4. SetQueueAttributes —— 设置队列策略,允许接收 SNS 消息 5. Subscribe —— 将队列订阅到主题 6. Publish —— 向主题发布消息 7. ReceiveMessage —— 轮询队列接收消息 8. DeleteMessageBatch —— 批量删除已处理消息 9. DeleteQueue —— 删除队列 10. Unsubscribe —— 取消订阅 11. DeleteTopic —— 删除主题场景运行时会通过控制台交互引导用户选择是否使用 FIFO 主题、是否启用基于内容的去重(Content-based deduplication)、是否设置订阅过滤器(按tone属性过滤,可选值cheerful、funny、serious、sincere)等。
5.2 FIFO 主题与去重机制
场景代码在SetupTopic中详细说明了 FIFO 语义(TopicsAndQueues.cs):
- FIFO 主题按序投递消息,支持去重与消息过滤。
- 使用 FIFO 主题时,主题名必须以
.fifo结尾。 - 去重 ID(Deduplication ID)可以显式设置,也可以启用内容去重(Content-based deduplication),由 SNS 对消息内容做哈希自动生成。
- 在 5 分钟去重窗口内,相同去重 ID 的消息会被接受但不会重复投递。
发布消息时(PublishMessages),FIFO 场景强制要求设置消息组 ID(MessageGroupId)——同一组内的消息按发布顺序接收:
Console.WriteLine("Because you are using a FIFO topic, you must set a message group ID." + "\r\nAll messages within the same group will be received in the order " + "they were published."); var messageGroupId = GetUserResponse("Enter a message group ID for this message:", "1"); if (!_useContentBasedDeduplication) { Console.WriteLine("Because you are not using content-based deduplication, " + "you must enter a deduplication ID."); deduplicationId = GetUserResponse("Enter a deduplication ID for this message.", "1"); } var messageID = await SnsWrapper.PublishToTopicWithAttribute( _topicArn, message, "tone", toneAttribute, deduplicationId, messageGroupId);5.3 订阅过滤器
SetupFilters与CreateFilterPolicy方法演示了如何为队列订阅附加过滤策略(TopicsAndQueues.cs):
var filters = new Dictionary<string, List<string>> { { "tone", filterSelections } }; string filterPolicy = JsonSerializer.Serialize(filters);序列化后的过滤策略形如{"tone":["cheerful","funny"]},只有带匹配tone属性的消息才会被投递到该队列。这实现了同一主题下的消息分流——不同队列可以订阅不同属性的消息。
5.4 清理资源
场景在结束或异常时都会执行CleanupResources:依次询问并删除队列、取消订阅、删除主题,避免遗留资源产生持续费用(TopicsAndQueues.cs)。这也是所有云资源示例应遵循的最佳实践。
六、场景二:AWS Message Processing Framework for .NET
第二个场景演示如何使用 AWS Message Processing Framework for .NET 构建发布与消费 SQS 消息的应用。该框架通过AWS.Messaging命名空间对 SQS 的底层细节(轮询、反序列化、路由)进行封装,开发者只需声明"消息类型 → 处理器"的映射即可。
项目结构包含两个应用与一个测试项目:
- Publisher/Program.cs:基于 ASP.NET Core 最小 API 的 Web 应用,通过 HTTP 端点接收请求并发布 SQS 消息。
- Handler/Program.cs:命令行应用,从 SQS 队列轮询并处理消息。
- Tests/HandlerTests.cs:处理器单元测试。
6.1 发布端(Publisher)
发布端通过AddAWSMessageBus注册消息总线,并将GreetingMessage类型映射到指定队列和消息标识符:
builder.Services.AddAWSMessageBus(builder => { // Check for input SQS URL. if ((args.Length == 1) && (args[0].Contains("https://sqs."))) { // 将 GreetingMessage 类型发布到指定队列,使用 "greetingMessage" 标识符 builder.AddSQSPublisher<GreetingMessage>(args[0], "greetingMessage"); } });随后定义一个 HTTP 端点,调用IMessagePublisher.PublishAsync发布消息(Publisher/Program.cs):
app.MapPost("/greeting", async ([FromServices] IMessagePublisher publisher, Publisher.GreetingMessage message) => { return await PostGreeting(message, publisher); }) .WithName("SendGreeting") .WithOpenApi(); public static async Task<IResult> PostGreeting(GreetingMessage greetingMessage, IMessagePublisher messagePublisher) { if (greetingMessage.SenderName == null || greetingMessage.Greeting == null) { return Results.BadRequest(); } await messagePublisher.PublishAsync(greetingMessage); return Results.Ok(); }注意队列 URL 通过命令行参数传入,参数需以https://sqs.开头才会注册发布器——这也是框架支持多队列/多类型映射的扩展点。
6.2 消费端(Handler)
消费端注册 SQS 轮询器,并声明消息类型与处理器的映射关系:
builder.ConfigureServices(services => { services.AddAWSMessageBus(builder => { if ((args.Length == 1) && (args[0].Contains("https://sqs."))) { // 注册轮询指定队列 builder.AddSQSPoller(args[0]); // 将 "greetingMessage" 类型消息反序列化为 GreetingMessage, // 并交给 GreetingMessageHandler 处理 builder.AddMessageHandler<GreetingMessageHandler, GreetingMessage>("greetingMessage"); } }); });处理器实现IMessageHandler<T>接口,通过MessageEnvelope<T>拿到消息内容后返回处理状态:
public class GreetingMessageHandler : IMessageHandler<GreetingMessage> { public Task<MessageProcessStatus> HandleAsync( MessageEnvelope<GreetingMessage> messageEnvelope, CancellationToken token = default) { Console.WriteLine( $"Received message {messageEnvelope.Message.Greeting} from {messageEnvelope.Message.SenderName}"); return Task.FromResult(MessageProcessStatus.Success()); } }框架核心价值在于:
- 消息路由:发布端与消费端使用相同的消息标识符
"greetingMessage"完成解耦,无需双方共享任何代码。 - 类型安全:消息自动反序列化为强类型对象(
GreetingMessage),无需手动解析 JSON。 - 可测试性:
IMessageHandler是纯接口抽象,单元测试可直接构造MessageEnvelope<T>调用处理器。仓库中的 HandlerTests.cs 正是这样做的:构造GreetingMessage与MessageEnvelope,调用HandleAsync后断言response.IsSuccess为真。
七、运行示例与测试
7.1 构建与运行
通用运行步骤参见 dotnetv3/README.md:
- 进入包含.sln文件的目录,执行
dotnet build SOLUTION.sln。 - 进入包含代码示例与.csproj文件的目录。
- 执行
dotnet run运行项目。
部分项目包含settings.json文件,编译前可按自己的账户与资源修改其中的值;也可以添加settings.local.json存放本地配置,应用运行时会被自动加载。示例编译后,也可以在 IDE 中直接运行。
对于需要传入参数的示例(如 Message Processing Framework 的 Publisher/Handler),通过命令行或 Debug launch profile 传入 SQS 队列 URL:
dotnet run -- "https://sqs.<region>.amazonaws.com/<account-id>/<queue-name>"7.2 运行测试
⚠️ 运行测试可能产生 AWS 账户费用。
测试运行说明同样参见 dotnetv3/README.md。在包含测试项目的目录中执行:
dotnet test如需更详细的输出:
dotnet test -l "console;verbosity=detailed"仓库测试按类别区分单元测试与集成测试,可单独筛选:
dotnet test --filter Category=Unit -l "console;verbosity=detailed" dotnet test --filter Category=Integration -l "console;verbosity=detailed"八、代码地图与延伸阅读
以下是本文涉及的核心文件索引,方便继续深入阅读:
| 用途 | 文件 |
|---|---|
| SQS 单操作封装(创建/属性/策略/接收/批量删除/删队列) | SQSWrapper.cs |
入门示例ListQueues | HelloSQS.cs |
| 创建队列 + 发送消息 | CreateSendExample.cs |
| 获取队列 URL | GetQueueUrl.cs |
| 接收并删除消息 | ReceiveDeleteExample.cs |
| SNS→SQS 消息发布场景 | TopicsAndQueues.cs |
| Message Processing Framework 发布端 | Publisher/Program.cs |
| Message Processing Framework 消费端 | Handler/Program.cs |
| 处理器单元测试 | HandlerTests.cs |
| .NET 示例总览与构建/测试指南 | dotnetv3/README.md |
延伸阅读资料:Amazon SQS 开发者指南、Amazon SQS API Reference、SDK for .NET 的 SQS API 参考(均可在 AWS 官方文档站检索获取)。仓库中 SNS 侧的配套封装位于 TopicsAndQueues 跨服务目录,可与本文 SQS 操作对照学习完整消息链路。
Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. SPDX-License-Identifier: Apache-2.0
- 示例工程
- 教程
- 后端
【免费下载链接】aws-doc-sdk-examples
Welcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below.
相关推荐
AWS SDK for Go 实现 Amazon SQS 队列与消息操作实战指南(go/sqs 示例全解)
AWS SDK for Go 实现 Amazon SQS 队列与消息操作实战指南(go/sqs 示例全解) 导读 本文围绕 aws doc sdk exampl
示例工程教程后端AWS SDK for .NET 操作 Amazon S3 完整指南:从基础动作到高级场景实战
AWS SDK for .NET 操作 Amazon S3 完整指南:从基础动作到高级场景实战 本文以开源仓库 aws doc sdk examples htt
示例工程教程后端使用 AWS SDK for C++ 操作 Amazon SQS:完整代码示例与实战指南
使用 AWS SDK for C++ 操作 Amazon SQS:完整代码示例与实战指南 导读 本文以 AWS 文档代码示例仓库(aws doc sdk exa
示例工程教程后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考