You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

跨服务器从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"}
  • 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 02:10:34