串流資料
在 AWS 上開發應用程式
Ricardo Sueiras
Principal Technologist
串流資料基礎
- 記錄隨到隨處理,持續進行。
- 生產者寫入到串流。
- 消費者可獨立且平行讀取。
- 吸收流量高峰,無須等待批次作業。
串流資料由哪些部分組成?
- 生產者:送出記錄的應用程式。
- 串流:可持久、具順序的緩衝區,解耦生產者與消費者。
- 消費者:讀取並處理記錄的應用程式。
- 儲存與保留:定義記錄可供重播的時間長度。
分片(Shards)
- 分片是吞吐與平行處理的單位。
- 區分鍵將記錄分散到各分片。
- 僅保證分片內的順序。
- 消費者用檢查點追蹤進度並從失敗中復原。
AWS 串流資料受管服務
- Amazon Kinesis Data Streams。
- Amazon Data Firehose。
- 串流資料應用的核心組件。
Kinesis Data Streams
- 低延遲串流。
- 可設定保留期:24 小時(預設)至 365 天。
- 可持久、可重播的串流。
- 多個消費者可平行讀取。
- 開發者用分片管理擴充。
何時使用 Kinesis Data Streams
- 需要即時低延遲時選用。
- 需要重播記錄。
- 有多個獨立消費者。
- 需要自訂轉換邏輯。
寫入記錄
- 區分鍵決定分片放置與順序。
PutRecords 每次請求可批次最多 500 筆或 5 MB。
- 失敗會以每筆為單位回報。
FailedRecordCount 協助找出失敗的記錄。
消費記錄:經典模式
- 消費者按需輪詢 Kinesis API。
- 寫入上限:1 MB/秒 或 1,000 筆/秒。
- 讀取上限:2 MB/秒 或每秒 5 次
GetRecords 呼叫。
消費記錄:串流起點
- 消費者指定從串流的哪裡開始讀。
TRIM_HORIZON:分片中最舊可用記錄。
LATEST:下一筆新記錄。
AT_SEQUENCE_NUMBER 或 AFTER_SEQUENCE_NUMBER:從特定序號開始。
AT_TIMESTAMP:在指定時間戳之時或之後。
消費記錄:增強型扇出(enhanced fan-out)
- 每個註冊的消費者在每個分片有專屬 2 MB/秒。
- 記錄以 HTTP/2 推送傳遞。
- 適用於多個獨立消費者需要低延遲時。
- 且同一串流需要高吞吐。
Lambda 事件來源對應
- Lambda 輪詢 Kinesis,並以批次觸發。
BatchSize 與 MaximumBatchingWindow 控制批次形狀。
ParallelizationFactor:每個分片的並行觸發數。
BisectBatchOnFunctionError:分割失敗批次以隔離問題記錄。
ReportBatchItemFailures:避免重試已成功的記錄。
OnFailure:將無法處理的批次送至 SQS 或 SNS 的 DLQ。
Data Firehose
- 受管服務,將串流資料送至目的地。
- 目標包含 Amazon S3、Redshift 與 HTTP 端點。
- 支援緩衝、格式轉換與 Lambda 轉換。
- 延遲較高,緩衝數秒到數分鐘。
- 不支援重播或長期保留。
- 適合簡單將資料送達目的地。
樣式:熱-冷
- 熱-冷樣式結合 Data Streams 與 Data Firehose。
- Data Streams 讓生產者可即時處理。
- Firehose 將同來源的記錄歸檔到 S3。
Kinesis 與 SQS
- 需要有序、可重播且多消費者的串流用 Kinesis。
- Kinesis 適合時間序列與有序資料。
- 任務分派則用 SQS。
擴充串流資料
- 資料串流會被切成分片。
- 區分鍵將記錄導向特定分片。
- 相同鍵永遠落在同一分片,以保順序。
- 區分鍵不均會造成熱分片。
- 可用高基數鍵緩解。
- 透過重分片向外或向內擴充。
容量模式:預置(provisioned)
- 提供兩種容量模式。
- 預置模式以固定分片數,手動管理。
- 成本與容量可預期。
- 需監控以避免節流與過度預置。
容量模式:隨需(on-demand)
- 隨需模式依流量自動擴縮。
- 無須手動管理分片。
- 適合負載難以預測的營運。
- 長期高流量下成本較高。
Kinesis Producer Library(KPL)
- AWS 函式庫簡化可擴充的生產者與消費者開發。
- Kinesis Producer Library(KPL)會批次與彙整記錄。
- 也會壓縮並在失敗時重試。
- 非同步緩衝以最大化生產者吞吐。
Kinesis Consumer Library(KCL)
- Kinesis Consumer Library(KCL)負責分片協調。
- 管理工作節點租用與檢查點。
- 狀態儲存在 DynamoDB 資料表。
- 消費者用該狀態查找續讀位置。
失敗處理:冪等性
- Kinesis 採至少一次交付,可能出現重複。
- 將消費者設計為冪等。
- 以序號或 DynamoDB 中的業務鍵去重。
- 消費者可自序號或時間戳重讀保留記錄。
失敗處理:吞吐錯誤
ProvisionedThroughputExceededException 表示被節流。
- 透過增加分片數解決。
- 改善區分鍵分佈。
- 以退避重試。
- 切換容量模式。
- 在 CloudWatch 追蹤
IncomingBytes 與 IncomingRecords。
- 監看
WriteProvisionedThroughputExceeded 以偵測節流。
- 當消費速度慢於生產時會發生延遲(lag)。
- 追蹤
GetRecords.IteratorAgeMilliseconds:數值高表示延遲。
安全性
- 伺服端靜態加密,使用 AWS KMS。
- 傳輸中加密,使用 TLS。
- 以 IAM 原則控管存取。
Preparing Video For Download...