通过Apache Livy提交PySpark批处理作业时配置不生效如何解决
问题根因说明
你怀疑的「新建SparkSession导致配置失效」的问题不成立,SparkSession.builder.getOrCreate()的逻辑是如果当前进程内已存在SparkSession实例就直接复用,不会新建。Livy在启动Spark批处理作业进程时,会提前初始化带请求体conf配置的上下文,正常调用getOrCreate()会自动继承这些配置,出现配置失效通常是以下几个原因:
- Livy服务端配置了配置覆盖策略,资源类配置被加入了用户自定义黑名单,会用集群默认值覆盖你提交的参数
job.py代码中在获取SparkSession后,调用了spark.conf.set()或者Builder的config()方法覆盖了传入的资源参数- 你使用的Livy版本存在已知Bug,批处理模式下
conf参数不会透传给Spark Driver进程
排查步骤
- 进入提交作业的Spark UI,打开
Environment标签页,搜索你设置的配置项(例如spark.executor.memory):- 如果完全没有你提交的配置,说明Livy没有透传参数到Spark进程
- 如果配置存在但值和你提交的不一致,说明配置被其他规则覆盖
- 检查Livy服务端
livy.conf配置文件,确认是否配置了livy.server.conf.blacklist将资源类配置加入了禁止用户自定义的黑名单,或者livy.spark.master强制指定了不限制资源的运行模式 - 检查
job.py全量代码,确认没有其他逻辑修改Spark运行时配置
解决方案
方案1:显式传入配置文件兜底
你可以将需要的配置写入spark-defaults.conf,和作业文件一起上传,在Livy请求体中新增files字段传入配置文件:
REQUEST_BODY = { 'file': '/spark/batch/job.py', 'files': ['/path/to/your/spark-defaults.conf'], 'conf': { 'spark.sql.legacy.allowCreatingManagedTableUsingNonemptyLocation': 'true', 'spark.driver.cores': 1, 'spark.driver.memory': '12g', 'spark.executor.cores': 1, 'spark.executor.memory': '8g', 'spark.dynamicAllocation.maxExecutors': 4, }, }
方案2:作业代码显式加载环境变量配置
Livy会把所有传入的conf配置写入Spark进程的系统环境变量中,你可以在创建SparkSession时显式读取加载:
from pyspark.sql import SparkSession import os builder = SparkSession.builder # 读取Livy透传的Spark配置环境变量 for env_key, env_val in os.environ.items(): if env_key.startswith("SPARK_"): conf_key = env_key.replace("SPARK_", "spark.").replace("_", ".") builder = builder.config(conf_key, env_val) spark = builder.getOrCreate()
方案3:调整Livy服务端配置
如果是Livy服务端拦截了你的配置,修改livy.conf移除资源类配置的黑名单限制,或者调整livy.spark.*开头的默认配置匹配你的资源需求。
内容的提问来源于stack exchange,提问作者Korntewin Boonchuay
相关产品推荐
相关产品推荐

