串流資料

在 AWS 上開發應用程式

Ricardo Sueiras

Principal Technologist

串流資料基礎

 

串流資料

  • 記錄隨到隨處理,持續進行。
  • 生產者寫入到串流。
  • 消費者可獨立且平行讀取。
  • 吸收流量高峰,無須等待批次作業。
在 AWS 上開發應用程式

串流資料由哪些部分組成?

 

串流系統元件

  • 生產者:送出記錄的應用程式。
  • 串流:可持久、具順序的緩衝區,解耦生產者與消費者。
  • 消費者:讀取並處理記錄的應用程式。
  • 儲存與保留:定義記錄可供重播的時間長度。
在 AWS 上開發應用程式

分片(Shards)

 

分片

 

  • 分片是吞吐與平行處理的單位。
  • 區分鍵將記錄分散到各分片。
  • 僅保證分片內的順序。
  • 消費者用檢查點追蹤進度並從失敗中復原。
在 AWS 上開發應用程式

AWS 串流資料受管服務

  • Amazon Kinesis Data Streams。
  • Amazon Data Firehose。
  • 串流資料應用的核心組件。

 

aws 服務

在 AWS 上開發應用程式

Kinesis Data Streams

  • 低延遲串流。
  • 可設定保留期:24 小時(預設)至 365 天。
  • 可持久、可重播的串流。
  • 多個消費者可平行讀取。
  • 開發者用分片管理擴充。

 

kinesis data streams

在 AWS 上開發應用程式

何時使用 Kinesis Data Streams

  • 需要即時低延遲時選用。
  • 需要重播記錄。
  • 有多個獨立消費者。
  • 需要自訂轉換邏輯。

 

kinesis data streams

在 AWS 上開發應用程式

寫入記錄

 

用戶端函式庫

  • 區分鍵決定分片放置與順序。
  • PutRecords 每次請求可批次最多 500 筆或 5 MB。
  • 失敗會以每筆為單位回報。
  • FailedRecordCount 協助找出失敗的記錄。
在 AWS 上開發應用程式

消費記錄:經典模式

 

經典

 

  • 消費者按需輪詢 Kinesis API。
  • 寫入上限:1 MB/秒 或 1,000 筆/秒。
  • 讀取上限:2 MB/秒 或每秒 5 次 GetRecords 呼叫。
在 AWS 上開發應用程式

消費記錄:串流起點

 

串流起點

  • 消費者指定從串流的哪裡開始讀。
  • TRIM_HORIZON:分片中最舊可用記錄。
  • LATEST:下一筆新記錄。
  • AT_SEQUENCE_NUMBERAFTER_SEQUENCE_NUMBER:從特定序號開始。
  • AT_TIMESTAMP:在指定時間戳之時或之後。
在 AWS 上開發應用程式

消費記錄:增強型扇出(enhanced fan-out)

  • 每個註冊的消費者在每個分片有專屬 2 MB/秒。
  • 記錄以 HTTP/2 推送傳遞。
  • 適用於多個獨立消費者需要低延遲時。
  • 且同一串流需要高吞吐。

 

樣式 - 增強型扇出

在 AWS 上開發應用程式

Lambda 事件來源對應

  • Lambda 輪詢 Kinesis,並以批次觸發。
  • BatchSizeMaximumBatchingWindow 控制批次形狀。
  • ParallelizationFactor:每個分片的並行觸發數。
  • BisectBatchOnFunctionError:分割失敗批次以隔離問題記錄。
  • ReportBatchItemFailures:避免重試已成功的記錄。
  • OnFailure:將無法處理的批次送至 SQS 或 SNS 的 DLQ。

 

lambda 事件來源對應

在 AWS 上開發應用程式

Data Firehose

  • 受管服務,將串流資料送至目的地。
  • 目標包含 Amazon S3、Redshift 與 HTTP 端點。
  • 支援緩衝、格式轉換與 Lambda 轉換。
  • 延遲較高,緩衝數秒到數分鐘。
  • 不支援重播或長期保留。
  • 適合簡單將資料送達目的地。

 

data firehose

在 AWS 上開發應用程式

樣式:熱-冷

 

熱/冷樣式

 

  • 熱-冷樣式結合 Data Streams 與 Data Firehose。
  • Data Streams 讓生產者可即時處理。
  • Firehose 將同來源的記錄歸檔到 S3。
在 AWS 上開發應用程式

Kinesis 與 SQS

 

sqs 與 kinesis

 

  • 需要有序、可重播且多消費者的串流用 Kinesis。
  • Kinesis 適合時間序列與有序資料。
  • 任務分派則用 SQS。
在 AWS 上開發應用程式

擴充串流資料

 

擴充

  • 資料串流會被切成分片。
  • 區分鍵將記錄導向特定分片。
  • 相同鍵永遠落在同一分片,以保順序。
  • 區分鍵不均會造成熱分片。
  • 可用高基數鍵緩解。
  • 透過重分片向外或向內擴充。
在 AWS 上開發應用程式

容量模式:預置(provisioned)

  • 提供兩種容量模式。
  • 預置模式以固定分片數,手動管理。
  • 成本與容量可預期。
  • 需監控以避免節流與過度預置。

 

容量模式

在 AWS 上開發應用程式

容量模式:隨需(on-demand)

  • 隨需模式依流量自動擴縮。
  • 無須手動管理分片。
  • 適合負載難以預測的營運。
  • 長期高流量下成本較高。

 

容量模式

在 AWS 上開發應用程式

Kinesis Producer Library(KPL)

  • AWS 函式庫簡化可擴充的生產者與消費者開發。
  • Kinesis Producer Library(KPL)會批次與彙整記錄。
  • 也會壓縮並在失敗時重試。
  • 非同步緩衝以最大化生產者吞吐。

 

用戶端函式庫

在 AWS 上開發應用程式

Kinesis Consumer Library(KCL)

 

用戶端函式庫

 

  • Kinesis Consumer Library(KCL)負責分片協調。
  • 管理工作節點租用與檢查點。
  • 狀態儲存在 DynamoDB 資料表。
  • 消費者用該狀態查找續讀位置。
在 AWS 上開發應用程式

失敗處理:冪等性

 

處理失敗

 

  • Kinesis 採至少一次交付,可能出現重複。
  • 將消費者設計為冪等。
  • 以序號或 DynamoDB 中的業務鍵去重。
  • 消費者可自序號或時間戳重讀保留記錄。
在 AWS 上開發應用程式

失敗處理:吞吐錯誤

 

處理失敗

 

  • ProvisionedThroughputExceededException 表示被節流。
  • 透過增加分片數解決。
  • 改善區分鍵分佈。
  • 以退避重試。
  • 切換容量模式。
在 AWS 上開發應用程式
  • 在 CloudWatch 追蹤 IncomingBytesIncomingRecords
  • 監看 WriteProvisionedThroughputExceeded 以偵測節流。
  • 當消費速度慢於生產時會發生延遲(lag)。
  • 追蹤 GetRecords.IteratorAgeMilliseconds:數值高表示延遲。

 

記錄與監控

在 AWS 上開發應用程式

安全性

  • 伺服端靜態加密,使用 AWS KMS。
  • 傳輸中加密,使用 TLS。
  • 以 IAM 原則控管存取。

 

安全性

在 AWS 上開發應用程式

一起來練習吧!

在 AWS 上開發應用程式

Preparing Video For Download...