View a markdown version of this page

AWS Glue 串流 - AWS Glue

本文為英文版的機器翻譯版本,如內容有任何歧義或不一致之處,概以英文版為準。

AWS Glue 串流

AWS Glue 串流是 的元件 AWS Glue,可讓您近乎即時地有效率地處理串流資料,讓您能夠執行關鍵任務,例如資料擷取、處理和機器學習。使用 Apache Spark 串流架構, AWS Glue 串流提供無伺服器服務,可大規模處理串流資料。 可在 Apache Spark 上 AWS Glue 提供各種最佳化,例如無伺服器基礎設施、自動擴展、視覺化任務開發、串流任務的即時筆記本,以及其他效能改善。

串流使用案例

AWS Glue 串流的一些常見使用案例包括:

Near-real-time的資料處理: AWS Glue 串流可讓組織近乎即時地處理串流資料,讓他們能夠衍生洞見,並根據最新資訊及時做出決策。

詐騙偵測:您可以使用 Streaming AWS Glue 進行串流資料的即時分析,這對於偵測信用卡詐騙、網路入侵或線上詐騙等詐騙活動很有價值。持續處理和分析傳入的資料,可讓您快速找出可疑的模式或異常情況。

社交媒體分析: AWS Glue 串流可以處理即時社交媒體資料,例如推文、文章或評論,讓組織能夠即時監控趨勢、情緒分析和管理品牌評價。

物聯網 (IoT) 分析: AWS Glue 串流適用於處理和分析 IoT 裝置、感應器和連線機器所產生的高速資料串流。可進行即時監控、異常偵測、預測性維護和其他 IoT 分析使用案例。

Clickstream 分析: AWS Glue 串流可以處理和分析來自網站或行動應用程式的即時 clickstream 資料。這可協助企業深入了解使用者行為、打造個人化使用者體驗,並根據即時點擊流資料將行銷活動最佳化。

日誌監控和分析: AWS Glue 串流可以持續即時處理和分析來自伺服器、應用程式或網路裝置的日誌資料。這有助於偵測異常、疑難排解問題,以及監控系統運作狀態和效能。

建議系統: AWS Glue 串流可以即時處理使用者活動資料,並動態更新建議模型。這可讓系統根據使用者的行為和偏好即時提供個人化的建議。

這些是可套用 AWS Glue 串流的各種使用案例範例。其與 AWS 生態系統和受管服務的整合,使其成為雲端中即時串流處理和分析的便利選擇。

使用 AWS Glue 串流有哪些好處?

使用 AWS Glue 串流的優點如下:

  • 無伺服器: AWS Glue 串流是無伺服器,無需管理基礎設施。此設計可減少營運成本,讓使用者專注於資料處理和分析工作,無須分神管理基礎設施。

  • Autoscaling: AWS Glue Streaming 提供自動擴展功能,可根據工作負載動態調整處理容量。此功能會自動擴展或縮減以應付資料量的波動,確保最佳效能和資源使用率。

  • 視覺化開發:串流任務開發可能很複雜。 AWS Glue 串流透過提供視覺化撰寫工具 AWS Glue Studio 來解決這項挑戰。 AWS Glue Studio 簡化了建立串流工作流程的程序,並可讓開發人員以視覺化方式設計和管理串流應用程式,進而降低學習曲線並提高生產力。

  • 符合成本效益:串流是無伺服器服務,無需佈建和維護基礎設施,即可 AWS Glue 提供成本效益。系統會根據執行串流任務期間所耗用的資源向使用者計費,以便根據實際使用情況進行成本最佳化和擴展。

  • 處理複雜的工作負載: AWS Glue 串流旨在處理複雜的串流工作負載。它可以處理和分析大量即時資料、支援進階轉換,並與其他 AWS 服務整合,從而實現複雜的串流資料管道和分析工作流程。

  • 無鎖定: AWS Glue 串流可提供彈性,並避免廠商鎖定。使用者可以利用 AWS Glue 串流作為更廣泛的 AWS 生態系統的一部分,將其與其他 AWS 服務無縫整合。這樣可以輕鬆地與現有的資料來源、應用程式和服務整合,而無須與特定技術或平台綁定在一起。

