
Apache Airflow Spark ProviderSpark 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 typespark_sql用于通过spark-sql命令行工具连接 Apache Spark 集群或本地环境。与通过 JDBC/ODBC 建立数据库连接不同该连接并不打开网络会话而是驱动 Hook 在本地调用spark-sql二进制由它向目标集群提交 SQL 作业。从源码看spark_sql.pySparkSqlHook的类文档明确说明它是一个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_defaultHook 展示名称为Spark SQL。2. 默认连接 ID官方文档明确SparkSqlHook 默认使用spark_sql_default。这一点在源码中同样得到印证——Hook 与 Operator 的默认参数一致Hook 中default_conn_name spark_sql_defaultspark_sql.py且__init__的conn_id参数默认值即为该类属性SparkSqlOperator的conn_id参数默认值直接写为字符串spark_sql_defaultspark_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:portPort可选仅当 Host 为 URL 形式时需要指定端口Hook 会将其与 Host 拼接为host:port作为 masterYARN Queue可选作业提交到的 YARN 队列名称存储在连接的 Extra 字段键名为queue中默认default其中 YARN Queue 是 Provider 为该连接类型注册的自定义表单控件。Hook 通过get_connection_form_widgets在 Airflow UI 的连接表单上追加了一个字符串输入框spark_sql.pyreturn { queue: StringField( lazy_gettext(YARN queue), widgetBS3TextFieldWidget(), descriptionDefault YARN queue to use, validators[Optional()], ) }对应的 Provider 元数据在 provider.yaml 中以conn-fields声明了queue字段类型为 string 或 null、标签为YARN queue保证 UI 渲染与元数据校验一致。3.1 默认值解析规则Host/Port → masterHook 在初始化时执行“显式参数优先、连接字段兜底”的解析逻辑spark_sql.pytry: 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_idspark_default, conn_typespark, hostyarn://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_commandSparkSqlHook._prepare_commandspark_sql.py负责把构造参数翻译为spark-sql的完整参数列表拼装顺序为spark-sql [--conf keyvalue ...] # conf 参数支持 dict 或逗号分隔的 kv,k2v2 字符串 [--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}字典也可以是keyvalue,PROPVALUE字符串后者按逗号拆分后逐项追加--confspark_sql.py。单元测试test_build_command与test_build_command_with_str_conf分别对两种形态做了断言验证test_spark_sql.py。额外参数透传run_query接受一个附加cmd字符串按空白拆分列表直接拼接可传入--deploy-mode cluster等spark-sql合法参数非法类型会抛出AirflowExceptionspark_sql.py。可观测性最终命令会以 debug 级别打印为Spark-Sql cmd: %s便于排障时核对实际执行的命令行spark_sql.py。5.2 SQL 内联与文件两种模式Hook 对sql参数做了区分处理spark_sql.pyif 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/killrun_queryspark_sql.py通过subprocess.Popen启动spark-sql进程将 stderr 合并到 stdoutstderrsubprocess.STDOUT并以文本模式逐行读取每行都以INFO级别写入任务日志——即 Spark 客户端的输出会直接出现在 Airflow 任务日志中self._sp subprocess.Popen( spark_sql_cmd, stdoutsubprocess.PIPE, stderrsubprocess.STDOUT, universal_newlinesTrue, **kwargs ) for line in iter(self._sp.stdout): self.log.info(line) returncode self._sp.wait() if returncode: raise AirflowException( fCannot execute {self._sql} on {self._master} (additional parameters: {cmd}). fProcess 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 中的使用SparkSqlOperatorspark_sql.py是连接类型的直接使用者其关键点模板化字段template_fields (sql,)、template_ext (.sql, .hql)并配置template_fields_renderers {sql: sql}因此sql支持 Jinja 模板与.sql/.hql文件渲染UI 中按 SQL 高亮展示spark_sql.pyexecute惰性创建 Hook 并调用run_query()on_kill调用hook.kill()终止进程spark_sql.py参数透传_get_hook把 Operator 的全部构造参数含conf、master、yarn_queue、keytab、principal等原样传入SparkSqlHookspark_sql.py。Provider 的系统测试 DAG 给出了最小可运行示例example_spark_dag.pyfrom airflow.providers.apache.spark.operators.spark_sql import SparkSqlOperator # [START howto_operator_spark_sql] spark_sql_job SparkSqlOperator( sqlSELECT COUNT(1) as cnt FROM temp_table, masterlocal, task_idspark_sql_job ) # [END howto_operator_spark_sql]这里显式传masterlocal会覆盖连接解析出的 master说明“显式参数优先”规则在 Operator 层面同样生效。更多参数说明可参考官方操作文档 operators.rst。7. UI 表单行为除了追加YARN queue控件外Hook 还通过get_ui_field_behaviourspark_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.rstHook 实现spark_sql.pyOperator 实现spark_sql.pyProvider 元数据provider.yamlHook 单元测试test_spark_sql.py系统测试 DAGexample_spark_dag.py操作文档operators.rstProvider 安全说明security.rst【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考