跨服务器从Airflow提交Spark Job的配置方法问询
Airflow跨服务器通过SparkSubmitOperator提交YARN任务配置步骤
1. 在Airflow服务器部署兼容版本的Spark二进制包
- 下载与Hadoop集群版本匹配的Spark安装包(例如Hadoop 3.3.x对应Spark 3.4.x),解压到Airflow服务器的固定目录,比如
/opt/spark - 配置环境变量:将
SPARK_HOME=/opt/spark添加到Airflow运行用户的~/.bashrc,或者在Airflow的配置文件airflow.cfg中通过env_vars字段指定,确保Airflow进程能读取到该变量
2. 同步Hadoop集群配置文件到Airflow服务器
- 从Hadoop集群的任意节点,复制以下配置文件到Airflow服务器的
$SPARK_HOME/conf目录:yarn-site.xml:包含YARN ResourceManager的地址、端口等核心配置core-site.xml:包含HDFS NameNode地址及文件系统配置hdfs-site.xml:包含HDFS的副本数、存储路径等配置
- 检查配置文件中的地址(如ResourceManager的
yarn.resourcemanager.address、NameNode的fs.defaultFS)是否为Airflow服务器可访问的内网/公网地址,同时确保集群对应端口(如8032、9000、8088)已开放防火墙规则
3. 配置Airflow的SparkSubmitOperator作业
方式1:通过Airflow连接配置(推荐)
- 在Airflow UI中进入「Admin」→「Connections」,新建Spark类型的连接:
- Host:填写YARN ResourceManager的地址(如
rm-cluster.example.com) - Port:填写ResourceManager的提交端口(默认8032)
- Extra:添加Spark路径及可选配置,例如:
{"spark.home": "/opt/spark", "spark.yarn.jars": "hdfs://namenode:9000/spark/jars/*.jar"}
- Host:填写YARN ResourceManager的地址(如
- DAG中调用Operator时指定
conn_id:from airflow import DAG from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1) } with DAG('spark_yarn_dag', default_args=default_args, schedule_interval='@daily') as dag: submit_spark_job = SparkSubmitOperator( task_id='run_spark_on_yarn', application='hdfs://namenode:9000/jobs/your_spark_job.jar', master='yarn', deploy_mode='cluster', conn_id='spark_yarn_conn', class='com.your.company.SparkMainClass', application_args=['--input-path', 'hdfs://namenode:9000/input', '--output-path', 'hdfs://namenode:9000/output'] )
方式2:直接在Operator中指定参数
- 无需创建Airflow连接,直接在Operator中配置所有必要参数:
submit_spark_job = SparkSubmitOperator( task_id='run_spark_on_yarn', application='/local/path/to/your_spark_job.jar', master='yarn', deploy_mode='client', spark_home='/opt/spark', conf={ 'spark.yarn.resourcemanager.address': 'rm-cluster.example.com:8032', 'spark.hadoop.fs.defaultFS': 'hdfs://namenode:9000', 'spark.driver.memory': '2g', 'spark.executor.instances': '3', 'spark.executor.memory': '4g' }, class='com.your.company.SparkMainClass' )
4. 验证配置与权限
- 在Airflow服务器上,切换到Airflow运行用户(如
airflow),执行测试命令验证Spark提交能力:$SPARK_HOME/bin/spark-submit --master yarn --deploy-mode cluster --class com.your.company.SparkMainClass hdfs://namenode:9000/jobs/your_spark_job.jar - 确保Airflow用户拥有以下权限:
- 读取HDFS上的作业jar包及输入数据的权限
- YARN集群的任务提交权限(若集群开启Kerberos认证,需在Airflow服务器配置Kerberos客户端,复制
krb5.conf,并在Operator中指定kerberos_principal和kerberos_keytab参数)
内容的提问来源于stack exchange,提问作者Amarjeet Kushwaha
相关产品推荐
相关产品推荐