何時使用 AWS Glue 串流?

說到串流使用案例,您可以有很多選擇。我們建議在下列情況下 AWS Glue 進行串流。

  1. 如果您已經使用 AWS Glue 或 Spark 進行批次處理, AWS Glue 串流是您的理想選擇。它可讓您順利轉換至建置串流任務,而無須學習新的語言或框架。 AWS Glue 串流利用現有的知識和基礎設施,簡化了任務開發程序,並可讓您輕鬆地將資料處理功能擴展到即時串流案例。

  2. 如果您需要統一的服務或產品來處理批次、串流和事件驅動工作負載, AWS Glue 串流是您的解決方案。使用 AWS Glue 串流,您可以將資料處理需求合併為單一架構,消除管理多個系統的複雜性。這樣可讓您有效地開發和維護各種資料工作流程,同時確保不同工作負載類型的一致性和相容性。

  3. AWS Glue 串流非常適合涉及極大型串流資料磁碟區和複雜轉換的案例,例如串流或關聯式資料庫之間的聯結。它可以有效地處理和分析大量資料串流,讓您能夠輕鬆處理高需求的工作負載。無論是高速資料擷取還是複雜的資料處理, AWS Glue 串流的可擴展性和進階處理功能都能確保最佳效能和準確的結果。

  4. 如果您偏好視覺化方法來建置串流任務, AWS Glue 則提供 AWS Glue Studio,可讓您以視覺化方式設計和管理串流應用程式,簡化開發程序。此工具的直覺式介面可讓開發人員使用視覺化介面來建立、設定和監控串流工作流程,進而減少學習曲線並提高生產力。

  5. 如果AWS Glue 嚴格的 SLAs (服務水準協議) 大於 10 秒,則串流是near-real-time的使用案例的絕佳選擇

  6. 如果您使用 Apache Iceberg、Apache Hudi 或 Delta Lake 建置交易資料湖, AWS Glue 串流會提供這些開放資料表格式的原生支援。這種無縫整合可讓您直接從這些交易資料湖處理串流資料,以確保資料的一致性、完整性和相容性。

  7. 需要擷取各種資料目標的串流資料時: AWS Glue 串流會將原生目標提供給各種資料目標,例如 Amazon Redshift、Amazon RDS、Amazon Aurora、Oracle、SQL Server 和其他目標。

支援的資料來源

AWS Glue 串流支援下列資料來源:

  • Amazon Kinesis

  • Amazon MSK (Managed Streaming for Apache Kafka)

  • 自我管理的 Apache Kafka

支援的資料目標

AWS Glue 串流支援各種資料目標,例如:

  • Data Catalog 支援 AWS Glue 的資料目標

  • Amazon S3

  • Amazon Redshift

  • MySQL

  • PostgreSQL

  • Oracle

  • Microsoft SQL Server

  • Snowflake

  • 任何可使用 JDBC 連接的資料庫

  • Apache Iceberg、Delta 和 Apache Hudi

  • AWS Glue Marketplace 連接器

啟用串流任務的即時模式

即時模式 (RTM) 是 AWS Glue 6.0 版中提供 Spark 結構化串流的新執行模型。RTM end-to-end延遲從秒或分鐘減少到次秒。即時模式僅適用於 Spark 結構化串流任務。它不適用於舊版 Spark 串流 (DStreams) 或其他任務類型。

RTM 使用 Trigger.RealTime。任務會在批次時段 (預設 5 分鐘) 內持續執行,並在記錄到達時處理記錄,而不是跨間隔累積資料。這與預設微批次模型不同,其中 forEachBatch/Trigger.ProcessingTime 輪詢、處理、遞交和每個間隔重新啟動任務。

重要

