Spark 声明式管道
Spark 声明式管道(SDP)是一个声明式框架,用于在 AWS Glue 6.0 中构建批处理和流式传输数据管道。使用 SDP,您可以通过 SQL 或 Python 定义数据的目标形态,而框架将自动确定执行计划、解析数据集之间的依赖关系,并行运行相互独立的分支。
SDP 省去了用于读取、写入、目录注册和执行排序的命令式样板代码,从而简化了管道开发过程。您可以专注于业务转换,而管道基础设施则由框架负责处理。
SDP 在 AWS Glue 6.0 及更高版本中可用。
SDP 概念
管道由一个 YAML 清单文件(spark-pipeline.yml)和一个或多个 SQL 或 Python 转换文件组成。SDP 自动执行以下操作:
-
通过从表引用推断 DAG 来解析依赖关系
-
确定执行顺序,而无需进行手动编排
-
并行运行独立分支,以实现最大吞吐量
-
通过检查点管理流式传输表的增量状态
-
实体化时将输出表注册到目录中
数据集类型
SDP 中提供三种数据集类型:
- 流式传输表
-
仅处理自上次运行以来的新数据。使用检查点跨作业执行维护状态。将流式传输表用于数据摄取、事件流、IoT 数据、变更数据捕获和仅追加数据来源。
- 实体化视图
-
每次运行时都会完全重新计算数据集。输出始终反映源数据的当前状态。使用实体化视图进行聚合、联接、摘要分析和报告操作。
- 临时视图
-
会话范围限定,且不进行持久化存储或编目登记。将临时视图用于中间转换和暂存逻辑。
重要
在当前版本中,实体化视图始终执行完全重新计算。它们不支持增量刷新。使用流式传输表处理增量工作负载。
先决条件
要使用 SDP,您需要以下内容:
-
AWS Glue 版本 6.0
-
用于存储管道的 Amazon S3 位置(检查点、元数据)
-
对于数据目录集成(可选):将
--enable-glue-datacatalog设置为true。或者,您也可以通过 Spark 配置直接配置目录设置。 -
对于永久表存储:将
spark.sql.warehouse.dir设置为 Amazon S3 路径,或在管道 YAML 中设置database:字段,并确保 AWS Glue 数据库已将LocationUri配置为 Amazon S3 路径 -
要使用流式传输表进行跨运行增量处理,流式传输表的数据和检查点状态必须在 Amazon S3 上持久保存。例如,Hive 或 AWS Glue 托管(非 Iceberg)表要求将数据库
LocationUri设置为 Amazon S3 路径,而 Apache Iceberg 表则自行管理其表元数据。
重要
如果您在管道 YAML 中为 Hive 或 AWS Glue 托管(非 Iceberg)表使用 database: 字段,则相应的 AWS Glue 数据库必须将其 LocationUri 设置为 Amazon S3 路径。LocationUri 将托管流式传输表(及其 _spark_metadata 日志)放置在 Amazon S3 上,从而为这些表提供持久的跨运行增量处理能力。使用 Data Catalog 且采用 catalog-impl=GlueCatalog(选项 1)的 Iceberg 表不需要数据库 LocationUri。使用明确的 Amazon S3 位置创建或更新数据库:
aws glue create-database --database-input '{ "Name":"my_pipeline_db", "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db" }'
创建管道
要创建 SDP 管道,请完成以下步骤。
步骤 1:创建管道 YAML
创建一个名为 spark-pipeline.yml 的文件:
name: my_analytics_pipeline catalog: spark_catalog database: analytics_db storage: s3://my-bucket/pipeline-storage/ libraries: - glob: include: transformations/** configuration: spark.sql.shuffle.partitions: "4"
下表描述了管道 YAML 字段。
| 字段 | 必填 | 描述 |
|---|---|---|
name |
是 | 管道的名称。 |
catalog |
否 | 要使用的目录。默认值为 spark_catalog。 |
database |
否 | 用于存放输出表的 AWS Glue 数据库。有关 LocationUri 要求,请参阅先决条件。 |
storage |
是 | 用于存储管道检查点和元数据的 Amazon S3 路径。 |
libraries |
是 | 要包含的转换文件的 Glob 模式。 |
configuration |
否 | Spark 配置属性。 |
步骤 2:编写转换
在 transformations/ 目录中创建转换文件。您可以在同一个管道中使用 SQL、Python 或此两者。
SQL 示例(transformations/silver.sql):
CREATE MATERIALIZED VIEW silver_sales AS SELECT *, UPPER(region) as clean_region FROM bronze_sales WHERE amount > 0; CREATE MATERIALIZED VIEW gold_summary AS SELECT clean_region, COUNT(*) as order_count, SUM(amount) as total_revenue FROM silver_sales GROUP BY clean_region;
Python 示例(transformations/bronze.py):
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() @dp.materialized_view(comment="Raw sales data from S3") def bronze_sales() -> DataFrame: return spark.read.format("csv").option("header", "true") \ .option("inferSchema", "true") \ .load("s3://source-bucket/raw-data/sales/")
Python 流式传输表示例(transformations/events.py):
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() dp.create_streaming_table( "streaming_events", comment="Incremental event ingestion", schema="event_id STRING, event_type STRING, timestamp LONG, payload STRING" ) @dp.append_flow(target="streaming_events") def ingest_events() -> DataFrame: return ( spark.readStream.format("json") .schema("event_id STRING, event_type STRING, timestamp LONG, payload STRING") .load("s3://source-bucket/events/") )
注意
对于 Python 中的流式传输表,请将 dp.create_streaming_table() 与 @dp.append_flow(target=...) 结合使用。
步骤 3:上传到 Amazon S3
将您的管道文件上传到 Amazon S3 作为以下任一内容:
-
一个包含
spark-pipeline.yml和transformations/目录的.zip文件 -
包含相同结构的 Amazon S3 前缀(目录)
步骤 4:创建并运行 AWS Glue 作业
使用以下参数创建 AWS Glue 作业:
-
--enable-spark-declarative-pipeline:true(必填;启用 SDP 模式) -
ScriptLocation:管道定义 zip 或 Amazon S3 前缀(SDP 管道的必填项) -
--enable-glue-datacatalog:true(可选;在 Data Catalog 中注册数据表)
以下示例使用 AWS CLI 创建 SDP 作业:
aws glue create-job \ --name my-sdp-pipeline \ --role arn:aws:iam::123456789012:role/MyGlueRole \ --glue-version 6.0 \ --worker-type G.1X --number-of-workers 2 \ --command '{"Name":"glueetl","ScriptLocation":"s3://my-bucket/pipelines/my_pipeline.zip"}' \ --default-arguments '{ "--enable-spark-declarative-pipeline": "true", "--enable-glue-datacatalog": "true" }'
运行管道
您可以使用 StartJobRun 运行 SDP 管道。您可以使用在运行时传递的作业参数来控制执行行为。
运行模式
将以下参数传递给 StartJobRun 以控制管道执行:
--conf spark.glue.sdp.jobMode-
控制执行模式:
RUN(默认):正常执行管道。VALIDATE:执行试运行,在不写入任何数据的情况下检查 YAML 语法、依赖关系解析和 SQL/Python 编译。
--conf spark.glue.sdp.runMode-
控制刷新哪些数据集:
- 默认(无运行模式标志)
-
运行所有数据集。实体化视图完全重新计算;流式传输表仅处理自上次检查点以来的新数据。
--refresh <dataset>-
仅刷新指定的数据集。流式传输表以增量方式处理新数据;实体化视图则进行完整重新计算。
--full-refresh <dataset>-
仅重置和重新计算指定的数据集。对于流式传输表,这会重置检查点并重新处理所有数据。
--full-refresh-all-
重置并重新计算所有数据集。
将 Iceberg 表与 SDP 结合使用
对于需要持久、跨运行增量处理的流式传输表,Apache Iceberg 是推荐的表格式,因为它不依赖于 Hive 或 AWS Glue 托管表所使用的基于文件的 _spark_metadata 日志。您可以通过两种方式使用 SDP 配置 Iceberg,具体取决于是否要在数据目录中注册输出表。
选项 1:使用数据目录的 Iceberg(推荐)
如果您希望将 Iceberg 表注册到数据目录中(使用 table_type=ICEBERG),以便通过 Amazon Redshift 和 Amazon EMR 等其他引擎查询这些表,请使用此选项。表数据和元数据存储在 Amazon S3 中,并保留跨运行增量状态。将以下内容添加到您 spark-pipeline.yml 的 configuration 部分:
configuration: spark.sql.catalog.glue_catalog: "org.apache.iceberg.spark.SparkCatalog" spark.sql.catalog.glue_catalog.catalog-impl: "org.apache.iceberg.aws.glue.GlueCatalog" spark.sql.catalog.glue_catalog.io-impl: "org.apache.iceberg.aws.s3.S3FileIO" spark.sql.catalog.glue_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
在您的管道 YAML 中,设置 catalog: glue_catalog 并将 database: 设置为 AWS Glue 数据库。创建 AWS Glue 作业时,将 --enable-spark-declarative-pipeline 设置为 true。请勿为 Iceberg 表设置 --enable-glue-datacatalog。
注意
Iceberg 目录使用自己的 warehouse 位置存储表数据和元数据。使用此选项时,无需为 Iceberg 表本身设置 spark.sql.warehouse.dir 或数据库 LocationUri。
选项 2:使用基于文件的(Hadoop)目录的 Iceberg
当您不需要数据目录注册时,请使用此选项。Iceberg 元数据在 Amazon S3 中是基于文件的,并且这些表未在数据目录中注册。将以下内容添加到您 spark-pipeline.yml 的 configuration 部分:
configuration: spark.sql.catalog.spark_catalog: "org.apache.iceberg.spark.SparkSessionCatalog" spark.sql.catalog.spark_catalog.type: "hadoop" spark.sql.catalog.spark_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
注意
请注意有关此选项的以下信息:
-
使用此选项创建的表不会注册到数据目录中,因此无法通过查询引擎进行查询。
-
此选项使用自己的
warehouse位置来存储表。
如果您需要在数据目录中注册表并可从其他引擎查询,请使用选项 1。仅当不需要注册数据目录时,才使用选项 2。
注意
--enable-glue-datacatalog 作业参数将 Spark Hive 元数据仓连接到 Hive(非 Iceberg)表的数据目录。对于 Iceberg 表,catalog-impl=GlueCatalog 通过 AWS SDK 直接在数据目录中注册这些表,因此不要为 Iceberg 设置 --enable-glue-datacatalog。在 AWS Glue 6.0 上,请勿同时使用 Iceberg 的默认 SparkSessionCatalog(type: hive)配置和 --enable-glue-datacatalog,试图将 Iceberg 表注册到数据目录;这种组合会失败。
配置 Iceberg 后,您可以获得以下优势:
-
流式传输表在作业运行之间将检查点状态保存在 Amazon S3 中
-
每次运行都会创建新的数据文件和 Iceberg 快照
-
后续运行从上次提交的偏移量处恢复
-
完整的表历史记录通过 Iceberg 的快照机制保留
您也可以将 Iceberg 表作为流式数据源进行增量读取。在奖章架构中,下游流式传输表在每次运行时只能处理新提交到上游 Iceberg 表的行。以下示例以增量方式将数据从 Iceberg bronze 表读取到 silver 流式传输表中:
from pyspark import pipelines as dp from pyspark.sql import SparkSession spark = SparkSession.active() dp.create_streaming_table( "silver", comment="Incremental silver layer built from the bronze Iceberg table" ) @dp.append_flow(target="silver") def from_bronze(): # Incremental read from the Iceberg bronze table; each run processes only new rows. return spark.readStream.table("bronze")
注意事项和限制
在使用 SDP 时,请考虑以下几点:
-
实体化视图始终完全重新计算。不支持增量刷新。使用流式传输表处理增量工作负载。
-
流式传输表 Python API。将
dp.create_streaming_table()与@dp.append_flow(target=...)结合使用。 -
对流式传输表进行跨运行增量处理。仅当流式传输表的数据和检查点状态持久保存在 Amazon S3 上时,才支持跨运行增量处理。Hive 或 AWS Glue 托管(非 Iceberg)表要求将数据库
LocationUri设置为 Amazon S3 路径,而 Apache Iceberg 表则自行管理其表元数据。 -
Hive 或 AWS Glue 托管(非 Iceberg)表需要设置数据库 LocationUri。使用数据目录且采用
catalog-impl=GlueCatalog的 Iceberg 表不需要设置该项。有关更多信息,请参阅 先决条件。 -
数据质量预期。当前 SDP 框架不支持内联数据质量标注。
-
避免在下游查询函数中使用
withColumn。当下游数据集(例如实体化视图)使用spark.table(...)从上游管道数据集读取数据并应用.withColumn(...)时,SDP 可能无法在第二次及后续运行中检测到数据集之间的依赖关系。这会导致下游读取上次运行中的陈旧数据(单次运行延迟)。为避免此问题,请在.select(...)内表达派生列,而不是使用.withColumn(...)。还要避免在查询函数中强制执行计划解析的任何操作(例如.schema或.collect)。 -
无迁移工具。不支持从其他管道框架自动迁移。以增量方式迁移表:SDP 可从现有目录表中读取数据。
-
调度。SDP 作业与其他 AWS Glue 作业使用相同的调度机制(AWS Glue 触发器、Amazon EventBridge、Apache Airflow)。
从命令式脚本迁移到 SDP
您可以将现有的命令式 Spark 脚本以增量方式迁移到 SDP:
-
从一张表开始:将单个
spark.sql(...).write.saveAsTable(...)调用转换为CREATE MATERIALIZED VIEWSQL 语句。 -
以增量方式添加表。SDP 处理混合依赖关系。SDP 表可以从不属于管道的现有目录表中读取数据。
-
在过渡期间并行运行两种模式。SDP 作业与命令式作业可以共存。
SDP 可以引用任何可通过 SparkSession 访问的表,包括现有的 Data Catalog 表、外部表以及跨数据库引用。