ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Airflow Spark Provider:Spark Submit Connection 连接配置详解

Apache Airflow Spark Provider:Spark Submit Connection 连接配置详解 Apache Airflow Spark ProviderSpark Submit Connection 连接配置详解【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文围绕 Apache Airflow 的 Spark Provider 中的Spark Submit 连接类型展开如何为SparkSubmitHook/SparkSubmitOperator正确配置spark_default连接Host、Deploy mode、Spark binary、YARN Queue、Kubernetes namespace、Standalone REST 参数等如何通过环境变量 URI 语法注入连接以及这些连接字段在源码中如何被解析并映射到最终的spark-submit命令行。读完本文你可以针对不同集群管理器Standalone、YARN、Kubernetes独立完成 Spark 连接配置并理解 cluster 模式下 driver 状态跟踪、REST 轮询与崩溃恢复durable execution的配置前提。连接类型定位与默认 Connection IDSpark Submit 连接类型用于通过spark-submit命令连接 Apache Spark。它是apache-airflow-providers-apache-spark中SparkSubmitHook与SparkSubmitOperator使用的连接连接类型标识为spark见 provider.yaml 中connection-types段hook-class-name指向 SparkSubmitHook默认连接 ID 为spark_defaultSpark Submit 与 Spark JDBC 的 hooks 和 operators 在未显式指定conn_id时都会使用它。在源码中可以直接验证——hook 与 operator 的定义如下spark_submit.py hookconn_name_attr conn_id default_conn_name spark_default conn_type spark hook_name Spark因此配置该连接后即使 operator 中不写conn_id也会自动读取spark_default如果该连接不存在hook 会回退到masteryarn的默认行为见下文源码解析。连接配置字段详解在 Airflow UI 的 Connections 页面中创建类型为spark的连接时各字段含义如下与 原始连接文档 一一对应字段是否必填说明Host必填要连接的 Spark 集群地址可以是local、yarn或一个 URL如spark://HOST:PORT、k8s://...Port可选当 Host 是 URL 时指定端口YARN Queue可选仅 Spark on YARN 生效应用提交到的 YARN 队列名称Deploy mode可选取值client或clusterdriver 部署在 worker 节点上cluster还是作为本地外部客户端运行clientSpark binary可选用于提交作业的命令某些发行版使用spark2-submit或spark3-submit。默认spark-submit只允许spark-submit、spark2-submit、spark3-submit三个取值Kubernetes namespace可选仅 Spark on Kubernetes 生效对应spark.kubernetes.namespace用于在多个用户之间通过 resource quota划分集群资源REST scheme可选仅 Spark Standalone 集群模式生效访问 Spark standalone REST API 的方案http或https默认http当 Spark master REST API 启用了 TLSspark.ssl.standalone.enabledtrue时应设为httpsREST port可选仅 Spark Standalone 集群模式生效Spark standalone REST API 端口对应spark.master.rest.port默认6066集群使用非默认 REST 端口时需要覆盖这些字段不只是文档约定它们在源码中有两处对应的强约束实现UI 表单控件SparkSubmitHook.get_connection_form_widgets() 定义了每个控件的标签、描述与校验器。例如deploy-mode字段使用any_of([client, cluster])校验、默认clientspark-binary使用any_of(ALLOWED_SPARK_BINARIES)校验其中DEFAULT_SPARK_BINARY spark-submit ALLOWED_SPARK_BINARIES [DEFAULT_SPARK_BINARY, spark2-submit, spark3-submit]同样的字段定义也在 provider.yaml 中以conn-fields声明含默认值与描述保证连接表单、provider 元数据与 hook 行为三者一致。字段解析hook 在初始化时通过_resolve_connection()把连接数据解析成一个内部字典源码。默认值与文档中的取值完全吻合conn_data: dict[str, Any] { master: yarn, queue: None, # yarn queue deploy_mode: None, spark_binary: self.spark_binary or DEFAULT_SPARK_BINARY, namespace: None, ... rest_scheme: http, rest_port: 6066, }解析逻辑的几个关键点Host 语义当连接由 URI 创建时schemespark://、k8s://等存放在conn_type字段conn.host只含主机名由 UI 创建时conn_type通常是sparkconn.host可能包含完整 master URL例如k8s://https://host。hook 会按这两种情况分别还原出完整的 master URL并在有 Port 时追加:PORTextra 覆盖规则queueYARN 队列、deploy-mode、spark-binary、namespace、rest-scheme、rest-port都从连接的extraJSON 读取spark-binary若不在允许列表内会直接抛出ValueError历史字段清理spark-homeextra 已不再被允许——源码中显式抛出异常提示必须保证spark-submit/spark2-submit/spark3-submit之一在 PATH 上并通过spark-binary指定spark.kubernetes.namespace兜底若conf中显式设置了spark.kubernetes.namespace它会覆盖连接里的namespace源码 L465-L466。字段如何进入 spark-submit 命令解析出的连接数据最终在 _build_spark_common_args() 中拼接为spark-submit参数--master来自 Host/Port、每个conf键值对转为--conf keyvalue、--queueYARN 队列、--deploy-modeclient/cluster、--keytab/--principal、--name、--class、--jars、--files、--py-files等。也就是说连接中的Host 决定--master取值Deploy mode 决定--deploy-mode取值Spark binary 决定被调用的可执行文件——这也解释了为什么spark-binary必须位于 worker 的 PATH 上hook 的类注释明确写着requires spark-submit binary in the PATH。通过环境变量配置连接URI 语法当使用环境变量注入连接时应使用 URI 语法并可以直接提供标准的 Spark master URI——master URL 会被正确解析不需要spark://spark://...这类重复前缀同时所有 URI 组件需要 URL 编码。文档给出的标准示例export AIRFLOW_CONN_SPARK_DEFAULTspark://mysparkcluster.com:80?deploy-modeclusterspark_binarycommandnamespacekubenamespace要点环境变量名AIRFLOW_CONN_CONN_ID中的CONN_ID使用大写?之后的 query 参数对应连接 extra 字段如deploy-mode、namespace等值需 URL 编码示例中 namespace 的空格写作该连接被SparkSubmitHook解析后master 为spark://mysparkcluster.com:80属于Spark standalone集群管理器且deploy_modecluster会触发 driver 状态跟踪逻辑见下文。安全注意事项文档中有一条明确的安全警告能够配置 Host 设置的用户应被视为可信用户。因为 Host 指向哪里worker 就会去哪里建立通信——将连接指向恶意服务器可能带来严重的安全漏洞包括远程代码执行RCE风险。在开放 Connections 配置权限时务必把这一风险纳入权限设计例如仅管理员可修改连接。源码纵深连接字段如何驱动三种集群模式连接配置不是孤立的静态数据它直接决定SparkSubmitOperator在运行时选择哪种跟踪后端。在 operators/spark_submit.py 中_backend属性根据 hook 解析出的连接信息选择后端if self._hook._is_yarn_cluster_mode: backend _YarnSparkSubmitBackend(self, self._hook) elif self._hook._is_kubernetes: backend _KubernetesSparkSubmitBackend(self, self._hook) else: backend _StandaloneSparkSubmitBackend(self, self._hook)判断依据正是连接中的 master 与 deploy_mode_is_yarn yarn in master、_is_kubernetes k8s in master、_should_track_driver_status仅当 master 含spark://且deploy_mode cluster时成立hook 源码 L338-L347。这解释了连接字段与实际行为的映射关系1. Standalone cluster 模式REST scheme / REST port 的真正用途连接里的REST scheme/REST port字段服务于 standalone 模式下的 driver 状态轮询。_StandaloneSparkSubmitBackend.get_job_status()会构造GET {scheme}://{host}:{rest_port}/v1/submissions/status/{driverId}请求源码 L214-L231并且支持 HA 场景——master URL 形如spark://m1:7077,m2:7077时按逗号拆分后逐个尝试。注意源码中的注释master URL 中的端口如 7077是 RPC 端口不是 REST API 端口REST 端口必须通过连接 extra 的rest-port指定。这正是文档把这两个字段标注为仅 Spark standalone cluster 模式的原因同时也是 durable 执行重试时重连已运行的 driver的前置条件。2. YARN cluster 模式queue、RM 地址与 REST 轮询queue字段最终变成spark-submit --queue参数见_build_spark_common_args当 operator 设置yarn_track_via_rm_apiTrue时hook 要求连接 extra 中必须配置yarn_resourcemanager_webapp_addressprovider.yaml 中的字段定义描述镜像 Hadoop 的yarn.resourcemanager.webapp.address。_get_yarn_rm_base_url()会读取该值裸host:port会自动补http://前缀hook 源码 L983-L1009提交后 hook 从日志正则提取application[0-9_]得到 YARN 应用 ID随后轮询GET /ws/v1/cluster/apps/{appId}直到状态进入SUCCEEDED/FAILED/KILLEDon_kill时则通过PUT /ws/v1/cluster/apps/{appId}/state发送KILLED源码 L1064-L1086。对应的连接与 operator 用法来自 operators 文档 中的 YARN 小节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_idspark_pi, conn_idspark_yarn_rm, application/path/to/spark-examples.jar, java_classorg.apache.spark.examples.SparkPi, deploy_modecluster, yarn_track_via_rm_apiTrue, )对于 Kerberized 集群安装requests-kerberos后只要连接同时配置了keytab与principalhook 会自动使用HTTPKerberosAuth()对 RM REST 请求鉴权源码_resolved_yarn_rm_auth仅当 RM 需要自定义requests鉴权对象时才显式传yarn_rm_auth。3. Kubernetes cluster 模式namespace 字段与 driver pod 跟踪namespace字段在提交时转为--conf spark.kubernetes.namespace...源码 L588-L592设置track_driver_via_k8s_apiTrue要求 master 为k8s://...且deploy_modecluster并强制存在 namespace后hook 会释放长期阻塞的spark-submitJVM改用 Kubernetes API 轮询 driver pod 的status.phase直至终态源码_poll_k8s_driver_via_api并在任务结束或on_kill时清理 driver pod该标志同时是 Kubernetes 上 durable 执行重试时重连既有 driver pod 而非重新提交的前置条件对应 operators 文档 中的示例run_spark SparkSubmitOperator( task_idrun_spark, applicationlocal:///opt/spark/examples/jars/spark-examples.jar, conn_idspark_k8s, deploy_modecluster, track_driver_via_k8s_apiTrue, durableTrue, )各集群管理器配置前提速查将上文归纳为一张配置决策表与 operators 文档中的 durable 表格 一致集群管理器需要在 Spark 连接中配置的关键字段Spark standalonecluster 模式Host spark://HOST:PORTdeploy-modecluster如需重试重连/HTTPS配置rest-scheme与rest-portYARNcluster 模式Host yarn或yarn加端口deploy-modecluster可选queueRM API 跟踪需yarn_resourcemanager_webapp_addressKubernetescluster 模式Host k8s://...deploy-modeclusternamespaceK8s API 跟踪需track_driver_via_k8s_apiTrue本地/客户端模式Host localspark-binary指向 worker PATH 中存在的提交命令补充两点运行前提从源码结构看Kubernetes 跟踪依赖apache-airflow-providers-cncf-kubernetes提供的kube_client未安装时 hook 在 master 含k8s时会直接抛出RuntimeErrorhook 源码 L314-L318durable 崩溃恢复依赖 Airflow 3.3 的 task state store低于 3.3 时durable仅产生警告、行为与旧版本一致。小结spark连接类型的核心字段是Host deploy-mode spark-binary其余字段queue、namespace、rest-scheme/rest-port、RM 地址按集群管理器选择性配置连接字段并非填了就行它们直接决定spark-submit的--master/--queue/--deploy-mode参数以及 operator 选择 Standalone / YARN / K8s 哪种状态跟踪后端配置连接前建议先核对源码约束provider.yaml、hook、operator并在集群启用 Kerberos、TLS、非默认 REST 端口时同步调整keytab/principal、rest-scheme等字段开放连接配置权限时牢记文档的 RCE 安全警告Host 指向恶意服务器会让 worker 与之通信。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表