RTM 需要透過任務引數明確選擇加入。如果沒有足夠的任務插槽來涵蓋所有來源分割區,RTM 會無提示地捨棄未指派的分割區。您必須佈建足夠的工作者來涵蓋所有 Kafka 分割區。

先決條件

啟用即時模式之前,請確認您的任務符合下列要求:

  • AWS Glue 6.0 版

  • 任務必須使用 Spark 結構化串流。即時模式不適用於舊版 Spark 串流 (DStreams) 或其他任務類型。

  • 任務類型必須是 Spark Streaming (gluestreaming 命令)

  • 任務語言必須是 Scala (--job-language scala)。在 Spark 4.2 之前,無法使用 PySpark RTM 支援。

  • 僅限 Kafka 來源。RTM in AWS Glue 6.0 不支援 Amazon Kinesis。

  • 僅限無狀態操作 (選取、篩選、專案、映射)。不支援有狀態的操作,例如彙總、聯結、重複資料刪除和視窗化操作。

  • 輸出模式必須為更新。RTM 不支援附加模式。

  • 自動擴展與即時模式不相容。請勿為 RTM 任務啟用自動擴展。設定足夠涵蓋來源主題中所有 Kafka 分割區的固定工作者數量。

何時使用即時模式

即時模式是專為特定類別的串流工作負載所設計。在下列情況下,請考慮使用即時模式:

  • 您需要一秒的end-to-end延遲,且微批次延遲 (1–2 秒或以上) 對您的使用案例而言太高。

  • 您的管道會執行無狀態轉換,例如篩選、投影、擴充或將記錄從 Kafka 路由到 Kafka 或其他接收器。

  • 您有固定且可預測的 Kafka 分割區數量,並且可以相應地佈建工作者。

  • 您的任務是以 Scala 撰寫。

在下列情況下繼續使用微批次模式:

  • 您需要有狀態的操作,例如彙總、聯結、重複資料刪除或視窗化運算。

  • 您可以使用 Amazon Kinesis 做為來源。

  • 您編寫 PySpark 任務。

  • 您依賴自動擴展來處理變數資料磁碟區。

  • 您可以使用 forEachBatch或 GlueContext 串流 API。

  • 第二級延遲適用於您的使用案例。

即時模式的運作方式

以下說明微型批次模型與即時模式之間的差異:

微型批次模式

每個間隔都會啟動任務、讀取累積的資料、處理資料、遞交檢查點、終止任務和重複。最小延遲約為 1–2 秒。

即時模式

任務會啟動一次,並在 期間執行 batchDurationMs(預設 5 分鐘)。任務會在記錄送達時處理記錄,延遲低於一秒。在截止日期時,任務會協同停止。驅動程式遞交檢查點,下一個批次會重新啟動任務。

兩種模式都使用相同的檢查點格式和復原機制。關鍵差異是任務生命週期。Micro-batch 模式會每隔間隔終止和重新啟動任務。即時模式可讓任務在較長的批次時段內持續執行。

重要

如果沒有足夠的任務槽來處理所有來源分割區,RTM 會無提示地捨棄未指派的分割區。請確定您佈建足夠的工作者來涵蓋所有分割區。

啟用即時模式

您可以將--enable-real-time-mode任務引數設定為 來啟用即時模式true。您可以在 AWS Glue 主控台或透過 API 設定此引數。

啟用即時模式 (主控台)

  1. 開啟 AWS Glue 主控台並開啟您的串流任務。

  2. 選擇 Job details (任務詳細資訊) 索引標籤。

  3. 針對 Glue 版本,選擇 Glue 6.0。針對類型,選擇 Spark 串流

  4. 捲動至任務參數區段。

  5. 選擇新增參數

  6. Key (索引鍵) 欄位,輸入 --enable-real-time-mode。針對數值,輸入 true

  7. 選擇儲存

注意

前置破折號為必要項目。任務參數是 的主控台檢視DefaultArguments

啟用即時模式 (API)

