Airflow跨机器提交PySpark任务至独立模式Spark集群:可行性与实现
问题解答
1. 该部署方案是否可行?
可行,但需要避开Spark独立模式对Python任务的限制:Spark独立集群不支持Python应用的cluster部署模式,但可以通过调整提交方式,让client模式下的Driver运行在Spark Master节点(而非Airflow所在的机器A),从而满足你“由Spark Master作为Driver”的需求。
2. 具体实现步骤
核心思路
放弃在Airflow容器内直接执行spark-submit(否则Driver会跑在Airflow节点),改为通过SSH连接到机器B的Spark Master节点/容器,在那里执行提交命令,让Driver运行在Spark Master上。同时在Spark Master环境中配置MinIO访问依赖,解决脚本读取问题。
步骤1:配置Spark Master环境,支持MinIO访问
在机器B的Spark Master容器中,添加访问MinIO所需的依赖包,并配置Spark参数:
- 下载与你的Spark/Hadoop版本匹配的
hadoop-aws和aws-java-sdk-bundlejar包,放入Spark安装目录下的jars文件夹。 - 在Spark的
conf/spark-defaults.conf中添加MinIO配置:spark.hadoop.fs.s3a.endpoint http://minio:9000 # MinIO的访问地址,若在同一Docker网络可直接用容器名 spark.hadoop.fs.s3a.access.key YOUR_MINIO_ACCESS_KEY spark.hadoop.fs.s3a.secret.key YOUR_MINIO_SECRET_KEY spark.hadoop.fs.s3a.path.style.access true spark.hadoop.fs.s3a.impl org.apache.hadoop.fs.s3a.S3AFileSystem
步骤2:配置Airflow与机器B的SSH连通
- 在Airflow的Web UI中添加SSH连接(Admin -> Connections),配置机器B的IP/主机名、SSH端口、用户名、密码或密钥,确保Airflow所在的机器A能通过SSH访问机器B的Spark Master节点/容器。
- 若机器B的Spark Master是Docker容器,需确保容器暴露了SSH端口,或者Airflow容器与Spark集群处于同一Docker网络,可直接通过容器名访问。
步骤3:编写Airflow任务,通过SSH提交Spark任务
使用Airflow的SSHOperator远程执行spark-submit命令,示例代码如下:
from airflow.providers.ssh.operators.ssh import SSHOperator from airflow.models import DAG from datetime import datetime default_args = { 'owner': 'airflow', 'retries': 1 } with DAG( dag_id='remote_spark_submit', default_args=default_args, start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: submit_spark_task = SSHOperator( task_id='submit_spark_to_master', ssh_conn_id='spark_master_b', # 对应Airflow中配置的SSH连接ID command=''' spark-submit \ --master spark://spark-master:7077 \ --deploy-mode client \ s3a://your-minio-bucket/path/to/your_spark_script.py ''' )
步骤4:验证与调试
- 先在机器B的Spark Master节点/容器中手动执行上述
spark-submit命令,确认能正常读取MinIO中的脚本并运行任务,Driver运行在Spark Master上。 - 再触发Airflow任务,检查任务日志,确认SSH连接正常,任务提交成功。
对之前错误的说明
- 直接在Airflow容器执行
spark-submit时,client模式下Driver会运行在Airflow容器内,因此需要Airflow容器配置hadoop-aws依赖,这会增加Airflow容器的复杂度,不是最优方案。 - 切换到
cluster模式报错是因为Spark独立集群本身不支持Python应用的cluster部署,这是官方限制,无法绕过。
内容的提问来源于stack exchange,提问作者Rizky Prasetya
相关产品推荐
相关产品推荐

