Amazon DVA-C02: 消息传递、流处理与事件驱动架构 (SNS, SQS, Kinesis, EventBridge, Step Functions) — 学习指南

属于 AWS Developer Associate DVA-C02 — 学习指南. 使用经过验证的答案练习: Amazon 考试中心, 或参加限时模拟考试: ExamRoll.io.

选择合适的消息和流处理原语

在 SNS、SQS(标准队列与 FIFO 队列)、Kinesis、EventBridge 和 Step Functions 之间做出选择,首先要从通信模式入手:发布/订阅、点对点、有序流、事件总线路由或工作流编排。SNS 是一个扇出(fan-out)模式的发布/订阅发布者;使用 Publish(SDK 调用:Publish/PublishBatch)发布消息,并让 SQS 端点、Lambda、HTTP/S 或移动端点进行订阅。SQS 是一个持久化的点对点缓冲区,具有 ReceiveMessage/DeleteMessage 语义;使用 CreateQueue 创建队列,并配置 VisibilityTimeout、ReceiveMessageWaitTimeSeconds(长轮询)、MessageRetentionPeriod 等属性,以及关联死信队列(DLQ)的重试策略(redrive policies)。FIFO 队列要求设置 FifoQueue=true,并使用 MessageGroupId 加上 MessageDeduplicationId(或 ContentBasedDeduplication)来实现排序和去重。Kinesis Data Streams 是基于分片(shard)的有序流服务;生产者调用 PutRecord/PutRecords,消费者使用 GetShardIterator(可指定 TRIM_HORIZON、LATEST、AT_SEQUENCE_NUMBER)然后调用 GetRecords。Kinesis Firehose 负责将数据托管交付到 S3/Redshift/OpenSearch,并提供 BufferingHints(SizeInMBs、IntervalInSeconds)和 Lambda 转换功能。EventBridge 使用 PutEvents 接收事件,并通过基于规则的筛选进行路由,支持 schema registry(模式注册)和跨账户事件总线。Step Functions 用于编排复杂流程;通过 StartExecution(标准工作流)或 StartSyncExecution(用于同步 Express 模式)启动执行,并与任务(Task)集成,例如

undefined

。当吞吐量、顺序性、交付保证、保留时间和编排需求之间存在冲突时,需要权衡这些服务的优缺点。

SQS 和 SNS 模式、去重及消费者扩展

当需要持久化解耦时,SQS 是首选;生产者实现 SendMessage/SendMessageBatch,消费者使用带有 WaitTimeSeconds 的 ReceiveMessage 来启用长轮询,以减少空接收次数。为实现严格排序和去重,使用 CreateQueue (FifoQueue=true) 创建一个 FIFO 队列,并设置 MessageGroupId 来定义有序分区;使用 MessageDeduplicationId 或启用 ContentBasedDeduplication,以便在去重窗口内抑制内容相同的消息。标准队列可能会传递重复消息——因此需要通过数据库条件写入(例如,使用 ConditionExpression attribute_not_exists(pk) 的 DynamoDB PutItem 操作)或在事务中使用 RDS 唯一约束和 upsert 操作,来确保消费者的幂等性。配置重试策略(redrive policies),在达到 maxReceiveCount 后将处理失败的消息路由到死信队列(DLQ);通过 GetQueueAttributes 监控 ApproximateNumberOfMessages 和 ApproximateNumberOfMessagesNotVisible 指标。Lambda 与 SQS 的集成使用 CreateEventSourceMapping:设置 BatchSize、MaximumBatchingWindowInSeconds,并启用 FunctionResponseTypes = [“ReportBatchItemFailures”] 以使用部分批处理失败报告语义,避免重新处理批次中已成功的记录。注意 FIFO 队列与 Lambda 集成的语义:消息组排序强制对每个 MessageGroupId 进行单线程处理,这限制了组内的并发性;可以通过分区到多个组 ID,或使用 SNS 将消息分发到多个队列以实现并行消费来扩展。同时要注意可见性超时(visibility timeout):当处理时间过长时,调用 ChangeMessageVisibility,否则可能导致消息被重复处理。

Kinesis Data Streams 和 Firehose:排序、保留与背压处理

Kinesis Data Streams 为流处理用例提供分片内有序和持久化保留功能。生产者使用一个映射到分片的 PartitionKey 来调用 PutRecord 或 PutRecords(批量);消费者调用 GetShardIterator 和 GetRecords,然后使用 KCL (Kinesis Client Library) 或自定义的 DynamoDB 检查点表来记录偏移量(checkpoint offsets)。默认保留期为 24 小时(可根据流的配置调整为更长的时间窗口,并在可用时使用扩展保留功能);使用 UpdateShardCount 规划分片数量,以匹配写入吞吐量和读取者并行度。消费者的扩展受到限制:单个 Lambda 事件源映射将一个分片映射到一个 Lambda 并发实例,因此需要增加分片来提高消费者并发度,或者启用增强型扇出(enhanced fan-out),通过 SubscribeToShard API(消费者注册)为每个消费者提供独立的 2 MB/秒管道和独立的扩展能力。使用 PutRecords 进行高效的批量写入;当消费者处理滞后时会产生背压(可通过监控 GetRecords.IteratorAgeMilliseconds 指标发现)。为应对流量尖峰,可以在 Kinesis 中缓冲数据或在前端使用 SQS,生产者应使用带指数退避的重试机制,并使用 KMS 对流进行加密以保护个人身份信息(PII)。Kinesis Data Firehose 简化了数据交付:配置 BufferingHints(SizeInMBs、IntervalInSeconds)、CompressionFormat 和一个用于数据转换的 Lambda 函数。Firehose 会自动处理到目标服务的重试/退避,并能将失败的记录写入一个备用的 S3 存储桶。一个常见的陷阱是分片预置不足:这会导致消费者“饥饿”和延迟飙升;应主动测量并进行扩展。