--enable-real-time-mode 旗標會存放在任務定義的DefaultArguments地圖中。您可以在建立或更新任務時加以設定。

建立新的任務 (AWS CLI)

執行以下命令:

aws glue create-job \ --name my-rtm-job \ --role arn:aws:iam::123456789012:role/MyGlueRole \ --glue-version 6.0 \ --worker-type G.1X --number-of-workers 4 \ --command '{"Name":"gluestreaming","ScriptLocation":"s3://my-bucket/scripts/rtm-job.scala"}' \ --default-arguments '{ "--enable-real-time-mode": "true", "--job-language": "scala", "--class": "GlueApp", "--TempDir": "s3://my-bucket/tmp/" }' \ --region us-east-2
建立新任務 (boto3)

使用下列程式碼:

import boto3 glue = boto3.client("glue", region_name="us-east-2") glue.create_job( Name="my-rtm-job", Role="arn:aws:iam::123456789012:role/MyGlueRole", GlueVersion="6.0", WorkerType="G.1X", NumberOfWorkers=4, Command={ "Name": "gluestreaming", "ScriptLocation": "s3://my-bucket/scripts/rtm-job.scala", }, DefaultArguments={ "--enable-real-time-mode": "true", "--job-language": "scala", "--class": "GlueApp", "--TempDir": "s3://my-bucket/tmp/", }, )
更新現有任務 (AWS CLI)

執行以下命令:

aws glue update-job \ --job-name my-existing-job \ --job-update '{ "GlueVersion": "6.0", "DefaultArguments": { "--enable-real-time-mode": "true", "--job-language": "scala" } }'

撰寫串流指令碼

任務引數宣告使用即時模式的意圖。您的指令碼會選取觸發條件。

下列 Scala 範例顯示使用 的串流查詢Trigger.RealTime

import org.apache.spark.sql.streaming.Trigger val query = df.writeStream .format("kafka") .outputMode("update") .trigger(Trigger.RealTime(60000L)) // checkpoint interval in milliseconds .start() query.awaitTermination()

Trigger.RealTime 需要以毫秒為單位的檢查點間隔。需要更新輸出模式。附加模式擲回 OUTPUT_MODE_NOT_SUPPORTED

只要標記已設定,您就可以在一個指令碼中混合模式:

dfA.writeStream.outputMode("update").trigger(Trigger.RealTime(60000L)).start() dfB.writeStream.outputMode("append").trigger(Trigger.ProcessingTime("30 seconds")).start()

缺少旗標時的行為

以下說明未設定--enable-real-time-mode旗標時任務的行為:

  • 在沒有--enable-real-time-mode旗標的情況下啟動即時查詢的任務,會在查詢開始時失敗。失敗訊息會指示您新增 引數。

  • 沒有此旗標,就不會影響Micro-batch-only的任務。

  • 設定旗標但僅使用微批次查詢的任務也不會受到影響。

考量和限制

當您使用即時模式時,請考慮下列事項:

分割區捨棄

如果沒有足夠的任務槽涵蓋所有來源分割區,則不會處理未指派的分割區。佈建工作者以涵蓋所有 Kafka 分割區。

無自動擴展

請勿啟用即時模式任務的自動擴展。自動擴展與 RTM 不相容,並引入可抵消低延遲優點的延遲。在來源主題中佈建等於或大於 Kafka 分割區數量的固定工作者數量。

僅限 Kafka

Amazon Kinesis 來源不支援 RTM in AWS Glue 6.0。

僅限 Scala

在 Spark 4.2 之前,RTM 不支援 PySpark。

僅限無狀態

transformWithState 不支援彙總、聯結、重複資料刪除、視窗化操作和 。

forEachBatch 不相容

RTM 不會使用forEachBatch模型。Trigger.RealTime 直接writeStream搭配 使用 。

檢查點復原

任務重新啟動時,RTM 會從最後一個檢查點復原。每個 都會發生檢查點batchDurationMs。最差情況重新處理是一個批次時段的持續時間 (at-least-once語意)。