Amazon SAA-C03: 应用集成、消息传递与流处理 — 学习指南
属于 AWS SAA-C03 — 完整学习指南. 使用经过验证的答案练习: Amazon 考试中心, 或参加限时模拟考试: ExamRoll.io.
Amazon SQS:解耦、排序与交付语义
Amazon SQS 是一种完全托管的、基于拉取的消息队列,其主要架构作用是解耦生产者和消费者。同步调用链会将生产者的延迟和可用性与每个下游依赖项绑定在一起;而插入一个 SQS 队列则能将其转换为异步交接。生产者根据流量到达的速率入队,而消费者则以其能够安全处理的速率出队。这是解决写入突发问题的经典方案,否则这些突发流量可能会压垮 RDS 实例——队列会吸收突发流量,而一个有界限的消费者集群会以受控的并发度来消耗队列,从而控制数据库的连接数。
队列有两种类型,类型的选择决定了吞吐量和交付保证。
| 特性 | 标准队列 (Standard) | FIFO 队列 |
|---|---|---|
| 排序 | 尽力而为 | 严格,按 MessageGroupId 划分 |
| 交付 | 至少一次(可能出现重复) | 在 5 分钟去重窗口内精确一次 |
| 吞吐量 | 几乎无限 | 300 TPS(批处理时为 3,000);启用高吞吐量模式可达 70,000 |
| 队列名称 | 任意 | 必须以 .fifo 结尾 |
标准队列以至少一次的方式交付消息,且仅提供尽力而为的排序。当消费者未能在可见性超时到期前删除消息,或当分布式后端在分片间重放消息时,可能会出现重复消息。即使你在测试中从未见过重复消息,该服务在架构上也被允许重新交付,尤其是在代理故障转移期间。假设标准队列“通常”只会交付一次是一种设计缺陷,而非运营风险——在大规模应用中,重复消息最终必然会出现。因此,消费者逻辑必须是幂等的:在 DynamoDB 中使用条件写入来跟踪 MessageId 或业务键,在下游 API 调用中使用幂等性密钥,或依赖更新插入 (upsert) 语义。
FIFO 队列在 MessageGroupId 内部提供严格的排序,并通过 MessageDeduplicationId(可以是显式提供,也可以是消息正文的 SHA-256 哈希值)实现精确一次处理,该 ID 能在 5 分钟的窗口期内抑制重复消息。MessageGroupId 是关键概念:共享同一个组 ID 的所有消息都会严格按顺序一次交付给单个消费者,而不同组 ID 的消息可以并行处理。对于一个订单处理系统,如果每个客户的事件必须按顺序处理,但不同客户之间是独立的,则应使用 MessageGroupId = customerId。为所有消息使用单一的组 ID 会序列化整个工作负载,从而摧毁吞吐量。当用户重新提交一个卡顿的结账请求时,去重 ID 是防止重复创建订单的正确原语:客户端生成一个确定性的幂等性令牌(一个与结账会话绑定的 UUID),SQS 会丢弃在窗口期内到达的任何重复提交。
PaymentsQueue:
Type: AWS::SQS::Queue
Properties:
QueueName: payments.fifo
FifoQueue: true
ContentBasedDeduplication: true
DeduplicationScope: messageGroup
FifoThroughputLimit: perMessageGroupId
VisibilityTimeout: 60
RedrivePolicy:
deadLetterTargetArn: !GetAtt PaymentsDLQ.Arn
maxReceiveCount: 5
当需求明确指出“需要排序”或“不允许重复”时,选择标准队列是典型的错误。任何应用程序逻辑都无法恢复队列从未保留的顺序,因为来自不同后端主机的消息会交错到达。当工作负载要求排序(如交易分类账、状态机转换)或精确一次语义(如支付捕获、库存扣减)时,应选择 FIFO 队列。
可见性超时、毒丸消息与负载限制
当消费者收到一条消息时,SQS 会在可见性超时期间(默认为 30 秒,最长 12 小时)使该消息对其他消费者不可见。如果消费者在超时到期前删除了消息,它就消失了;如果未能删除——因为消费者崩溃了,或者处理时间过长——消息会重新出现并被再次交付。将可见性超时设置得比实际处理时间短是导致重复处理的主要原因:一个处理需要 45 秒的 Lambda 函数,如果其队列的可见性超时是默认的 30 秒,那么每条消息都将被至少重复处理两次。应将超时时间设置为至少 p99 的处理时间(AWS 对 Lambda 驱动的队列的指导建议是至少 6 倍的函数超时时间),对于处理时长不确定的作业,应动态延长超时时间:
sqs.change_message_visibility(
QueueUrl=queue_url,
ReceiptHandle=handle,
VisibilityTimeout=300 # extend by 5 minutes
)
死信队列 (DLQ) 用于捕获毒丸消息。源队列上的 RedrivePolicy 指定一个 maxReceiveCount(通常为 3-5 次);一旦超过该次数,SQS 会将消息移动到 DLQ 以供离线检查。DLQ 必须与源队列的类型匹配(FIFO ↔ FIFO)。如果没有 DLQ,格式错误的消息会无限循环,这在 FIFO 队列中尤其具有破坏性——排序机制会阻止同一组中的后续消息被交付,直到问题消息被处理掉,因此一条坏消息会阻塞整个组。
SQS 消息的上限为 256 KB。对于更大的负载——例如,一个携带渲染后文档的作业——应使用 SQS 扩展客户端库,它会将负载写入 S3,而只将存储桶/键的引用入队。消费者库在接收时会透明地获取负载。不要将负载拆分到多个消息中(这样会失去原子性和顺序性),也不要期望通过 Base64 编码一个 2 MB 的二进制大对象能使其符合大小限制。
队列驱动的自动扩展
对于位于 SQS 队列之后的 EC2 或 ECS 上的消费者集群,正确的扩展信号不是 CPU,而是队列积压。CPU 指标滞后于消息到达率,并且会误将一个饱和的消费者解读为“繁忙但能应对”。典型的扩展指标是 ApproximateNumberOfMessagesVisible,但直接根据原始队列深度进行扩展粒度太粗。推荐的方法是使用每个实例的积压量这个自定义指标:
backlogPerInstance = ApproximateNumberOfMessagesVisible / RunningInstances
将此指标发布到 CloudWatch,并用它来驱动 Auto Scaling 组或 ECS 服务的追踪目标策略,以便每个工作单元都维持一个有界限的积压量(例如 10 条消息)。这可以在突发期间实现平滑的横向扩展,并防止在队列深度较小但消费者已经饱和时出现振荡。对于缩容,可以结合 ApproximateAgeOfOldestMessage 指标,以避免在仍有旧消息残留时终止容量。
Amazon SNS:扇出、过滤和跨账户交付
SNS 是一种基于推送的发布/订阅服务。发布者向主题写入消息;SNS 将消息推送给每个订阅:SQS 队列、Lambda 函数、HTTP(S) 端点、电子邮件、短信、Kinesis Data Firehose 或移动推送。主流的持久化模式是 SNS → SQS 扇出:一个主题有多个 SQS 队列订阅,这样每个下游服务都有自己的持久化缓冲区、重试策略和死信队列 (DLQ),而发布者只需知道主题即可。如果一个消费者服务宕机数小时,其队列会累积消息并在恢复后进行处理——单独使用 SNS 缺乏这种缓冲能力,并且会耗尽其重试策略。
Producer ──▶ SNS topic ──┬──▶ SQS Queue A ──▶ Service A
├──▶ SQS Queue B ──▶ Service B
└──▶ SQS Queue C ──▶ Service C
消息过滤允许每个订阅声明一个 JSON 过滤策略,这样 SNS 只交付匹配的消息,从而避免了每个消费者都接收所有消息然后在客户端进行过滤的反模式:
{
"eventType": ["order_placed", "order_cancelled"],
"region": ["us-east-1", "us-west-2"]
}
有两个行为特性很重要。首先,标准 SNS 主题不保证消息的顺序——每个订阅者的重试计时器和独立的网络路径使得顺序重排成为常态。如果顺序很重要,请使用订阅了 SQS FIFO 队列的 SNS FIFO 主题;消息组 ID 会在整个过程中传播。否则,订阅者必须是幂等的,并且能够容忍顺序重排。其次,HTTP(S) 订阅会根据交付策略进行重试(默认是:立即重试三次,然后进行长达一小时的指数退避,最后丢弃)。订阅者必须在 15 秒内以 2xx 状态码响应,验证 x-amz-sns-message-type 签名,并且——对于不可靠的端点——始终配置一个 SNS 死信队列(DLQ)(将消息重新驱动到 SQS),以便捕获未送达的消息,而不是静默丢弃。
跨账户调用是一个常见的陷阱。当账户 A 向一个主题发布消息,该主题扇出到账户 B 中的一个 Lambda 时,需要两个策略:SNS 主题策略(或订阅方向)必须允许该订阅,并且 Lambda 的基于资源的策略必须允许来自 sns.amazonaws.com 的 lambda:InvokeFunction 操作,并带有与该主题匹配的 SourceArn 条件。缺少 Lambda 资源策略是最常见的失败模式——订阅看起来是健康的,但调用却被 403 拒绝。如果主题使用客户管理的 KMS 密钥加密,则密钥策略还必须向发布主体和 sns.amazonaws.com 授予 kms:Decrypt 和 kms:GenerateDataKey 权限。
{
"Effect": "Allow",
"Principal": {"Service": "sns.amazonaws.com"},
"Action": "lambda:InvokeFunction",
"Resource": "arn:aws:lambda:us-east-1:222222222222:function:ProcessOrder",
"Condition": {"ArnLike": {"AWS:SourceArn": "arn:aws:sns:us-east-1:111111111111:orders"}}
}
Amazon EventBridge:路由事件总线
EventBridge(前身为 CloudWatch Events)通过基于内容的路由、Schema 发现、SaaS 合作伙伴事件源以及归档/重播功能,扩展了发布/订阅模型。事件流经事件总线(默认、自定义或合作伙伴总线),并与规则进行匹配,规则的事件模式会根据 JSON 结构进行过滤。规则可以通过输入路径和输入模板转换负载,附加死信目标,并交付到超过 20 种原生目标,包括 Lambda、Step Functions、ECS 任务、SQS、SNS、Kinesis 和 API 目的地。
{
"source": ["com.acme.orders"],
"detail-type": ["OrderPlaced"],
"detail": {"amount": [{"numeric": [">", 500]}]}
}
它与 SNS 的区别在于架构层面。SNS 专为对同构订阅者进行高吞吐量广播而优化,具有简单的属性过滤和较低的延迟。EventBridge 则专为异构事件驱动架构而优化:许多生产者发出不同的事件 Schema,消费者根据模式而非主题进行订阅。对于正在分解为微服务的单体应用——特别是当生产者包括 SaaS 合作伙伴或原生发出事件的 AWS 服务(Config、GuardDuty、CodePipeline、CloudTrail)时——EventBridge 通常是正确的选择。对于向同构订阅者进行极高容量、低延迟的扇出,SNS 仍然胜出,因为 EventBridge 的单事件延迟稍高,且默认吞吐量上限较低。
Amazon MQ:用于现有协议的代理消息传递
Amazon MQ 是一个运行 ActiveMQ 或 RabbitMQ 的托管消息代理。它的存在是为了迁移那些依赖 AMQP 0-9-1、AMQP 1.0、MQTT、STOMP、OpenWire 或 JMS 的本地工作负载,而无需重写应用程序。如果一个支付系统使用具有事务性精确一次语义的第三方 JMS 代理,将其迁移到 Amazon MQ 可以保留线路协议和交付保证,同时消除基础设施管理。对于全新的 AWS 原生设计,应选择 SQS/SNS/EventBridge;仅当协议兼容性是约束条件时,才选择 Amazon MQ。
Kinesis Data Streams
Kinesis Data Streams (KDS) 是一个用于高吞吐量流式摄取的持久、有序、分区日志——例如点击流、物联网遥测、日志聚合。记录通过 PartitionKey 被放入分片中;顺序只在分片内部得到保证,而不在整个流中保证。每个分片支持 1 MB/s 或 1,000 条记录/秒的写入,以及 2 MB/s 的读取(使用增强型扇出时更高)。记录默认保留 24 小时,可延长至 365 天,因此多个独立的消费者可以重播相同的历史记录——这一点 SQS 做不到,因为 SQS 在确认后会删除消息。
按需模式通过自动扩展至每个流 200 MiB/s 的写入能力,免去了分片计算的麻烦,非常适合不可预测的流量。在容量已知的稳定状态下,预置模式更便宜。
当工作负载要求有序、可重播的摄取,且吞吐量超出了 FIFO 的能力(FIFO 的上限远低于 KDS 处理的每秒数百万条记录),或者当多个独立消费者必须读取同一个流,或者当需求中提到“在整个处理过程中保持原始顺序”并伴有高数据量时,应选择 KDS 而非 SQS FIFO。
Kinesis Data Firehose
Kinesis Data Firehose 是一项完全托管的交付服务。它从 Kinesis 流或直接 PUT 读取数据,按大小(1–128 MB)或时间(60–900 秒,以先达到的为准)进行缓冲,可选择性地调用 Lambda 进行逐条记录转换(PII 清理、格式规范化),可以使用 Glue 架构动态地将 JSON 转换为 Parquet 或 ORC 格式,通过 KMS 加密,并交付到 S3、Redshift、OpenSearch 或 Splunk。它没有分片,无需运行消费者,并采用按 GB 付费的定价模式。
将数据可扩展地采集到数据湖的典型模式是将 Data Streams(按需模式)作为持久化缓冲区,并与 Firehose 配对以交付到 S3:
Producers → Kinesis Data Streams (on-demand) → Firehose (60s buffer, Parquet) → S3 → Athena/Glue
对于采集数百万移动端事件、对其进行加密并以 Parquet 格式存入 S3 的场景,正确的答案是使用 Firehose 并配置 Parquet 转换和一个 KMS 密钥——而不是 KDS 加上自定义消费者再加上手写的 Parquet 写入程序,后者的代码和基础设施要多得多。Firehose 是近乎实时的,不支持消费者端重放;当需要重放时,应在路径中保留 KDS。
Kinesis Data Analytics(现为 Managed Service for Apache Flink)针对流运行 SQL 或 Flink 作业,以进行窗口化聚合。
Lambda 集成和重试语义
Lambda 与这些服务集成时,重试行为有本质上的不同:
| 源 | 批处理 | 排序 | 失败时 |
|---|---|---|---|
| SQS Standard | 最多 10,000 条消息 | 无 | 可见性超时后返回;达到 maxReceiveCount 后进入 DLQ |
| SQS FIFO | 按组 | 按组 | 组被阻塞,直到成功或进入 DLQ |
| Kinesis Streams | 最多 10,000 条记录 | 按分片 | 重试会阻塞分片,直到成功、记录过期或达到 MaximumRetryAttempts/OnFailure 目的地 |
| Firehose | 不适用(转换) | 不适用 | 失败的记录存入 S3 的错误前缀中 |
对于 SQS,应保持 Lambda 函数超时 ≤ 队列可见性超时,并将可见性超时设置为至少是函数超时的 6 倍。对于 Kinesis,应启用 BisectBatchOnFunctionError 并配置一个 OnFailure 目的地(SQS 或 SNS),这样单个毒丸记录就不会无限期地阻塞整个分片。
选型决策表
| 需求 | 正确选择 | 替代方案为何不适用 |
|---|---|---|
| 有序、恰好一次的应用消息传递,运维最少 | SQS FIFO | 标准 SQS 缺乏排序/去重功能;MQ 增加了代理管理负担 |
| 保留现有的 AMQP/JMS/MQTT 客户端 | Amazon MQ | SQS/SNS 使用专有 API |
| 将一个事件持久地扇出到多个 AWS 消费者 | SNS → 多个 SQS | 生产者与消费者直接耦合会重新引入单体问题;仅使用 SNS 在消费者宕机时会丢失消息 |
| 使用筛选器/转换路由异构事件 | EventBridge | SNS 筛选策略缺乏转换功能、合作伙伴源和架构注册表 |
| 向相同的订阅者进行超高吞吐量的扇出 | SNS | EventBridge 延迟更高,默认吞吐量更低 |
| 采集并重放海量有序流 | Kinesis Data Streams | SQS 保留期上限为 14 天,且无法按偏移量重放 |
| 无需代码即可将流交付到 S3/Redshift/OpenSearch | Firehose | 单独使用 Data Streams 需要一个消费者应用程序 |
| 在 S3 中将流式 JSON 转换为 Parquet | 使用带 Glue 架构的 Firehose | 自定义 KDS 消费者需要编写/运维一个 Parquet 写入程序 |
← 分析、数据湖、ML 与专用工作负载 · 所有领域 · 安全、IAM、KMS 与治理 →
练习这些题目 → · 在 ExamRoll.io 上限时练习 →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
通过考试 →