Apache Airflow 集成 Apache Spark 完全指南:apache-airflow-providers-apache-spark 6.3.2 安装、连接、算子与持久化执行
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
本篇技术指南基于当前仓库中的 Apache Airflow Spark Provider(apache-airflow-providers-apache-spark6.3.2)编写,完整介绍该 Provider 包在已有 Airflow 环境中的安装方式、版本依赖、三种连接类型(Spark Submit / Spark Connect / Spark SQL)的配置方法,以及五个核心算子(SparkSubmit、SparkJDBC、SparkSql、PySpark、SparkPipelines)与@task.pyspark装饰器的实战用法,并深入剖析基于 Task State Store 的持久化执行(durable execution)在 Standalone、Kubernetes、YARN 三种集群模式下的崩溃恢复原理。读完本文,你将能够在 Airflow 中以声明式 DAG 编排 Apache Spark 作业,并在 Worker 崩溃时实现不重复提交的作业恢复。
一、Provider 包概览
apache-airflow-providers-apache-spark是 Apache Airflow 官方维护的 Provider 发行版,用于对接 Apache Spark),主要能力包括:
- 通过
spark-submit脚本向 Spark 集群提交应用(支持 Standalone、YARN、Kubernetes 等多种集群管理器与部署模式); - 通过
spark-sql脚本在 Spark Hive Metastore 服务上执行 SQL 查询; - 通过 Spark Connect 接口远程使用 DataFrame API 运行 PySpark 任务;
- 通过 JDBC 在 Spark 与关系型数据库之间双向传输数据;
- 支持 Spark Declarative Pipelines 声明式数据管道;
- 在 Airflow 3.3+ 上支持持久化执行(durable execution),实现 Worker 崩溃后的作业恢复。
完整的类清单、变更日志与包信息可参考 providers/apache/spark/docs/index.rst 中自动生成的包文档索引。
二、安装
该 Provider 包设计为安装于已有的 Airflow 环境之上。推荐方式是通过 PyPI 安装:
pip install apache-airflow-providers-apache-spark包支持以下 Python 版本:3.10、3.11、3.12、3.13、3.14。
如需从源码安装(例如基于本仓库开发调试),可参考 providers/apache/spark/docs/installing-providers-from-sources.rst 中的步骤;包内的pyproject.toml(见 providers/apache/spark/pyproject.toml)与provider.yaml(见 providers/apache/spark/provider.yaml)是构建与元数据声明的权威来源。
官方发布的 sdist/wheel 包及其
asc签名、sha512校验文件可从 Apache 官方下载站获取,安装前建议校验签名与校验和,保证供应链安全。
三、版本要求与依赖
3.1 最低 Airflow 版本
该 Provider 发行版支持的最低 Apache Airflow 版本为2.11.0(apache-airflow>=2.11.0)。需要注意:持久化执行(durable execution)特性需要 Airflow 3.3+(依赖task_state_store支持),在低于 3.3 的版本上durable参数不会生效,仅会在显式设置时输出警告,算子行为与旧版本完全一致——每次重试都重新提交全新作业。
3.2 核心 PIP 依赖
| PIP 包 | 版本要求 |
|---|---|
apache-airflow | >=2.11.0 |
apache-airflow-providers-common-compat | >=1.12.0 |
pyspark-client | >=4.0.0 |
grpcio-status | >=1.67.0 |
requests | >=2.32.0 |
tenacity | >=8.3.0 |
其中pyspark-client与grpcio-status服务于 Spark Connect 客户端连接;tenacity被用于驱动状态查询的重试逻辑(见 spark_submit.py 中_fetch_driver_status上的@retry装饰器:最多尝试 3 次、每次间隔 1 秒)。
3.3 跨 Provider 可选依赖(Cross Provider Dependencies)
使用部分高级特性时,需要额外安装其他 Provider 发行版。安装时可通过 extra 语法一并拉取:
pip install apache-airflow-providers-apache-spark[cncf.kubernetes]| 依赖的 Provider 包 | Extra |
|---|---|
apache-airflow-providers-cncf-kubernetes | cncf.kubernetes |
该 extra 用于 Kubernetes 集群模式下通过 Python Kubernetes 客户端跟踪 Driver Pod 状态(见下文「Kubernetes 集群模式」小节)。
3.4 可选依赖(Optional Dependencies)
| Extra | 依赖 |
|---|---|
cncf.kubernetes | apache-airflow-providers-cncf-kubernetes>=7.4.0 |
openlineage | apache-airflow-providers-openlineage |
pyspark | pyspark>=4.0.0 |
openlineageextra 启用 OpenLineage 血缘注入。SparkSubmitOperator支持通过 Airflow 配置openlineage.spark_inject_parent_job_info与openlineage.spark_inject_transport_info将父作业信息与传输信息注入 Spark 属性(见 spark_submit.py 中execute()的注入逻辑)。pysparkextra 提供本地 PySpark 运行时,供 PySparkOperator 的独立(standalone)模式与@task.pyspark装饰器使用。
四、连接(Connection)配置
该 Provider 提供三种连接类型,分别服务于不同的执行通道。配置入口均为 Airflow 管理界面的 Connections 页面或airflow connectionsCLI。
4.1 Spark Submit 连接
Spark Submit 连接通过spark-submit命令连接 Apache Spark,默认连接 ID 为spark_default(Spark Submit / Spark JDBC 的 Hook 与算子默认使用该 ID)。字段说明如下:
| 字段 | 必填 | 说明 |
|---|---|---|
| Host | 是 | 要连接的主机,可以是local、yarn或一个 URL |
| Port | 否 | 当 Host 为 URL 时指定端口 |
| YARN Queue | 否 | YARN 上提交应用时使用的队列名(仅适用于 Spark on YARN) |
| Deploy mode | 否 | driver 部署在 Worker 节点上(cluster)还是作为外部客户端本地运行(client) |
| Spark binary | 否 | 用于提交的命令,部分发行版使用spark2-submit;默认spark-submit,仅允许spark-submit、spark2-submit、spark3-submit三个取值 |
| Kubernetes namespace | 否 | Kubernetes 命名空间(spark.kubernetes.namespace),用于通过资源配额划分集群资源(仅适用于 Spark on Kubernetes) |
| REST scheme | 否 | 访问 Spark Standalone REST API 使用的协议(http或https),默认http;当集群开启 TLS(spark.ssl.standalone.enabled=true)时设为https |
| REST port | 否 | Spark Standalone REST API 端口(对应spark.master.rest.port),默认6066 |
当通过环境变量配置连接时,必须使用 URI 语法。可以直接提供标准的 Spark master URI,URL 编码后各组件会被正确解析,无需重复spark://spark://...前缀:
export AIRFLOW_CONN_SPARK_DEFAULT='spark://mysparkcluster.com:80?deploy-mode=cluster&spark_binary=command&namespace=kube+namespace'安全警告:必须信任有权限配置 Host 设置的用户。将连接指向恶意服务器可能带来严重安全漏洞,包括远程代码执行(RCE)风险。
4.2 Spark Connect 连接
Spark Connect 连接通过 Spark Connect 接口连接 Spark,默认连接 ID 为spark_connect_default。字段说明:
| 字段 | 必填 | 说明 |
|---|---|---|
| Host | 是 | 有效的主机名 |
| Port | 否 | 当 Host 为 URL 时指定端口 |
| User ID | 否 | 用于向代理认证的用户 ID(仅 Spark Connect) |
| Token | 否 | 用于向代理认证的令牌(仅 Spark Connect) |
| Use SSL | 否 | 连接时是否使用 SSL(仅 Spark Connect) |
Spark Connect 本身没有内置认证,但其 gRPC HTTP/2 接口允许通过认证代理与 Spark Connect 服务器通信。需要认证时,创建 Spark Connect 连接并填入正确的凭据即可。
4.3 Spark SQL 连接
Spark SQL 连接通过spark-sql命令连接 Spark,默认连接 ID 为spark_sql_default(SparkSqlHook 使用)。字段说明:
| 字段 | 必填 | 说明 |
|---|---|---|
| Host | 是 | 要连接的主机,可以是local、yarn或一个 URL |
| Port | 否 | 当 Host 为 URL 时指定端口 |
| YARN Queue | 否 | 提交应用使用的 YARN 队列名 |
4.4 使用 JDBC 连接
SparkJDBCOperator除 Spark 连接外,还需要一个 JDBC 连接(jdbc-default为默认 ID)来访问关系型数据库。在算子的jdbc_conn_id参数中指定。
五、算子详解
该 Provider 提供五个核心算子,完整的参数定义参见 providers/apache/spark/docs/operators.rst。各算子的前置条件如下:
| 算子 | 前置条件 |
|---|---|
SparkSubmitOperator | 配置 Spark Submit 连接(connections/spark-submit.rst) |
SparkJDBCOperator | 配置 Spark Submit 连接 + JDBC 连接 |
SparkSqlOperator | 所有配置均来自算子参数 |
PySparkOperator | 可选配置 Spark Connect 连接(connections/spark-connect.rst) |
SparkPipelinesOperator | 配置 Spark Submit 连接,且环境中存在spark-pipelinesCLI |
5.1 SparkSubmitOperator —— 通用应用提交
SparkSubmitOperator封装了spark-submit脚本:该脚本负责配置 Spark 及其依赖的 classpath,并支持 Spark 支持的各种集群管理器与部署模式。要求spark-submit二进制位于 PATH 中。核心参数(实现见 spark_submit.py):
| 参数 | 说明 | 默认值 |
|---|---|---|
application | 提交的应用(jar 或 py 文件),支持模板化 | 必填 |
conf | 任意 Spark 配置属性(dict),支持模板化 | {} |
conn_id | Spark 连接 ID,无效连接默认回退到yarn | spark_default |
files/py_files/jars | 上传给执行器的附加文件 / Python 文件 / jar | None |
java_class | Java 应用的主类 | None |
packages/exclude_packages/repositories | Maven 坐标、排除项与附加仓库 | None |
executor_cores | 每个执行器核数(Standalone & YARN) | 2 |
executor_memory | 每个执行器内存(如1000M、2G) | 1G |
driver_memory | driver 内存 | 1G |
num_executors | 启动的执行器数量 | None |
name | 作业名称 | arrow-spark |
keytab/principal/proxy_user | Kerberos 与代理用户配置,会覆盖连接 extra 中的同名配置 | None |
deploy_mode | cluster或client,覆盖连接配置 | None |
yarn_queue | YARN 队列名,覆盖连接配置 | None |
status_poll_interval | 集群模式下轮询 driver 状态的间隔(秒);YARN RM API 轮询最少 10 秒 | 1 |
application_args/env_vars | 传给应用的参数 / spark-submit 的环境变量(支持 yarn 与 k8s 模式) | None |
verbose | 是否向 spark-submit 进程传递 verbose 标志 | False |
spark_binary | 提交命令,覆盖连接配置 | None |
post_submit_commands | 作业结束后执行的 shell 命令列表(如清理 Istio 等 sidecar),失败仅告警不使任务失败 | None |
durable | 是否将外部作业 ID 持久化到 Task State Store 以便重试时重连 | True |
track_driver_via_k8s_api | K8s 集群模式下释放 spark-submit JVM,改由 Kubernetes API 轮询 Pod 阶段(轮询间隔最少 20 秒) | False |
yarn_track_via_rm_api | YARN 集群模式下释放 spark-submit JVM,改由 YARN ResourceManager REST API 轮询 | False |
yarn_rm_auth | 自定义requests.auth.AuthBase,用于 RM REST 调用 | None |
reconnect_on_retry | 已弃用,是durable的旧名称,仍会发出弃用警告并映射到durable | None |
模板化字段(template_fields)包括:application、conf、files、py_files、jars、driver_class_path、packages、exclude_packages、keytab、principal、proxy_user、name、application_args、env_vars、post_submit_commands、properties_file。
最简单的提交示例(来自 example_spark_dag.py):
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator submit_job = SparkSubmitOperator( application="${SPARK_HOME}/examples/src/main/python/pi.py", task_id="submit_job" )5.2 SparkJDBCOperator —— Spark 与 JDBC 数据库双向传输
SparkJDBCOperator继承自SparkSubmitOperator(见 spark_jdbc.py),内部使用SparkSubmitOperator执行 Spark 与 JDBC 数据库之间的数据传输。通过cmd_type参数控制数据流向:
spark_to_jdbc:Spark 将数据从 metastore 写入 JDBC 表;jdbc_to_spark:Spark 从 JDBC 表读取数据写入 metastore(使用 Spark 命令saveAsTable)。
关键参数:
| 参数 | 说明 |
|---|---|
cmd_type | 数据流向:spark_to_jdbc/jdbc_to_spark |
jdbc_table | JDBC 表名 |
jdbc_conn_id | JDBC 连接 ID(默认jdbc-default) |
jdbc_driver | JDBC 驱动类名(驱动 jar 需通过spark_jars传入) |
metastore_table | metastore 表名 |
jdbc_truncate | (spark_to_jdbc 专用)save_mode为 Overwrite 时,Spark 是截断还是删表重建;schema 不同时无法截断会删表重建 |
save_mode | Spark save-mode(如overwrite、append) |
save_format | (jdbc_to_spark 专用)保存格式(如parquet) |
batch_size | (spark_to_jdbc 专用)每轮写入 JDBC 的批大小,默认 1000 |
fetch_size | (jdbc_to_spark 专用)每轮从 JDBC 抓取的批大小,默认取决于驱动 |
num_partitions | Spark 可同时使用的最大分区数,同时限制可打开的 JDBC 连接数 |
partition_column/lower_bound/upper_bound | (jdbc_to_spark 专用)按数值列分区抓取,三者需与num_partitions同时指定 |
create_table_column_types | (spark_to_jdbc 专用)建表时的列类型(如"name CHAR(64), comments VARCHAR(1024)") |
完整示例:
from airflow.providers.apache.spark.operators.spark_jdbc import SparkJDBCOperator jdbc_to_spark_job = SparkJDBCOperator( cmd_type="jdbc_to_spark", jdbc_table="foo", spark_jars="${SPARK_HOME}/jars/postgresql-42.2.12.jar", jdbc_driver="org.postgresql.Driver", metastore_table="bar", save_mode="overwrite", save_format="JSON", task_id="jdbc_to_spark_job", ) spark_to_jdbc_job = SparkJDBCOperator( cmd_type="spark_to_jdbc", jdbc_table="foo", spark_jars="${SPARK_HOME}/jars/postgresql-42.2.12.jar", jdbc_driver="org.postgresql.Driver", metastore_table="bar", save_mode="append", task_id="spark_to_jdbc_job", )5.3 SparkSqlOperator —— Spark SQL 查询执行
SparkSqlOperator在 Spark 服务器上启动应用,要求spark-sql脚本位于 PATH 中。算子会在 Spark Hive Metastore 服务上执行 SQL 查询,sql参数支持模板化,且可以是.sql或.hql文件(template_ext为(".sql", ".hql"),见 spark_sql.py)。关键参数包括sql(必填)、conn_id(默认spark_sql_default)、master(spark://host:port、mesos://host:port、yarn或local,默认取连接中的 host/port,否则yarn)、executor_cores(默认 2)、executor_memory(默认 1G)、yarn_queue(默认连接中的 queue 或default)。
示例:
from airflow.providers.apache.spark.operators.spark_sql import SparkSqlOperator spark_sql_job = SparkSqlOperator( sql="SELECT COUNT(1) as cnt FROM temp_table", master="local", task_id="spark_sql_job" )5.4 PySparkOperator —— Spark Connect / 独立模式运行 PySpark
PySparkOperator将 PySpark 作业提交到外部 Spark Connect 服务,或在独立(standalone)模式下直接运行(见 spark_pyspark.py)。它以PythonOperator为基础,自动完成 SparkSession 的构建与销毁:
- 未指定
conn_id时 master 默认为local[*]; - 连接类型为 Spark Connect 时,使用
SparkConnectHook解析出sc://形式的 URL 并设置spark.remote; - 连接的 extra JSON 会逐项写入 SparkConf;
config_kwargs参数可覆盖连接中的 Spark 配置;- 调用可执行函数前,
spark会话对象会注入op_kwargs,函数结束后spark_session.stop()释放资源。
示例(来自 example_spark_dag.py):
from airflow.providers.apache.spark.operators.spark_pyspark import PySparkOperator def my_pyspark_job(spark): df = spark.range(100).filter("id % 2 = 0") print(df.count()) spark_pyspark_job = PySparkOperator( python_callable=my_pyspark_job, conn_id="spark_connect", task_id="spark_pyspark_job" )5.5 SparkPipelinesOperator —— 声明式数据管道
SparkPipelinesOperator使用spark-pipelinesCLI 执行 Spark Declarative Pipelines,同时支持管道执行与 dry-run 校验(见 spark_pipelines.py)。
from airflow.providers.apache.spark.operators.spark_pipelines import SparkPipelinesOperator # 执行管道 run_pipeline = SparkPipelinesOperator( task_id="run_pipeline", pipeline_spec="/path/to/pipeline.yml", pipeline_command="run", conn_id="spark_default", num_executors=2, executor_cores=4, executor_memory="2G", driver_memory="1G", )pipeline_spec指向定义声明式管道的 YAML 文件,例如:
name: my_pipeline storage: file:///path/to/pipeline-storage libraries: - glob: include: transformations/**支持两种管道命令:
run:执行管道(默认);dry-run:仅校验管道,不执行。
六、持久化执行(Durable Execution)与崩溃恢复
在集群部署模式(--deploy-mode cluster)下,Spark Driver 独立于 Airflow Worker 运行。若 Worker 在 Spark 作业运行期间崩溃,Driver 仍会继续运行,但 Airflow 会丢失对它的跟踪;此时若重试并重新提交全新作业,会浪费已完成的计算,甚至可能因作业非幂等而引发冲突。
SparkSubmitOperator通过将Driver 标识符持久化到 Task State Store解决这一问题:提交后立即持久化,重试时读取标识符并重连到已在运行的 Driver,而不是重新提交(同步路径,见 spark_submit.py 中的execute_resumable流程:submit_job→poll_until_complete→get_job_result)。
durable默认True,所有集群管理器均支持;- 需要Airflow 3.3+(
task_state_store支持);低于 3.3 时durable不生效; - 清理(Clear)任务与重试等价处理:若 Driver 已成功,清理不会删除存储的 Driver ID,下次尝试会读取并立即返回而不重新提交。可通过
[state_store] clear_on_success配置恢复"清理总是重新提交"的行为(详见 Airflow 核心概念中的 resumable-tasks 文档); - 对延迟执行任务(
deferrable=True)最可靠:同步轮询期间清理任务可能先通过on_kill取消 Driver,使下次尝试来不及重连; - 各集群管理器的重试前置条件不同:
| 集群管理器 | durable=True的前置条件 |
|---|---|
| Spark Standalone | 连接 extra 中配置REST scheme与REST port |
| Kubernetes cluster mode | track_driver_via_k8s_api=True |
| YARN cluster mode | yarn_track_via_rm_api=True与yarn_resourcemanager_webapp_address |
6.1 Spark Standalone:REST API 重连
重连运行中的 Driver 需要调用 Spark Standalone REST API(GET /v1/submissions/status/{driverId})。需确保 Spark 连接的REST scheme与REST portextra 与集群配置一致:
REST scheme:集群 REST 端口启用 TLS(spark.ssl.standalone.enabled=true)时设为https,默认http;REST port:对应集群的spark.master.rest.port,默认6066。
从源码看(_StandaloneSparkSubmitBackend),对 HA 形态的多 master URL(如spark://m1:7077,m2:7077)会逐个尝试;master URL 中的端口(如 7077)是 RPC 端口而非 REST API 端口。is_job_active认为SUBMITTED、RUNNING、RELAUNCHING、UNKNOWN四种状态仍在运行,FINISHED视为成功;若轮询结束时状态非FINISHED,会抛出RuntimeError,且无论成功失败都会执行post_submit_commands。
6.2 Kubernetes 集群模式:Kubernetes API 跟踪
K8s 集群模式下,spark-submit默认会阻塞整个作业周期,JVM 仅用于轮询 Pod 阶段并持有堆内存,这对长作业(尤其是 Driver 长时间空闲)并不理想。
设置track_driver_via_k8s_api=True后,算子改为通过 Python Kubernetes 客户端跟踪 Driver Pod 状态,释放长时间占用的spark-submitJVM。同一标志也使算子能在崩溃后重新找到 Driver:Driver Pod 名在轮询开始前持久化到 Task State Store,重试时重连该 Pod 而非重新提交。从_KubernetesSparkSubmitBackend的实现看:提交时会强制设置spark.kubernetes.submission.waitAppCompletion=false,取到 Pod 名后以namespace:pod_name作为外部作业 ID;轮询结束时若 Pod 到达Succeeded且durable=True,会把k8s_driver_status=Succeeded缓存到 Task State Store——即使 Driver Pod 已被垃圾回收,重试也能识别出已完成作业(404/消失路径不会缓存,避免重试误判成功而跳过重新提交)。
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator run_spark = SparkSubmitOperator( task_id="run_spark", application="local:///opt/spark/examples/jars/spark-examples.jar", conn_id="spark_k8s", deploy_mode="cluster", track_driver_via_k8s_api=True, durable=True, )使用要求:
- Spark 连接的
master必须为k8s://...且deploy_mode为cluster; - 不要在
conf中设置spark.kubernetes.submission.waitAppCompletion=true,否则任务启动时会抛出ValueError; - Airflow Worker 必须能访问 Kubernetes API Server,并具有读取与删除 Driver 命名空间内 Pod 的权限;
- Pod 完成状态依据
pod.status.phase判断;若 Driver Pod 存在 sidecar 容器(如启用了 Istio 注入),Pod 阶段可能不会推进到Succeeded,轮询将无限等待——应设置execution_timeout作为硬上限; - 设置
durable=False可让重试总是提交全新 Driver。
6.3 YARN 集群模式:ResourceManager REST API 跟踪
YARN 集群模式下,默认 Spark 提交路径会让本地spark-submitJVM 在整个应用生命周期内存活于 Airflow Worker 上,长时间占用 Worker 内存。
设置yarn_track_via_rm_api=True后,YARN 接受应用后即释放本地spark-submitJVM,改由轮询 YARN ResourceManager REST API(GET /ws/v1/cluster/apps/{appId})直到应用到达终态。轮询间隔由status_poll_interval控制,最小 10 秒,避免对 ResourceManager 造成压力。该标志同样是 YARN 上持久化执行的前置条件:由于durable默认True,不设置该标志会在任务启动时抛出ValueError,而不是静默退化为 fire-and-forget 提交。
提交前需要在 Spark 连接的 extra 中配置yarn_resourcemanager_webapp_address:
airflow connections add spark_yarn_rm \ --conn-type spark \ --conn-host yarn \ --conn-extra '{ "deploy-mode": "cluster", "yarn_resourcemanager_webapp_address": "http://rm.example.com:8088" }'SparkSubmitOperator( task_id="spark_pi", conn_id="spark_yarn_rm", application="/path/to/spark-examples.jar", java_class="org.apache.spark.examples.SparkPi", deploy_mode="cluster", yarn_track_via_rm_api=True, )Kerberos 集群:需要在 Airflow 环境中安装requests-kerberos。当 Spark 连接同时配置了keytab与principal时,Airflow 自动使用HTTPKerberosAuth()发起 ResourceManager REST 请求。
自定义认证:仅当 ResourceManager 需要自定义requests认证对象时使用yarn_rm_auth:
import requests SparkSubmitOperator( task_id="spark_pi", conn_id="spark_yarn_rm", application="/path/to/spark-examples.jar", java_class="org.apache.spark.examples.SparkPi", deploy_mode="cluster", yarn_track_via_rm_api=True, yarn_rm_auth=requests.auth.HTTPBasicAuth("user", "password"), )从_YarnSparkSubmitBackend的实现看:提交时会强制spark.yarn.submit.waitAppCompletion=false以立即退出 spark-submit 并捕获应用 ID;若显式设置了该配置为true会抛出ValueError。on_kill时若已拿到应用 ID,则通过 RM REST API 杀死应用(此时 CLI 方式的 kill 已无目标);is_job_active依据 Hadoop ResourceManager REST 文档,将NEW、NEW_SAVING、SUBMITTED、ACCEPTED、RUNNING视为活跃。
七、@task.pyspark 装饰器
@task.pyspark装饰器(见 decorators/pyspark.rst 与 pyspark.py)会把 SparkSession(spark)与 SparkContext(sc)对象注入被装饰的 Python 可调用对象。参数:
| 参数 | 说明 |
|---|---|
conn_id | 连接 Spark 集群使用的连接 ID;未指定时 master 设为local[*] |
config_kwargs | 初始化 SparkConf 使用的 kwargs,会覆盖连接中设置的 Spark 配置项 |
示例(来自 example_pyspark.py,spark对象自动注入):
@task.pyspark(conn_id="spark-local") def spark_task(spark: SparkSession) -> pd.DataFrame: df = spark.createDataFrame( [ (1, "John Doe", 21), (2, "Jane Doe", 22), (3, "Joe Bloggs", 23), ], ["id", "name", "age"], ) df.show() return df.toPandas()推荐使用 Spark Connect:Apache Spark 3.4 引入的 Spark Connect 是解耦的客户端-服务器架构,允许通过 DataFrame API 远程连接 Spark 集群。在 Airflow 中使用 PySpark 装饰器时,Spark Connect 是首选方式,因为无需在 Airflow 所在主机上运行 Spark Driver。使用方式是在主机 URL 前加sc://前缀,例如sc://spark-cluster:15002。认证方面,Spark Connect 无内置认证,可通过认证代理利用 gRPC HTTP/2 接口完成认证(创建 Spark Connect 连接并设置正确凭据)。
八、完整示例 DAG 与系统测试
仓库提供了两个可直接运行的系统测试示例:
- example_spark_dag.py:同时演示
SparkSubmitOperator、SparkJDBCOperator(两种流向)、SparkSqlOperator与PySparkOperator的完整 DAG; - example_pyspark.py:演示
@task.pyspark装饰器,将 Spark DataFrame 转 Pandas 后交给下游普通任务处理。
单元测试覆盖了全部 Hook 与算子,包括 test_spark_submit.py、test_spark_jdbc.py、test_spark_sql.py、test_spark_pyspark.py、test_spark_pipelines.py 以及 test_spark_submit_post_commands.py 等,可作为行为规范的参考依据。
九、源码结构导航
- 算子实现:operators 目录,含
spark_submit.py、spark_jdbc.py、spark_sql.py、spark_pyspark.py、spark_pipelines.py; - Hook 实现:hooks 目录,含
spark_submit.py、spark_jdbc.py、spark_connect.py、spark_sql.py、spark_pipelines.py、spark_jdbc_script.py; - 装饰器实现:decorators/pyspark.py;
- 包元数据:provider.yaml 与 pyproject.toml;
- 官方文档源:docs 目录(连接、算子、装饰器、安全、变更日志等)。
十、变更日志与安全公告
- 变更日志:providers/apache/spark/docs/changelog.rst;
- 安全公告:providers/apache/spark/docs/security.rst。
在升级 Provider 版本前,务必阅读变更日志以了解行为变更(例如reconnect_on_retry到durable的参数更名、K8s/YARN 跟踪模式对conf的限制等),并遵循安全公告中的最佳实践。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考