当本地训练的python环境与 Spark 集群其他节点环境不一致时,核心解决方案是:将本地 Python 虚拟环境打包为压缩包,通过 Spark 的–archives参数分发到所有 Executor 节点,强制所有节点使用统一的环境(Python 解释器 + 依赖包)
以下是具体实现步骤,支持conda和virtualenv两种虚拟环境(覆盖绝大多数本地训练场景),且无需 root 权限、不影响集群原有环境。
一、核心原理
环境打包:将本地虚拟环境(含 Python 解释器、所有依赖包如scikit-learn、pyarrow、自定义类脚本)压缩为归档文件(.zip或.tar.gz)。
Spark 分发:通过spark-submit的–archives参数,让 Spark 自动将压缩包分发到所有 Executor 节点的工作目录。
指定解释器:通过–conf spark.pyspark.python和–conf spark.pyspark.driver.python,强制 Driver 和 Executor 使用打包环境中的 Python 解释器,而非集群默认 Python。
1、激活本地 conda 环境,确认依赖完整,打包 conda 环境为压缩包

# 激活本地conda环境
conda activate sklearn-spark-env

# 打包本地conda环境(替换为你的环境路径,可通过conda env list查看)
conda pack -n sklearn-spark-env -o sklearn-spark-env.tar.gz --ignore-missing-files

-n sklearn-spark-env:指定要打包的 conda 环境名;
-o sklearn-spark-env.tar.gz:输出压缩包名称;
–ignore-missing-files:忽略部分系统依赖(集群通常已包含),减小包体积。
2、通过 Spark 提交任务,分发环境并运行,分两种情况:
2.1、client模式

#!/bin/bash

# 配置环境变量
export JAVA_HOME=/opt/hadoop-cluster/java  # java路径 
export SPARK_HOME=/opt/hadoop-cluster/spark  # spark路径
export PATH=$JAVA_HOME/bin:$SPARK_HOME/bin:$PATH  
#  client模式需要配置PYSPARK_DRIVER_PYTHON 即客户端本地python路径
export PYSPARK_DRIVER_PYTHON=/home/ll/miniconda3/envs/pyspark_env/bin/python  # client模式才需要配置 cluster模式不需要配置
export PYSPARK_PYTHON=./env/bin/python  # 解压后的python路径


# 任务配置参数
SPARK_MASTER="yarn"  # yarn
DEPLOY_MODE="client"  # client or cluster
APP_NAME="spark_py_job_client"
NUM_EXECUTORS=2
EXECUTOR_MEM="1g"  # 根据任务情况 一般 3~8g 太大会影响队列其他人的任务资源
DRIVER_MEM="1g"  # 默认1g 基本不动
EXECUTOR_CORES=2
QUEUE="ll"  # 队列

# 文件配置参数
PY_SCRIPT="/home/ll/py_code/spark/00_spark_test.py"  # 执行的python文件路径
VENV_ARCHIVE="/home/ll/py_code/spark/pyspark_env.tar.gz" # 打包好的tar.gz文件路径

# 执行命令
spark-submit \
--master $SPARK_MASTER \
--deploy-mode $DEPLOY_MODE \
--name $APP_NAME \
--num-executors $NUM_EXECUTORS \
--executor-memory $EXECUTOR_MEM \
--driver-memory $DRIVER_MEM \
--executor-cores $EXECUTOR_CORES \
--queue $QUEUE \  # 队列名称
--archives $VENV_ARCHIVE#env \  # #env要加上 表示加压后的名称
--conf spark.default.parallelism=8 \  # spark.default.parallelism 一般设置为NUM_EXECUTORS*EXECUTOR_CORES的2~3倍
--conf spark.driver.log.level=WARN \  # 调整日志级别 减少日志输出
--conf spark.executor.log.level=WARN \
$PY_SCRIPT

2.1、cluster模式

#!/bin/bash

# 配置环境变量
export JAVA_HOME=/opt/hadoop-cluster/java
export SPARK_HOME=/opt/hadoop-cluster/spark
export PATH=$JAVA_HOME/bin:$SPARK_HOME/bin:$PATH  
#  cluster模式不需要配置PYSPARK_DRIVER_PYTHON
export PYSPARK_PYTHON=./env/bin/python  # 解压后的python路径

# 任务配置参数
SPARK_MASTER="yarn"
DEPLOY_MODE="cluster"  # client or cluster
APP_NAME="spark_py_job_cluster"
NUM_EXECUTORS=2
EXECUTOR_MEM="1g"
DRIVER_MEM="1g"  
EXECUTOR_CORES=2
QUEUE="ll"

# 文件配置参数
PY_SCRIPT="/home/ll/py_code/spark/00_spark_test.py"
VENV_ARCHIVE="/home/ll/py_code/spark/pyspark_env.tar.gz"

# 执行命令
spark-submit \
--master $SPARK_MASTER \
--deploy-mode $DEPLOY_MODE \
--name $APP_NAME \
--num-executors $NUM_EXECUTORS \
--executor-memory $EXECUTOR_MEM \
--driver-memory $DRIVER_MEM \
--executor-cores $EXECUTOR_CORES \
--queue $QUEUE \
--archives $VENV_ARCHIVE#env \
--conf spark.default.parallelism=8 \
--conf spark.driver.log.level=ERROR \
--conf spark.executor.log.level=ERROR \
$PY_SCRIPT

更多推荐