如何在导入Databricks Notebook的自定义Python包中引用Spark实例?
解决Databricks自定义包中Spark未定义的问题
问题核心
Databricks Notebook中自动初始化的spark对象属于全局作用域,自定义Python包的类无法直接访问,导致调用sf_query()时出现未定义错误。以下是两种可行的解决方法:
方案一:初始化类时传入SparkSession对象
直接将Notebook中的spark实例传入类的构造函数,保存为类的实例变量后使用。
修改后的自定义包类代码:
class snowflake_read(): def __init__(self, format, options, spark): self.options = options self.format = format self.spark = spark # 保存外部传入的SparkSession实例 def sf_query(self, query): df = self.spark.read.format(self.format).options(self.options).options("query", query).load() return df
Notebook中的调用代码:
import read_package sf_read = read_package.snowflake_read(<format>, <options>, spark) sf_read.sf_query('<query>')
方案二:在类内部获取SparkSession
通过SparkSession.builder.getOrCreate()自动获取Databricks已初始化的SparkSession,无需外部传参。
修改后的自定义包类代码:
from pyspark.sql import SparkSession class snowflake_read(): def __init__(self, format, options): self.options = options self.format = format self.spark = SparkSession.builder.getOrCreate() # 获取或创建SparkSession实例 def sf_query(self, query): df = self.spark.read.format(self.format).options(self.options).options("query", query).load() return df
Notebook中的调用代码保持不变:
import read_package sf_read = read_package.snowflake_read(<format>, <options>) sf_read.sf_query('<query>')
方案对比
- 方案一:依赖关系明确,适合需要指定特定SparkSession的场景,灵活性更高。
- 方案二:代码更简洁,无需额外传参,适配Databricks默认的Spark初始化环境。
内容的提问来源于stack exchange,提问作者docphilstone
相关产品推荐
相关产品推荐

