Amazon DVA-C02: 訊息傳遞、串流與事件驅動架構 (SNS, SQS, Kinesis, EventBridge, Step Functions) — 學習指南
屬於 AWS Developer Associate DVA-C02 — 學習指南. 使用經過驗證的解答練習: Amazon 考試中心, 或參加限時模擬考試: ExamRoll.io.
選擇正確的訊息傳遞與串流基礎元件
在 SNS、SQS (標準 vs FIFO)、Kinesis、EventBridge 和 Step Functions 之間做選擇,首先要從通訊模式著手:發布/訂閱 (pub/sub)、點對點、有序串流、事件匯流排路由或工作流程協調。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 (Standard) 或 StartSyncExecution (用於同步的 Express 模式),並透過 Task 與
undefined
等服務整合。當吞吐量、順序性、交付保證、保留時間和協調需求發生衝突時,請考量這些權衡取捨。
- SNS: 高吞吐量扇出,無順序性,Publish/Subscribe 模式,使用 MessageAttributes 進行路由。
- SQS Standard: 至少一次 (at-least-once) 交付,盡力而為 (best-effort) 的順序性,支援長輪詢,是實現解耦的較低成本選擇。
- SQS FIFO: 在 MessageGroupId 內保證僅一次 (exactly-once) 且有序,用於嚴格要求順序和去重複的場景。
- Kinesis Data Streams: 每個 shard 內有序,高吞吐量串流,使用 PutRecord/PutRecords,需要進行 shard 擴展。
- Kinesis Firehose: 託管式交付與緩衝,支援伺服器端加密和 Lambda 轉換。
- EventBridge: 具備路由規則的事件匯流排,支援 schema registry、封存與重播 (archive & replay),使用 PutEvents API。
- Step Functions: 有狀態的協調,支援重試/Catch,Standard 與 Express 模式在持久性和吞吐量之間有不同的權衡取捨。
SQS 與 SNS 模式、去重複化及消費者擴展
當您需要持久性的解耦時,SQS 是首選;生產者應實作 SendMessage/SendMessageBatch,消費者則使用帶有 WaitTimeSeconds 的 ReceiveMessage 來啟用長輪詢,以減少空接收的次數。若要嚴格的順序性與去重複,請使用 CreateQueue (FifoQueue=true) 建立一個 FIFO 佇列,並設定 MessageGroupId 來劃分有序分區;使用 MessageDeduplicationId 或啟用 ContentBasedDeduplication,以便在去重複時間窗內抑制相同的 payload。標準佇列可能會傳遞重複的訊息——因此,消費者必須設計成冪等的 (idempotent),例如透過資料庫的條件式寫入 (使用 ConditionExpression attribute_not_exists(pk) 的 DynamoDB PutItem) 或在交易中使用 RDS 的唯一約束 (unique constraints) 和 upsert 操作。設定重新驅動策略 (redrive policies),在達到 maxReceiveCount 後將失敗的訊息路由到 DLQ;透過 GetQueueAttributes 監控 ApproximateNumberOfMessages 和 ApproximateNumberOfMessagesNotVisible。與 Lambda 整合時,使用 CreateEventSourceMapping 來對應 SQS:設定 BatchSize、MaximumBatchingWindowInSeconds,並啟用 FunctionResponseTypes = [“ReportBatchItemFailures”] 以使用部分批次回應 (partial-batch-response) 的語意,避免重新處理已成功的紀錄。注意 FIFO 與 Lambda 整合的語意:訊息群組的順序性會強制每個 MessageGroupId 進行單線程處理,這限制了每個群組的並行性;可以透過將訊息分割到多個 group ID,或使用 SNS 將訊息分發到多個佇列以供平行消費者處理,來進行擴展。同時也要注意可見性逾時 (visibility timeout):當處理時間過長時,應設定 ChangeMessageVisibility,否則您將面臨重複處理的風險。
Kinesis Data Streams 與 Firehose:順序性、保留與背壓處理
Kinesis Data Streams 為串流應用場景提供每個 shard 內的順序性以及持久的資料保留。生產者使用一個會對應到 shard 的 PartitionKey 來呼叫 PutRecord 或 PutRecords (批次);消費者則呼叫 GetShardIterator 和 GetRecords,然後使用 KCL (Kinesis Client Library) 或自訂的 DynamoDB checkpoint 表來記錄檢查點偏移量 (checkpoint offsets)。預設的保留時間為 24 小時 (可根據串流設定調整為更長的時間,並在可用時使用擴展保留功能);應使用 UpdateShardCount 來規劃 shard 數量,以匹配寫入吞吐量和讀取者的平行處理能力。消費者的擴展是受限的:單一的 Lambda 事件來源映射 (event source mapping) 會將一個 shard 對應到一個 Lambda 並行執行個體,因此需要增加 shard 數量來提高消費者的並行性,或者啟用增強型扇出 (enhanced fan-out),讓每個消費者都有自己 2 MB/sec 的通道,並使用 SubscribeToShard API (消費者註冊) 進行獨立擴展。使用 PutRecords 進行高效率的批次處理;當消費者處理落後時會產生背壓 (back-pressure) (可監控 GetRecords.IteratorAgeMilliseconds)。為了應對流量高峰,可以在 Kinesis 上進行緩衝或在前端使用 SQS,生產者應使用帶有指數退避 (exponential backoff) 的重試機制,並對包含 PII 的資料流使用 KMS 進行串流層級的加密。Kinesis Data Firehose 簡化了交付過程:設定 BufferingHints (SizeInMBs、IntervalInSeconds)、CompressionFormat 以及一個用於資料轉換的 Lambda。Firehose 會處理對目標端點的重試/退避,並能將失敗的紀錄寫入一個備份用的 S3 儲存貯體。一個常見的陷阱是 shard 佈建不足:這會導致消費者飢餓 (starve) 和延遲飆升;應主動測量並進行擴展。
使用 EventBridge 與 Step Functions 進行路由與協調
EventBridge 擅長於 schema 驅動的事件路由、跨帳戶/事件合作夥伴整合,它使用 PutEvents 來注入事件,並用 PutRule/PutTargets 將事件路由到 SQS、Lambda、Kinesis、Step Functions 或 HTTP 端點。EventBridge 使用事件模式 (event patterns) 進行篩選,並支援封存與重播 (replay) 以重建狀態。請為規則使用死信佇列 (dead-letter queues) (透過 Target 的 SqsParameters 或 DeadLetterConfig),並注意 EventBridge 在失敗時會提供具備指數退避 (exponential backoff) 的重試機制,之後才會送到 DLQ。若需協調,請選擇 Step Functions:Standard 狀態機適用於需要執行歷史記錄、內建重試/Catch 機制的長時間執行、持久性工作流程;而 Express 則適用於成本較低、盡力而為 (best-effort) 執行的高吞吐量、短時間工作流程。使用 Task 與服務整合 (例如 arn:aws:states:::lambda:invoke 或 arn:aws:states:::aws-sdk:apigateway:invoke) 以及使用 “waitForTaskToken” 的回呼模式 (callback patterns) 來實作非同步的外部核准。透過指數退避來實作重試與 Catch,並對長時間執行的任務使用 HeartbeatSeconds。使用 Map 狀態來平行處理大型集合,但要注意並行性 (concurrency) 和下游的節流 (throttles)。一個常見的陷阱是,對於需要「僅一次」(exactly-once) 持久性歷史記錄的工作流程,誤選了 Express——為了可稽核性 (auditability),應選擇 Standard。此外,透過傳遞冪等性權杖 (idempotency token) 並讓目標服務在寫入時強制執行唯一性,來確保由 Step Functions 叫用的任務具有冪等性 (idempotency)。
實務問題:使用案例情境
情境:StreamlyGames 在一個多帳戶的 AWS 環境中營運一個全球遊戲後端。玩家會上傳 10 MB 的遊戲片段到 S3;一個處理管道必須對影片進行轉碼、執行機器學習 (ML) 分析,並將結果以有序、去重複的方式寫入 Aurora Serverless,同時消費者 (consumers) 需具備可擴展性。
挑戰:確保每個上傳的檔案都能為每位玩家觸發「僅一次」且有序的處理,處理流量高峰時不能遺失事件,並擴展用於 ML 推論的消費者,同時防止重複的資料庫寫入。
建議方法:
- 建立一個 S3 事件通知,將物件建立事件發佈到一個 EventBridge custom bus (使用 PutEvents),同時也發佈到一個 SQS FIFO 佇列 (使用 CreateQueue 搭配 FifoQueue=true)。該佇列以 PlayerID 作為 MessageGroupId,並使用基於 S3 ETag 的 MessageDeduplicationId。
- 設定一個 Lambda 消費者,搭配一個 SQS 事件來源映射 (CreateEventSourceMapping),使用 BatchSize=1、FunctionResponseTypes=[“ReportBatchItemFailures”],並將 VisibilityTimeout 設定為大於最長處理時間;在呼叫第三方 ML API 時,使用 ChangeMessageVisibility。
- Lambda 使用一個確定性的冪等性金鑰 (deterministic idempotency key) (例如 INSERT … ON CONFLICT DO NOTHING 或唯一性約束) 對 Aurora 執行冪等性的資料庫寫入,並對進度進行檢查點 (checkpoints);對於長時間的 ML 呼叫,使用帶有任務權杖 (task tokens) 的非同步 Step Functions (arn:aws:states:::lambda:invoke.waitForTaskToken),或使用 Step Functions Express 來實現高吞吐量。
- 為了擴展推論能力,將中繼事件寫入 Kinesis Data Streams (每個區域每個分區),以供高吞吐量的消費者使用,並為專用的 ML 工作者叢集 (worker fleets) 啟用增強型扇出消費者 (enhanced fan-out consumers) (SubscribeToShard);監控 IteratorAgeMilliseconds 並使用 UpdateShardCount 來進行擴展。
理由:使用 SQS FIFO 來保證每個玩家的順序性及在接收時的去重複化,使用冪等性的資料庫寫入來強制執行「僅一次」語意,並使用 Kinesis 加上增強型扇出或 Step Functions 來處理突發性、高吞吐量的 ML 處理,同時保持消費者的解耦與可擴展性。
← 資料庫與快取 (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.
通過考試 →