- 示例工程
- 教程
- 后端
【免费下载链接】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 官方代码示例仓库中kotlin/usecases/topics_and_queues示例的完整技术指南。该示例基于 AWS SDK for Kotlin,通过命令行交互程序演示如何创建 Amazon SNS 主题、将 Amazon SQS 队列订阅到主题、启用 FIFO(先进先出)队列、配置基于属性的过滤订阅,并发布/接收消息的完整链路。阅读本文后,你将掌握 SNS 与 SQS 组合使用时的核心配置逻辑、FIFO 去重与消息分组的原理,以及使用 Kotlin 调用SnsClient与SqsClient的完整代码实现。
一、发布/订阅架构:为什么需要 SNS + SQS 组合
发布/订阅(Publish and Subscribe)是一种信息传递机制,广泛用于社交媒体以及软件内部的模块解耦。生产者(Producer)发布消息,订阅者(Subscriber)接收消息。在这一模型中,消息生产者与消息消费者被解耦,因此消息传递更加灵活、健壮。
该示例不会构建一个完整的端到端应用程序,而是提供一个可以亲手"把玩"发布/订阅架构的交互式命令行程序,帮助你理解两种 AWS 消息服务的分工差异:
- Amazon SNS 是推送(Push)服务:它主动将消息推送到各类端点,如电子邮件地址、移动应用端点或 SQS 队列。生产者只负责把消息发布到主题(Topic),由 SNS 负责投递。
- Amazon SQS 是轮询(Poll)服务:订阅者通过调用"接收消息"API(ReceiveMessage)主动轮询队列。任何代码都可以轮询队列;消息会一直保存在队列中,直到你显式删除它,这给消息的处理方式带来了更大的灵活性。
单独使用 SNS 即可实现发布/订阅,但将 SNS 与 SQS 组合使用,可以获得更灵活的消息消费方式:SNS 负责推送与过滤,SQS 负责缓存与异步消费。
二、FIFO 主题与队列的核心概念
FIFO(先进先出)主题
示例程序启动后,首先询问是否使用 FIFO 主题:
Would you like to work with FIFO topics? (y/n)FIFO 主题在创建时配置,启用后还会解锁其他选项(去重、消息分组、消息过滤等)。FIFO 主题保证消息严格有序投递。
消息去重(Deduplication)
Use content-based deduplication instead of a deduplication ID? (y/n)去重仅对 FIFO 主题可用。去重机制防止订阅者针对被判定为重复的事件响应多次。规则如下:
- 若一条消息发布到 SNS FIFO 主题后,在5 分钟去重时间窗口内发现与已有消息具有相同的去重 ID,则该消息会被接收但不会被投递。
- 基于内容的去重(Content-based deduplication):使用消息内容的哈希值作为去重 ID。选择此项后,发布消息时无需手动提供去重 ID。
- 显式去重 ID:若未启用基于内容的去重,则每条消息都必须携带去重 ID(Deduplication ID)。
消息分组(Message Grouping)
对于 FIFO 主题,每条消息必须携带消息组 ID(Group ID)。组 ID 决定消息在主题内的排序范围,相同组 ID 的消息按发布顺序严格排列。组 ID 最多 128 个字符,可包含字母数字字符(a-z, A-Z, 0-9)以及标点符号(!"#$%&'()*+,-./:;<=>?@[\]^_{|}~)`。
消息去重 ID 的字符限制与组 ID 相同,同样最多 128 个字符。
命名规则
示例对主题与队列名称有明确约束:
| 对象 | 长度 | 允许字符 | 特殊规则 |
|---|---|---|---|
| SNS 主题名称 | 1–256 字符 | 大小写 ASCII 字母、数字、下划线、连字符 | 选择 FIFO 后程序自动追加.fifo后缀(FIFO 主题必需) |
| SQS 队列名称 | 1–80 字符 | 大小写 ASCII 字母、数字、下划线、连字符 | 选择 FIFO 后程序自动追加.fifo后缀(FIFO 队列必需) |
三、示例的完整交互流程
程序通过标准输入(Scanner)逐项收集配置,全程共 11 个步骤。以下按 README 与实际源码 SNSWorkflow.kt 还原完整界面。
1. 创建 SNS 主题
Enter a name for your SNS topic:若前面选择了 FIFO 主题,程序会自动在名称后追加.fifo。选择 FIFO 后,程序会立即询问去重方式并收集组 ID / 去重 ID(源码第 90–113 行):
- 选择基于内容去重(y):提示
Enter a group id value - 选择显式去重 ID(n):依次提示
Enter deduplication Id value与Enter a group id value
2. 创建 SQS 队列
Enter a name for an SQS queue.为每个订阅者单独建队列是有益的——你可以为不同订阅者定制消息消费方式与消息过滤规则。FIFO 场景下队列名称同样自动追加.fifo后缀。
3. 为订阅添加过滤(仅 FIFO 场景)
Filter messages for "<queue name>.fifo"s subscription to the topic "<topic name>.fifo"? (y/n)若选择为订阅添加过滤器,则按预定属性集合过滤:
You can filter messages by one or more of the following "tone" attributes. 1. cheerful 2. funny 3. serious 4. sincere Enter a number (or enter zero to stop adding more).可连续选择多个tone属性,输入0结束选择。
4. 发布消息
Enter a message text to publish.之后进入发布阶段(FIFO 场景):
- 询问是否为消息附加
tone属性(Add an attribute to this message? (y/n)),随后从cheerful / funny / serious / sincere中选择一个; - 发布到主题。
README 中描述了重复发布选项Post another message? (y/n),允许连续发布多条消息;当前仓库源码的main流程在发布一条消息后即进入接收阶段(SNSWorkflow.kt 第 207–234 行),读者可按需自行扩展循环发布逻辑。
5. 接收、显示并清理
发布结束后,程序轮询队列并显示消息内容(Message Id 与完整 Body),随后依次删除已接收消息、退订、删除队列与主题,完成资源的完整回收。
四、源码实现深度解析
4.1 创建主题:FIFO 与非 FIFO 的分支
非 FIFO 主题直接通过CreateTopicRequest创建(SNSWorkflow.kt 第 568–577 行):
suspend fun createSNSTopic(topicName: String?): String? { val request = CreateTopicRequest { name = topicName } SnsClient { region = "us-east-1" }.use { snsClient -> val result = snsClient.createTopic(request) return result.topicArn } }FIFO 主题则在CreateTopicRequest中通过attributes声明两个关键属性(第 579–597 行):
suspend fun createFIFO(topicName: String?, duplication: String): String? { val topicAttributes: MutableMap<String, String> = HashMap() if (duplication.compareTo("n") == 0) { topicAttributes["FifoTopic"] = "true" topicAttributes["ContentBasedDeduplication"] = "false" } else { topicAttributes["FifoTopic"] = "true" topicAttributes["ContentBasedDeduplication"] = "true" } ... }| 主题属性 | 值 | 含义 |
|---|---|---|
FifoTopic | true | 启用 FIFO,消息严格有序 |
ContentBasedDeduplication | true/false | 是否基于消息内容哈希自动去重 |
可见duplication选项直接映射到ContentBasedDeduplication属性,这就是"选择内容去重后无需再输入去重 ID"的底层原因。
4.2 创建队列:FIFO 队列属性
createQueue根据selectFIFO分支创建队列(第 527–566 行)。FIFO 队列通过QueueAttributeName.FifoQueue属性声明:
if (selectFIFO) { val attrs = mutableMapOf<String, String>() attrs[QueueAttributeName.FifoQueue.toString()] = "true" val createQueueRequest = CreateQueueRequest { queueName = queueNameVal attributes = attrs } ... }创建成功后,通过GetQueueUrlRequest获取队列 URL(后续所有 SQS 操作都以queueUrl为寻址依据)。
4.3 获取队列 ARN 并附加 IAM 策略
getSQSQueueAttrs通过GetQueueAttributesRequest请求QueueAttributeName.QueueArn属性,得到队列 ARN(第 506–525 行)。
要让 SNS 能把消息投递到 SQS 队列,必须为队列附加一条 IAM 资源策略(SetQueueAttributesRequest+Policy属性,第 491–504 行)。源码中策略内容如下:
{ "Statement": [ { "Effect": "Allow", "Principal": { "Service": "sns.amazonaws.com" }, "Action": "sqs:SendMessage", "Resource": "<SQS 队列 ARN>", "Condition": { "ArnEquals": { "aws:SourceArn": "<SNS 主题 ARN>" } } } ] }策略要点:
Principal限定为sns.amazonaws.com,即只允许 SNS 服务访问;Action仅授予sqs:SendMessage(最小权限);Condition通过ArnEquals将来源限定为当前主题的 ARN,防止其他主题向该队列投递。
测试代码 AWSSNSTest.kt 在非 FIFO 用例中展示了等价的 ARN 拼接写法(arn:aws:sqs:us-east-1:<accountId>:<queueName>)。
4.4 订阅队列并配置过滤策略
subQueue是示例的核心方法(第 437–489 行),通过SubscribeRequest完成订阅:
request = SubscribeRequest { protocol = "sqs" endpoint = queueArnVal returnSubscriptionArn = true topicArn = topicArnVal }关键参数:
protocol = "sqs":声明订阅端点是 SQS 队列;endpoint = queueArnVal:目标队列的 ARN;returnSubscriptionArn = true:立即返回订阅 ARN,供后续退订使用。
当用户选择了tone过滤属性时,示例通过 Gson 构造{"tone": [...]}JSON,再调用setSubscriptionAttributes将FilterPolicy写入订阅(第 468–485 行):
val attributeNameVal = "FilterPolicy" val jsonString = "{\"tone\": []}" val jsonObject = gson.fromJson(jsonString, JsonObject::class.java) val toneArray = jsonObject.getAsJsonArray("tone") for (value: String? in filterList) { toneArray.add(JsonPrimitive(value)) } val updatedJsonString: String = gson.toJson(jsonObject) val attRequest = SetSubscriptionAttributesRequest { subscriptionArn = result.subscriptionArn attributeName = attributeNameVal attributeValue = updatedJsonString } snsClient.setSubscriptionAttributes(attRequest)例如选择cheerful与funny后,最终写入的过滤策略为{"tone":["cheerful","funny"]}。此后只有携带匹配tone消息属性的消息才会被投递到该队列——过滤由 SNS 在服务端完成,队列只会收到符合规则的消息。
4.5 发布消息:普通与 FIFO 的差异
普通主题的发布(第 353–363 行):
val request = PublishRequest { message = messageVal topicArn = topicArnVal }FIFO 主题的发布(pubMessageFIFO,第 365–434 行)则依据"是否使用内容去重 / 是否附加消息属性"组合出四种分支:
- 内容去重 + 无属性:携带
messageGroupId即可; - 显式去重 ID + 无属性:携带
messageDeduplicationId与messageGroupId; - 内容去重 + 附加属性:携带
messageGroupId(注意:当前源码此分支未把messageAttributes附加到请求中,若你的业务依赖消息属性过滤可留意这一细节); - 显式去重 ID + 附加属性:同时携带
messageDeduplicationId、messageGroupId与messageAttributes。
消息属性通过MessageAttributeValue构造,声明数据类型为String:
val messAttr = aws.sdk.kotlin.services.sns.model.MessageAttributeValue { dataType = "String" stringValue = "true" } val mapAtt: Map<String, aws.sdk.kotlin.services.sns.model.MessageAttributeValue> = mapOf(msgAttValue to messAttr)mapOf(msgAttValue to messAttr)中的键(如cheerful)就是过滤策略匹配的依据,与订阅端FilterPolicy中的tone数组值一一对应。
4.6 接收与删除消息
receiveMessages(第 332–351 行)通过ReceiveMessageRequest轮询队列,maxNumberOfMessages = 5;当消息携带属性(msgAttValue非空)时额外设置waitTimeSeconds = 1启用短轮询等待:
val receiveRequest = ReceiveMessageRequest { queueUrl = queueUrlVal waitTimeSeconds = 1 maxNumberOfMessages = 5 }收到消息后,deleteMessages(第 312–330 行)将MessageId组装为DeleteMessageBatchRequestEntry,通过deleteMessageBatch批量删除,避免消息被重复消费。
4.7 收尾清理
示例在流程末尾依次执行(第 275–310 行):
unSub:调用UnsubscribeRequest退订;deleteSQSQueue:先按队列名解析 URL,再deleteQueue;deleteSNSTopic:调用DeleteTopicRequest删除主题。
整个生命周期从创建到删除闭环,可安全反复运行。
五、构建与运行
构建配置
示例使用 Gradle 管理依赖,构建脚本位于 kotlin/usecases/topics_and_queues/build.gradle.kts,核心配置如下:
plugins { kotlin("jvm") version "1.9.0" application } dependencies { implementation("aws.sdk.kotlin:sns:1.0.0") implementation("aws.sdk.kotlin:sqs:1.0.0") implementation("aws.smithy.kotlin:http-client-engine-okhttp:0.30.0") implementation("aws.smithy.kotlin:http-client-engine-crt:0.30.0") implementation("com.google.code.gson:gson:2.10.1") testImplementation("org.junit.jupiter:junit-jupiter:5.9.2") implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core:1.7.1") }依赖说明:
aws.sdk.kotlin:sns/aws.sdk.kotlin:sqs:AWS SDK for Kotlin 的 SNS、SQS 客户端;http-client-engine-okhttp/http-client-engine-crt:Smithy Kotlin 提供的 HTTP 传输引擎(CRT 引擎支持更长超时与流式场景);gson:用于构造FilterPolicyJSON;kotlinx-coroutines-core:SDK 的 suspend API 依赖协程;- 测试框架为 JUnit 5(
useJUnitPlatform()),JVM 目标版本为 17。
前置条件与运行
运行前需要完成环境准备:配置 AWS 凭据(SnsClient与SqsClient均显式指定region = "us-east-1"),并确保 IAM 权限覆盖 Amazon SNS 与 Amazon SQS 的对应操作。
在kotlin/usecases/topics_and_queues目录下执行:
gradle run程序启动后按上文交互流程逐步操作即可。注意:创建主题、队列等操作会在你的 AWS 账号中产生实际资源与费用,示例运行结束后会自动完成清理。
六、自动化测试验证
示例附带完整的端到端测试 AWSSNSTest.kt,通过 JUnit 5 按@Order顺序执行两个用例:
testWorkflowFIFO:走 FIFO 全链路——创建 FIFO 主题(显式去重 IDdup100、组 IDgroup100)、创建 FIFO 队列、附加 IAM 策略、订阅、发布消息、接收、删除、退订、清理;testWorkflowNonFIFO:走标准主题链路,验证非 FIFO 场景下的同名操作。
测试用(1..10000).random()生成随机后缀(如topic123.fifo、queue456.fifo),避免资源命名冲突;发布后使用delay(1000)等待消息投递,再轮询队列断言接收结果。这两个用例同时验证了 README 中两种分支路径的代码可运行性。
执行测试:
gradle test七、注意事项
- 运行该代码可能产生 AWS 账号费用(创建主题、队列及消息投递均按量计费),运行测试同样可能产生费用;
- 建议为代码授予最小权限(Least Privilege):仅授予完成任务所需的最低权限。本示例的策略即遵循此原则,只授予
sqs:SendMessage且限定来源 ARN; - 该代码未经所有 AWS 区域验证,如需在生产区域运行请以 AWS 区域服务可用性为准。
延伸阅读
- 本示例对应的入门文档:kotlin/usecases/topics_and_queues/README.md
- 完整主流程实现:SNSWorkflow.kt
- 端到端测试用例:AWSSNSTest.kt
- Gradle 构建脚本:build.gradle.kts
- 同一仓库中面向 .NET 的等价示例:dotnetv3/cross-service/TopicsAndQueues 可作跨语言对照参考
- 示例工程
- 教程
- 后端
【免费下载链接】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 C++ 实战:使用 SNS 主题与 SQS 队列实现带过滤与 FIFO 的发布订阅
AWS SDK for C++ 实战:使用 SNS 主题与 SQS 队列实现带过滤与 FIFO 的发布订阅 本篇文章基于 aws doc sdk example
示例工程教程后端AWS .NET 示例:用 SNS 主题 + SQS 队列实现发布订阅、FIFO 去重与过滤实战(TopicsAndQueues)
AWS .NET 示例:用 SNS 主题 + SQS 队列实现发布订阅、FIFO 去重与过滤实战(TopicsAndQueues) 导读 本文基于 AWS 官方
示例工程教程后端使用 Amazon SNS 主题与 Amazon SQS 队列实现发布订阅:Go 版 FIFO 与消息过滤实战指南
使用 Amazon SNS 主题与 Amazon SQS 队列实现发布订阅:Go 版 FIFO 与消息过滤实战指南 导读 本文基于 AWS SDK for Go
示例工程教程后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考