用于路由和编排的 EventBridge 和 Step Functions

EventBridge 擅长于基于 schema 的事件路由和跨账户/事件合作伙伴集成,它使用 PutEvents 注入事件,并使用 PutRule/PutTargets 将事件路由到 SQS、Lambda、Kinesis、Step Functions 或 HTTP 端点。EventBridge 使用事件模式进行筛选,并支持归档和重放以重建状态。为规则使用死信队列(带有 SqsParameters 或 DeadLetterConfig 的 Target),并注意 EventBridge 提供指数退避重试,失败后会发送到 DLQ。对于编排,选择 Step Functions:Standard 状态机适用于需要执行历史和内置重试/Catch 的长时间运行的持久工作流;Express 状态机适用于成本更低、尽力而为执行的高吞吐量短时工作流。使用服务集成(arn:aws:states:::lambda:invoke 或 arn:aws:states:::aws-sdk:apigateway:invoke)的任务集成,以及使用 “waitForTaskToken” 的回调模式来实现异步外部审批。通过指数退避实现重试和 Catch,并为长任务使用 HeartbeatSeconds。使用 Map 状态来并行处理大型集合,但要注意并发性和下游的节流限制。一个典型的陷阱是为需要精确一次持久化历史的工作流错误地选择了 Express——为保证可审计性应选择 Standard。此外,通过传递幂等性令牌并让目标服务在写入时强制唯一性,来确保由 Step Functions 调用的任务具有幂等性。

实践问题:用例场景

场景:StreamlyGames 在一个多账户 AWS 环境中运营一个全球游戏后端。玩家将 10 MB 的游戏片段上传到 S3;一个处理管道必须转码视频、运行机器学习分析,并将结果以有序、去重的方式写入 Aurora Serverless,同时消费者需要能够扩展。

挑战:确保每个上传的文件为每个玩家触发精确一次的有序处理,在不丢失事件的情况下处理流量峰值,并扩展消费者以进行机器学习推理,同时防止重复的数据库写入。

推荐方法:

  1. 创建一个 S3 事件通知,将对象创建事件发布到 EventBridge 自定义总线(PutEvents),同时也发布到一个 SQS FIFO 队列(CreateQueue 时设置 FifoQueue=true),以 PlayerID 作为 MessageGroupId,并使用基于 S3 ETag 的 MessageDeduplicationId。
  2. 配置一个 Lambda 消费者,使用 SQS 事件源映射(CreateEventSourceMapping),设置 BatchSize=1、FunctionResponseTypes=[“ReportBatchItemFailures”],并将 VisibilityTimeout 设置为大于最大处理时间;在调用第三方机器学习 API 时使用 ChangeMessageVisibility。
  3. Lambda 使用确定性的幂等键(INSERT … ON CONFLICT DO NOTHING 或唯一约束)向 Aurora 执行幂等的数据库写入,并为进度设置检查点;对于长时间的机器学习调用,使用带任务令牌的异步 Step Functions(arn:aws:states:::lambda:invoke.waitForTaskToken)或使用 Step Functions Express 来实现高吞吐量。
  4. 为了扩展推理能力,将中间事件按每个区域的每个分片写入 Kinesis Data Streams 以供高吞吐量消费者使用,并为专用的机器学习工作集群启用增强型扇出消费者(SubscribeToShard);监控 IteratorAgeMilliseconds 并使用 UpdateShardCount 进行扩展。

理由:使用 SQS FIFO 保证每个玩家的顺序并在接收时去重,使用幂等的数据库写入来强制实现精确一次的语义,并使用 Kinesis 加增强型扇出或 Step Functions 来处理突发的高吞吐量机器学习处理,同时保持消费者解耦和可扩展。


数据库与缓存 (RDS · 所有领域

练习这些题目 → · 在 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.

通过考试 →

浏览 Amazon →

Related guides

一体化访问

一次订阅。所有考试。

所有计划均可无限制搜索答案、进行模拟测试、获取AI解释以及访问完整的资源库 — 支持20多种语言。

每月
24.87
Just €0.83/day
包含所有内容:
  • 无限答案搜索
  • 无限模拟测试
  • AI驱动的解释
  • 完整资源库
  • 20多种语言
  • 每周内容更新
  • 奖励与推荐
  • 优先支持
开始免费试用

无需信用卡*

最具价值
12个月
179.87
Just €0.49/daySave 40%
包含所有内容:
  • 无限答案搜索
  • 无限模拟测试
  • AI驱动的解释
  • 完整资源库
  • 20多种语言
  • 每周内容更新
  • 奖励与推荐
  • 优先支持
开始免费试用

无需信用卡*

✓ 包含免费计划 · ✓ 随时取消 · ✓ 所有计划均解锁完整产品