如何在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
相关产品推荐
相关产品推荐

