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

如何在Python脚本中获取后台运行Spark应用的Spark Session对象

实现Python脚本复用同一个常驻Spark Session的方案

核心思路是将Spark应用的启动、销毁逻辑和客户端业务脚本解耦,把Spark应用作为常驻服务在后台运行,客户端通过RPC/REST接口和该服务交互,不需要每次运行脚本都重新启动Spark进程。

方案一:Spark Thrift JDBC/ODBC Server(适合SQL为主的场景)

这是成本最低的实现方案,无需额外部署第三方组件:

  • 第一步在集群后台启动Spark Thrift Server,启动后对应的Spark应用会持续运行,命令示例:
    $SPARK_HOME/sbin/start-thriftserver.sh --master yarn --deploy-mode client --name persistent-spark-app --num-executors 2 --executor-memory 4g
    参数可根据自身资源需求调整,默认服务端口为10000。
  • 第二步Python客户端无需初始化SparkSession,直接通过JDBC客户端连接Thrift Server即可复用后台Spark Session的资源,示例代码:
from pyhive import hive

# 连接后台常驻的Spark Thrift服务
conn = hive.Connection(
    host="你的Thrift Server部署节点IP",
    port=10000,
    username="集群用户名"
)
cursor = conn.cursor()

# 直接执行SQL,所有计算都在后台常驻的Spark应用中运行
cursor.execute("select * from test_db.test_table limit 10")
result = cursor.fetchall()
print(result)

# 关闭客户端连接不会停止后台的Spark应用
cursor.close()
conn.close()

该方案的优点是客户端不需要安装完整Spark环境,改造现有SQL类业务脚本的成本极低。

方案二:Livy REST Server(适合需要提交自定义PySpark代码的场景)

如果你的业务除了SQL之外还要执行复杂的DataFrame操作、自定义UDF等逻辑,可以用Livy服务实现:

  • 第一步部署Livy服务后,提前创建一个常驻的PySpark Session,命令示例:
    curl -X POST -H "Content-Type: application/json" -d '{"kind": "pyspark", "name": "persistent-session", "executorMemory": "4g", "numExecutors": 2}' http://你的Livy部署节点IP:8998/sessions
    请求返回的id就是你的常驻Session ID。
  • 第二步Python客户端调用Livy的REST接口,给指定的常驻Session提交代码即可,示例代码:
import requests
import time

LIVY_SERVICE_URL = "http://你的Livy部署节点IP:8998"
# 替换为你提前创建的常驻Session ID
PERSISTENT_SESSION_ID = 1

# 提交PySpark代码到常驻Session执行
req_data = {
    "code": """
# 这里可以写任意PySpark代码,spark变量就是后台常驻Session的SparkSession对象
df = spark.sql("select * from test_db.test_table limit 10")
df.collect()
"""
}
resp = requests.post(f"{LIVY_SERVICE_URL}/sessions/{PERSISTENT_SESSION_ID}/statements", json=req_data)
statement_id = resp.json()["id"]

# 轮询获取执行结果
while True:
    resp = requests.get(f"{LIVY_SERVICE_URL}/sessions/{PERSISTENT_SESSION_ID}/statements/{statement_id}")
    state = resp.json()["state"]
    if state in ("available", "error", "cancelled"):
        break
    time.sleep(1)

print(resp.json()["output"]["data"])

该方案的优点是灵活性极高,支持提交任意PySpark逻辑,适合复杂计算场景。

注意事项

  • 两种方案都不需要在客户端Python脚本中初始化本地SparkSession,所有计算逻辑都在后台常驻的Spark应用中执行
  • 要给常驻Spark应用配置合理的资源上限,避免长期占用过多集群资源
  • 按需配置Session超时时间避免空闲被回收:Thrift Server可设置spark.session.timeout参数,Livy可设置livy.server.session.timeout参数
  • 如果需要跨脚本共享DataFrame/RDD,提前将数据注册成临时视图即可,后续脚本直接通过视图名访问。

内容的提问来源于stack exchange,提问作者zwithouta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 22:30:05