本文為英文版的機器翻譯版本,如內容有任何歧義或不一致之處,概以英文版為準。
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 進行串流。
如果您已經使用 AWS Glue 或 Spark 進行批次處理, AWS Glue 串流是您的理想選擇。它可讓您順利轉換至建置串流任務,而無須學習新的語言或框架。 AWS Glue 串流利用現有的知識和基礎設施,簡化了任務開發程序,並可讓您輕鬆地將資料處理功能擴展到即時串流案例。
如果您需要統一的服務或產品來處理批次、串流和事件驅動工作負載, AWS Glue 串流是您的解決方案。使用 AWS Glue 串流,您可以將資料處理需求合併為單一架構,消除管理多個系統的複雜性。這樣可讓您有效地開發和維護各種資料工作流程,同時確保不同工作負載類型的一致性和相容性。
AWS Glue 串流非常適合涉及極大型串流資料磁碟區和複雜轉換的案例,例如串流或關聯式資料庫之間的聯結。它可以有效地處理和分析大量資料串流,讓您能夠輕鬆處理高需求的工作負載。無論是高速資料擷取還是複雜的資料處理, AWS Glue 串流的可擴展性和進階處理功能都能確保最佳效能和準確的結果。
如果您偏好視覺化方法來建置串流任務, AWS Glue 則提供 AWS Glue Studio,可讓您以視覺化方式設計和管理串流應用程式,簡化開發程序。此工具的直覺式介面可讓開發人員使用視覺化介面來建立、設定和監控串流工作流程,進而減少學習曲線並提高生產力。
如果AWS Glue 嚴格的 SLAs (服務水準協議) 大於 10 秒,則串流是near-real-time的使用案例的絕佳選擇。
如果您使用 Apache Iceberg、Apache Hudi 或 Delta Lake 建置交易資料湖, AWS Glue 串流會提供這些開放資料表格式的原生支援。這種無縫整合可讓您直接從這些交易資料湖處理串流資料,以確保資料的一致性、完整性和相容性。
需要擷取各種資料目標的串流資料時: 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 設定此引數。
啟用即時模式 (主控台)
-
開啟 AWS Glue 主控台
並開啟您的串流任務。 -
選擇 Job details (任務詳細資訊) 索引標籤。
-
針對 Glue 版本,選擇 Glue 6.0。針對類型,選擇 Spark 串流。
-
捲動至任務參數區段。
-
選擇新增參數。
-
在 Key (索引鍵) 欄位,輸入
--enable-real-time-mode。針對數值,輸入true。 -
選擇儲存。
注意
前置破折號為必要項目。任務參數是 的主控台檢視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語意)。