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

如何在Spark单元测试中启用多Catalog?

问题描述

我正在开发将在开启Unity Catalog的Databricks上运行的PySpark应用,所有读取的表名称格式为catalog.database.table。我希望在单元测试中也能启用多Catalog,虽已知可通过配置表前缀的变通方案(如本地单元测试读取database.table,Databricks上读取catalog.database.table),但开源Spark 3.x理应支持多Catalog。

我使用pytest并在conftest.py中定义了如下fixture:

@pytest.fixture(scope="session")
def spark():
    builder = (
        SparkSession.builder.appName("MyApp")
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
        .config(
            "spark.sql.catalog.spark_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog",
        ).config(
            "spark.sql.catalog.my_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog",
        )
    )

    spark = configure_spark_with_delta_pip(builder).getOrCreate()
    yield spark
    spark.stop()

但当尝试使用my_catalog时:

@pytest.fixture
def some_fixture(spark):
    spark.catalog.setCurrentCatalog("my_catalog")
    spark.sql("CREATE DATABASE IF NOT EXISTS db")

出现SparkException错误:[INTERNAL_ERROR] The Spark SQL phase analysis failed with an internal error. You hit a bug in Spark or the Spark plugins you use. Please, report this bug to the corresponding communities or vendors, and provide the full stack trace.

请问如何在Spark中使用多Catalog?

解决方案

开源Spark 3.x的多Catalog配置需要注意几个关键细节,针对你的问题,调整如下:

1. 为自定义Catalog指定根路径

DeltaCatalog作为自定义Catalog时,必须配置对应的根存储路径,否则无法正常创建数据库/表。在SparkSession builder中添加路径配置:

.config("spark.sql.catalog.my_catalog.path", "/tmp/my_catalog")  # 可指定本地临时路径或其他存储路径

2. 完整的SparkSession配置示例

修改后的conftest.py fixture如下:

@pytest.fixture(scope="session")
def spark():
    builder = (
        SparkSession.builder.appName("MyApp")
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
        .config(
            "spark.sql.catalog.spark_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog",
        )
        .config(
            "spark.sql.catalog.my_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog",
        )
        .config("spark.sql.catalog.my_catalog.path", "/tmp/my_catalog")  # 新增路径配置
    )

    spark = configure_spark_with_delta_pip(builder).getOrCreate()
    yield spark
    spark.stop()

3. 使用自定义Catalog的正确方式

创建数据库时,可显式指定Catalog避免歧义,或确保当前Catalog已正确切换:

@pytest.fixture
def some_fixture(spark):
    spark.catalog.setCurrentCatalog("my_catalog")
    # 方式1:显式指定Catalog创建数据库
    spark.sql("CREATE DATABASE IF NOT EXISTS my_catalog.db")
    # 方式2:切换当前Catalog后直接创建,默认使用当前Catalog
    spark.sql("CREATE DATABASE IF NOT EXISTS db")

4. 额外注意事项

  • 确保Delta Lake版本与Spark版本兼容,Delta 2.x+对Spark 3.2+的多Catalog支持更完善
  • 若使用本地文件系统,需保证指定的路径有读写权限
  • 单元测试完成后可清理临时路径下的文件,避免占用空间

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 21:06:09