Apache Airflow Spark Provider:Spark SQL 连接类型配置与执行机制详解
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
本文基于 Airflow 仓库中 Apache Spark Provider 的官方连接文档 spark-sql.rst,系统讲解spark_sql连接类型的配置字段、默认值解析规则,以及SparkSqlHook如何把连接参数拼装为真实的spark-sql命令行并在SparkSqlOperator中落地执行。读完后你可以独立完成 Spark SQL 连接的创建与调试,并理解命令拼装、SQL 内联/文件两种执行模式及进程生命周期管理的底层实现。
1. Spark SQL 连接类型概述
Apache Spark SQL 连接类型(connection type:spark_sql)用于通过spark-sql命令行工具连接 Apache Spark 集群或本地环境。与通过 JDBC/ODBC 建立数据库连接不同,该连接并不打开网络会话,而是驱动 Hook 在本地调用spark-sql二进制,由它向目标集群提交 SQL 作业。
从源码看(spark_sql.py),SparkSqlHook的类文档明确说明它是一个spark-sql二进制的包装器,要求spark-sql可执行文件位于 Airflow 组件所在环境的PATH中。该 Hook 的核心元信息定义如下:
class SparkSqlHook(BaseHook): conn_name_attr = "conn_id" default_conn_name = "spark_sql_default" conn_type = "spark_sql" hook_name = "Spark SQL"这段定义决定了三件事:Airflow UI 中新建该类型连接时使用的类型标识是spark_sql;未显式指定conn_id时默认读取spark_sql_default;Hook 展示名称为Spark SQL。
2. 默认连接 ID
官方文档明确:SparkSqlHook 默认使用spark_sql_default。这一点在源码中同样得到印证——Hook 与 Operator 的默认参数一致:
- Hook 中
default_conn_name = "spark_sql_default"(spark_sql.py),且__init__的conn_id参数默认值即为该类属性; SparkSqlOperator的conn_id参数默认值直接写为字符串"spark_sql_default"(spark_sql.py)。
因此在 Airflow UI 的 Connections 页面创建一条spark_sql类型、ID 为spark_sql_default的连接后,Operator/Hook 无需再显式传conn_id即可使用;若需要区分多套 Spark 环境(如本地调试与生产集群),可创建多条连接并按需传入conn_id。
3. 连接字段配置
官方文档定义了三个连接字段,结合 Hook 的表单扩展代码与 Provider 元数据 provider.yaml,各字段的含义与行为如下:
| 字段 | 是否必填 | 说明 |
|---|---|---|
| Host | 必填 | 要连接的目标,可为local、yarn或一个 URL(例如yarn://yarn-master、spark://host:port) |
| Port | 可选 | 仅当 Host 为 URL 形式时需要指定端口;Hook 会将其与 Host 拼接为host:port作为 master |
| YARN Queue | 可选 | 作业提交到的 YARN 队列名称,存储在连接的 Extra 字段(键名为queue)中,默认default |
其中 YARN Queue 是 Provider 为该连接类型注册的自定义表单控件。Hook 通过get_connection_form_widgets在 Airflow UI 的连接表单上追加了一个字符串输入框(spark_sql.py):
return { "queue": StringField( lazy_gettext("YARN queue"), widget=BS3TextFieldWidget(), description="Default YARN queue to use", validators=[Optional()], ) }对应的 Provider 元数据在 provider.yaml 中以conn-fields声明了queue字段(类型为 string 或 null、标签为YARN queue),保证 UI 渲染与元数据校验一致。
3.1 默认值解析规则(Host/Port → master)
Hook 在初始化时执行“显式参数优先、连接字段兜底”的解析逻辑(spark_sql.py):
try: conn = self.get_connection(conn_id) except AirflowNotFoundException: conn = None if conn: options = conn.extra_dejson # Set arguments to values set in Connection if not explicitly provided. if master is None: if conn is None: master = "yarn" elif conn.port: master = f"{conn.host}:{conn.port}" else: master = conn.host if yarn_queue is None: yarn_queue = options.get("queue", "default")即:
- 未找到连接且未显式传
master时,master 回退为"yarn"; - 连接带端口时,master 拼接为
host:port; - 队列从连接 Extra JSON 的
queue键读取,缺省为default。
单元测试中的示例连接Connection(conn_id="spark_default", conn_type="spark", host="yarn://yarn-master")展示了带 scheme 的 URL 写法(test_spark_sql.py),最终生成的命令中--master yarn://yarn-master与之一致。
4. 安全警示:Host 字段的 RCE 风险
官方文档对该连接类型给出了一条重要安全警告,必须原样重视:
警告:请谨慎授予用户修改 host 设置的权限,因为它可能使连接与外部服务器建立通信。需要明确认识到,将连接指向恶意服务器可能引发严重的安全漏洞,包括遭遇远程代码执行(RCE)攻击的风险。
从实现原理上理解这条警告:spark_sql连接的 Host 最终被原样传入--master参数,决定spark-sql客户端会连接并信任哪个资源管理器。若 DAG 作者(或能改连接配置的用户)可以任意指定 Host,就等价于把“向任意服务端提交作业/加载配置”的能力交给了不可信方——这正是文档所指 RCE 风险的来源。因此生产环境中应:
- 将连接创建权限收敛给管理员,普通 DAG 作者只能引用既有连接;
- 审查 DAG 中是否允许用户通过模板渲染注入
master、conf等参数; - 与 security.rst 所述的 Provider 级安全说明配合使用,评估 Spark 相关集成的整体暴露面。
5. Hook 实现:连接参数如何变成 spark-sql 命令
5.1 命令拼装(_prepare_command)
SparkSqlHook._prepare_command(spark_sql.py)负责把构造参数翻译为spark-sql的完整参数列表,拼装顺序为:
spark-sql [--conf key=value ...] # conf 参数,支持 dict 或逗号分隔的 "k=v,k2=v2" 字符串 [--total-executor-cores N] # 仅 Standalone & Mesos 适用 [--executor-cores N] # 每个 executor 的核心数 [--executor-memory SIZE] # 每个 executor 的内存,如 1000M、2G [--keytab PATH] # Kerberos keytab 文件完整路径 [--principal PRINCIPAL] # Kerberos principal [--num-executors N] # 启动的 executor 数量 [-f FILE | -e SQL] # SQL 文件或内联 SQL(见 5.2) [--master MASTER] # 连接解析出的目标 [--name NAME] # 作业名,默认 "default-name" [--verbose] # verbose 为 True 时追加(默认开启) [--queue QUEUE] # YARN 队列 [... 调用方额外传入的 cmd 参数]几个值得注意的实现细节:
- conf 双形态支持:
conf既可以是{"key": "value"}字典,也可以是"key=value,PROP=VALUE"字符串,后者按逗号拆分后逐项追加--conf(spark_sql.py)。单元测试test_build_command与test_build_command_with_str_conf分别对两种形态做了断言验证(test_spark_sql.py)。 - 额外参数透传:
run_query接受一个附加cmd(字符串按空白拆分,列表直接拼接),可传入--deploy-mode cluster等spark-sql合法参数;非法类型会抛出AirflowException(spark_sql.py)。 - 可观测性:最终命令会以 debug 级别打印为
Spark-Sql cmd: %s,便于排障时核对实际执行的命令行(spark_sql.py)。
5.2 SQL 内联与文件两种模式
Hook 对sql参数做了区分处理(spark_sql.py):
if self._sql: sql = self._sql.strip() if sql.endswith((".sql", ".hql")): connection_cmd += ["-f", sql] else: connection_cmd += ["-e", sql]- 以
.sql/.hql结尾时,先strip()去掉首尾空白,再以-f <path>方式执行 SQL 文件; - 否则作为内联 SQL 以
-e <sql>方式执行。
测试用例中" /path/to/sql/file.sql "这类带空白的路径,正是用来验证strip()后-f参数值与原始路径去空白结果一致的(test_spark_sql.py)。这也意味着:在 Operator 中把sql指向模板渲染后的.sql文件路径即可实现“文件型 SQL 作业”。
5.3 执行与进程管理(run_query/kill)
run_query(spark_sql.py)通过subprocess.Popen启动spark-sql进程,将 stderr 合并到 stdout(stderr=subprocess.STDOUT)并以文本模式逐行读取,每行都以INFO级别写入任务日志——即 Spark 客户端的输出会直接出现在 Airflow 任务日志中:
self._sp = subprocess.Popen( spark_sql_cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, universal_newlines=True, **kwargs ) for line in iter(self._sp.stdout): self.log.info(line) returncode = self._sp.wait() if returncode: raise AirflowException( f"Cannot execute '{self._sql}' on {self._master} (additional parameters: '{cmd}'). " f"Process exit code: {returncode}." )进程退出码非 0 时抛出AirflowException,错误信息包含 SQL 内容、目标 master、附加参数与退出码,测试test_spark_process_runcmd_and_fail精确断言了这一错误文案(test_spark_sql.py)。
kill方法(spark_sql.py)在任务被终止时调用:若进程仍在运行(poll() is None),直接Popen.kill()杀掉本地spark-sql客户端进程。
6. 在 SparkSqlOperator 中的使用
SparkSqlOperator(spark_sql.py)是连接类型的直接使用者,其关键点:
- 模板化字段:
template_fields = ("sql",)、template_ext = (".sql", ".hql"),并配置template_fields_renderers = {"sql": "sql"},因此sql支持 Jinja 模板与.sql/.hql文件渲染,UI 中按 SQL 高亮展示(spark_sql.py); - execute:惰性创建 Hook 并调用
run_query();on_kill:调用hook.kill()终止进程(spark_sql.py); - 参数透传:
_get_hook把 Operator 的全部构造参数(含conf、master、yarn_queue、keytab、principal等)原样传入SparkSqlHook(spark_sql.py)。
Provider 的系统测试 DAG 给出了最小可运行示例(example_spark_dag.py):
from airflow.providers.apache.spark.operators.spark_sql import SparkSqlOperator # [START howto_operator_spark_sql] spark_sql_job = SparkSqlOperator( sql="SELECT COUNT(1) as cnt FROM temp_table", master="local", task_id="spark_sql_job" ) # [END howto_operator_spark_sql]这里显式传master="local",会覆盖连接解析出的 master,说明“显式参数优先”规则在 Operator 层面同样生效。更多参数说明可参考官方操作文档 operators.rst。
7. UI 表单行为
除了追加YARN queue控件外,Hook 还通过get_ui_field_behaviour(spark_sql.py)隐藏了与spark_sql类型无关的通用字段:
return { "hidden_fields": ["schema", "login", "password", "extra"], "relabeling": {}, }这与 provider.yaml 中的ui-field-behaviour.hidden-fields声明一致——因为该连接不通过账号密码认证,而是依赖spark-sql客户端自身的 Kerberos/环境配置,所以 Schema、Login、Password 字段被隐藏,避免用户在 UI 上误填。
8. 关键机制的测试佐证
单元测试 test_spark_sql.py 与 test_spark_sql.py 覆盖了本文所述全部核心行为,可作为行为契约参考:
test_spark_process_runcmd:断言默认场景下(无显式参数)完整命令为spark-sql -e SELECT 1 --master yarn://yarn-master --name default-name --verbose --queue default,印证了 master 来自连接、name 默认值与 queue 缺省逻辑;test_spark_process_runcmd_with_str/with_list:验证字符串与列表两种形式追加参数(如--deploy-mode cluster)的正确性;test_spark_process_runcmd_and_fail:验证非零退出码时抛出带退出码的AirflowException。
9. 小结与相关文件索引
spark_sql连接类型的核心心智模型可以概括为:连接只描述“去哪”(host/port/queue),“怎么跑”由 Hook 的构造参数决定,二者按“显式参数优先、连接字段兜底”规则合并,最终拼装成一条spark-sql命令行执行。配置该连接时建议牢记 Host 的必填性与 RCE 风险,Kerberos 环境优先使用 keytab/principal 参数而非连接账号字段。
本文引用的仓库文件:
- 连接文档:spark-sql.rst
- Hook 实现:spark_sql.py
- Operator 实现:spark_sql.py
- Provider 元数据:provider.yaml
- Hook 单元测试:test_spark_sql.py
- 系统测试 DAG:example_spark_dag.py
- 操作文档:operators.rst
- Provider 安全说明:security.rst
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考