You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过Apache Livy设置Spark配置属性及传递SparkSession参数?

通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 08:58:22