Docker部署Airflow无法调用Windows本地Spark,报pyspark模块缺失
解决Docker部署的Airflow无法调用PySpark的问题
你用Docker部署的Airflow,Windows本地的Spark能正常通过spark-submit运行脚本,但把逻辑写成Airflow DAG时,出现ModuleNotFoundError: No module named 'pyspark'——这是因为Airflow容器和Windows本地环境完全隔离,容器里没装PySpark依赖,也没法直接访问本地的Spark环境。下面给几种实用解决办法:
方法一:给Airflow容器安装PySpark依赖
临时测试方案(重启容器失效)
- 先找到Airflow worker容器名:运行
docker ps,找到名称类似airflow-worker的容器 - 进入容器:
docker exec -it <容器名> bash - 安装PySpark:
pip install pyspark
持久化方案(自定义镜像)
- 创建一个
Dockerfile,基于你使用的官方Airflow镜像:
FROM apache/airflow:2.x.x # 替换成你实际用的Airflow版本号 RUN pip install pyspark
- 修改
docker-compose.yml,把对应服务(airflow-worker、airflow-scheduler、airflow-webserver)的image字段替换成你要构建的镜像名,或者添加build字段指向Dockerfile路径:
services: airflow-worker: build: . # 假设Dockerfile和docker-compose.yml在同一目录 # 其他原有配置...
- 重新部署Airflow:
docker-compose up -d --build
方法二:用SparkSubmitOperator替代PythonOperator(推荐)
这种方式不需要在Airflow容器里装PySpark,直接调用Windows本地的Spark集群提交任务,更符合Spark任务的运行规范。
- 把PySpark逻辑单独抽成脚本文件(比如
spark_job.py),放在Airflow的dags目录下:
from pyspark.sql import SparkSession from createconnection import connection from datetime import datetime import os def main(): url, uname, pwd = connection() spark = SparkSession.builder.appName("MyApp").getOrCreate() try: jdbcDF = spark.read \ .format("jdbc") \ .option("url", url) \ .option("driver", "com.microsoft.sqlserver.jdbc.SQLServerDriver") \ .option("query", '''SELECT TOP 100 t1.TBNAME as Experiment, t2.TBYEAR as EXP_YEAR, t4.TBFLDNAME as TRAITNAME FROM dbo.EXP01 as t1 INNER JOIN dbo.EXP05 as t2 ON t1.TBID=t2.TBEXPTID INNER JOIN dbo.EXP02 as t3 ON t1.TBID=t3.TBEXPTID INNER JOIN dbo.TRD01 as t4 ON t3.TBTRAITID=t4.TBID ''') \ .option("user", uname) \ .option("password", pwd) \ .load() # 注意:输出路径要改成容器和Windows都能访问的共享目录(比如Docker挂载的volume) output_file = f"/opt/airflow/output/output_{datetime.now().strftime('%Y-%m-%d_%H-%M-%S')}.csv" if not os.path.exists(output_file): with open(output_file, 'w') as f: f.write('Experiment,EXP_YEAR,TRAITNAME\n') jdbcDF.write.csv(output_file, header=True, mode='append') spark.stop() except ValueError as error: print("Connector write failed", error) if __name__ == "__main__": main()
- 修改DAG代码,使用
SparkSubmitOperator:
from datetime import datetime, timedelta from airflow import DAG from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2023, 5, 9), 'retries': 0, 'retry_delay': timedelta(minutes=5) } dag = DAG( 'pyspark_dag', default_args=default_args, description='A simple DAG to run a PySpark script every 2 minutes', schedule_interval=timedelta(minutes=2) ) run_script_task = SparkSubmitOperator( task_id='run_pyspark_script', application='/opt/airflow/dags/spark_job.py', # 容器内的脚本路径 conn_id='spark_default', # 需先在Airflow UI配置Spark连接 dag=dag )
- 在Airflow UI配置Spark连接:
- 进入Airflow UI → Admin → Connections
- 新建连接,Conn Id设为
spark_default,Conn Type选Spark - Host填Windows主机的IP(容器要访问本地Spark),Port填Spark默认端口7077
- 确保Windows防火墙允许容器访问7077端口,且Spark配置允许远程连接(修改
spark-defaults.conf,设置spark.driver.host为Windows主机IP,spark.driver.bindAddress为0.0.0.0)
方法三:挂载Windows的Spark目录到容器
把Windows本地的Spark安装目录挂载到Airflow容器,再设置环境变量让容器识别:
- 修改
docker-compose.yml,添加volume挂载和环境变量:
services: airflow-worker: volumes: - ./dags:/opt/airflow/dags - ./logs:/opt/airflow/logs - ./plugins:/opt/airflow/plugins - //c/Spark:/opt/spark # 替换成你的Windows Spark路径,注意Docker挂载Windows路径的格式(//盘符/目录) environment: - SPARK_HOME=/opt/spark - PYTHONPATH=/opt/spark/python:/opt/spark/python/lib/py4j-0.10.9.5-src.zip:$PYTHONPATH # 替换成你Spark对应版本的py4j文件名
- 重启Airflow容器:
docker-compose up -d
内容的提问来源于stack exchange,提问作者Souradip Roy
相关产品推荐
相关产品推荐

