如何通过Apache Livy设置Spark配置属性及传递SparkSession参数?
我来帮你梳理下怎么通过Apache Livy设置Spark配置,以及如何给SparkSession传递参数,都是实际使用中常用的方法:
一、提交作业时设置Spark配置属性
你有三种常用方式来配置Spark参数:
1. 通过REST API提交时指定conf字段
这是最灵活的方式,每次提交作业都可以自定义配置。用curl提交的话,在请求的JSON体里加入conf对象,里面放你需要的Spark配置:
curl -X POST -H "Content-Type: application/json" --data '{ "file": "/path/to/your/spark-job.py", "conf": { "spark.executor.memory": "2g", "spark.driver.memory": "1g", "spark.sql.shuffle.partitions": "200", "spark.app.name": "Livy-Custom-Job" } }' http://<livy-server-ip>:8998/batches
这里的conf里的键就是Spark原生的配置项,和你在spark-defaults.conf里写的完全一致。
2. 编程方式提交(比如Python)
如果用代码调用Livy的API提交作业,同样在请求的JSON payload里加入conf字段,以Python的requests库为例:
import requests livy_server = "http://<livy-server-ip>:8998/batches" submit_payload = { "file": "/path/to/your/spark-job.py", "conf": { "spark.executor.cores": "2", "spark.submit.deployMode": "cluster" } } response = requests.post(livy_server, json=submit_payload) print("提交结果:", response.json())
3. 设置Livy全局默认配置
如果你希望所有通过Livy提交的作业都使用统一的默认配置,可以修改Livy的livy.conf文件,添加对应的Spark配置前缀为livy.spark.:
# livy.conf 示例配置 livy.spark.master = yarn livy.spark.submit.deployMode = cluster livy.spark.executor.memory = 1g
注意:提交作业时指定的conf会覆盖全局默认配置,优先级更高。
二、向SparkSession传递参数(编程方式)
这里分两种场景,根据你的需求选择:
1. 让SparkSession自动加载Livy传递的配置
当你通过Livy的conf字段提交配置后,作业里创建的SparkSession会自动读取这些配置,不需要额外代码。比如你的作业代码可以这样写:
from pyspark.sql import SparkSession # 直接创建SparkSession,会自动加载Livy传递的所有配置 spark = SparkSession.builder.getOrCreate() # 验证配置是否生效 print("当前应用名称:", spark.conf.get("spark.app.name")) print("Executor内存:", spark.conf.get("spark.executor.memory"))
2. 在作业代码中显式设置SparkSession参数
如果需要在作业内部动态配置SparkSession,可以在SparkSession.builder里用.config()方法指定,或者在Session创建后用spark.conf.set()修改(部分运行时可修改的配置):
from pyspark.sql import SparkSession # 初始化时指定配置 spark = SparkSession.builder \ .appName("Custom-Session-Job") \ .config("spark.sql.warehouse.dir", "/user/hive/warehouse") \ .config("spark.driver.maxResultSize", "1g") \ .getOrCreate() # 运行时修改支持动态调整的配置 spark.conf.set("spark.sql.adaptive.enabled", "true")
注意:有些Spark配置是静态的,只能在Session初始化时设置,运行时修改会报错,比如spark.executor.memory这类资源相关的配置。
三、完整示例:提交+验证配置
举个实际的例子,假设你的测试作业是test-job.py,内容如下:
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() # 打印Livy传递的配置 print("App Name:", spark.conf.get("spark.app.name")) print("Shuffle Partitions:", spark.conf.get("spark.sql.shuffle.partitions")) # 简单的Spark操作 df = spark.range(100) df.show(10) spark.stop()
然后用curl提交这个作业并传递配置:
curl -X POST -H "Content-Type: application/json" --data '{ "file": "/tmp/test-job.py", "conf": { "spark.app.name": "Test-Livy-Job", "spark.sql.shuffle.partitions": "100", "spark.executor.memory": "1.5g" } }' http://localhost:8998/batches
提交后,你可以通过Livy的API查看作业日志,就能看到打印出的配置参数,验证是否生效。
内容的提问来源于stack exchange,提问作者Markus

