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

Azure Databricks Iceberg多数据源冲突及独占运行时诉求

问题描述

在使用Azure Databricks Runtime 14.3 LTS配置Hadoop Catalog测试Apache Iceberg存储实现时,执行建表操作遇到以下错误:

Multiple sources found for iceberg (com.databricks.sql.transaction.tahoe.uniform.sources.IcebergBrowseOnlyDataSource, org.apache.iceberg.spark.source.IcebergSource), please specify the fully qualified class name.

请问是否可以覆盖Databricks的Uniform Reader,让集群独占使用Apache Iceberg运行时?

集群配置信息

  • 策略:unrestricted,多节点
  • 访问模式:no isolation shared
  • 1个Worker:16GB内存、4核
  • 1个Driver:16GB内存、4核
  • Runtime版本:14.3.x-scala2.12
  • 启用Photon
  • 实例类型:Standard_D4ds_v5
  • 4 DBU/小时

Spark配置

  • spark.sql.catalog.spark_catalog.warehouse:/test_data/
  • spark.hadoop.fs.wasbs.impl:org.apache.hadoop.fs.azure.NativeAzureFileSystem
  • spark.sql.catalog.spark_catalog:org.apache.iceberg.spark.SparkCatalog
  • spark.hadoop.fs.abfss.impl:org.apache.hadoop.fs.azurebfs.AzureBlobFileSystem
  • spark.databricks.sql.iceberg.handle-tables:false
  • spark.sql.extensions:org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
  • spark.sql.catalog.spark_catalog.type:hadoop

集群已安装库

  • com.microsoft.azure:azure-storage:8.6.6
  • dbfs:/FileStore/jars/bf401b66_66f7_46c8_9964_b6ad07ad8099/data_lake_library-0.2.1.dev59+unknown-py3-none-any.whl(自定义团队构建的wheel包)
  • numpy==1.26.4
  • org.apache.hadoop:hadoop-azure:3.4.0
  • org.apache.hadoop:hadoop-client-runtime:3.4.0
  • org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.5.2
  • pandas==2.2.2
  • pyarrow==7.0.0

测试代码

import logging
from pyspark.sql import SparkSession
from data_lake_library.config.config import AzureKeyVaultParams
from data_lake_library.config.secrets_manager import AzureKeyVaultSecrets

# 配置日志
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

def main():
    # Key Vault名称和用于获取凭证的密钥名称
    keyvault_name = ""  # 替换为实际Key Vault名称
    secret_name = ""  # 替换为Key Vault中的实际密钥名称

    # 初始化Azure Key Vault
    logger.info(f"初始化Azure Key Vault,名称: {keyvault_name}")
    akv_params = AzureKeyVaultParams(keyvault_name=keyvault_name)
    akv = AzureKeyVaultSecrets(params=akv_params)

    # 从Key Vault获取SAS令牌
    logger.info(f"从Azure Key Vault获取密钥: {secret_name}")
    access_token = akv.get_secret(secret_name)
    logger.info("成功从Azure Key Vault获取SAS令牌")

    # 设置ADLS Gen2配置(服务主体、OAuth令牌或账户密钥)
    storage_account_name = "" # ADLS存储账户名称
    container_name = "" # 用于测试的容器名称

    # 初始化Spark会话
    logger.info("初始化Spark会话。")
    spark = SparkSession.builder \
        .appName("Iceberg ADLS Gen2 Example") \
        .config(f"fs.azure.account.key.{storage_account_name}.dfs.core.windows.net", access_token) \
        .getOrCreate()

    logger.info("Spark会话创建成功。")

    # 定义Iceberg表在ADLS Gen2中的存储路径
    table_path = f"abfss://{container_name}@{storage_account_name}.dfs.core.windows.net/my_database/my_iceberg_table"

    # 在基于Hadoop的Iceberg Catalog下的ADLS Gen2中创建新数据库
    logger.info("在基于Hadoop的Iceberg Catalog下的ADLS Gen2中创建数据库。")
    spark.sql("CREATE DATABASE IF NOT EXISTS spark_catalog.my_database")
    logger.info("数据库创建成功(如果不存在)。")

    # 在ADLS Gen2中创建带明确路径的Iceberg表
    logger.info("在ADLS Gen2中创建Iceberg表。")
    spark.sql(f"""
        CREATE TABLE IF NOT EXISTS spark_catalog.my_database.my_iceberg_table (
            col1 INT,
            col2 STRING
        )
        USING iceberg
    """)
    logger.info(f"Iceberg表创建成功,存储位置: {table_path}")

    # 停止Spark会话
    logger.info("停止Spark会话。")
    spark.stop()

if __name__ == "__main__":
    main()

完整报错信息

Multiple sources found for iceberg
(com.databricks.sql.transaction.tahoe.uniform.sources.IcebergBrowseOnlyDataSource,
org.apache.iceberg.spark.source.IcebergSource), please specify the
fully qualified class name. File , line 62
59 spark.stop()
61 if name == "main":
---> 62 main() File , line 48, in main()
46 # Create an Iceberg table with an explicit path for storage in ADLS Gen2
47 logger.info("Creating Iceberg table in ADLS Gen2.")
---> 48 spark.sql(f"""
49 CREATE TABLE IF NOT EXISTS spark_catalog.my_database.my_iceberg_table (
50 col1 INT,
51 col2 STRING
52 )
53 USING iceberg
54 """)
55 logger.info(f"Iceberg table created successfully at location: {table_path}")
57 # Stop the Spark session File /databricks/spark/python/pyspark/instrumentation_utils.py:47, in
_wrap_function..wrapper(*args, **kwargs)
45 start = time.perf_counter()
46 try:
---> 47 res = func(*args, **kwargs)
48 logger.log_success(
49 module_name, class_name, function_name, time.perf_counter() - start, signature
50 )
51 return res File /databricks/spark/python/pyspark/sql/session.py:1748, in
SparkSession.sql(self, sqlQuery, args, **kwargs) 1744
assert self._jvm is not None 1745 litArgs =
self._jvm.PythonUtils.toArray( 1746
[_to_java_column(lit(v)) for v in (args or [])] 1747 )
-> 1748 return DataFrame(self._jsparkSession.sql(sqlQuery, litArgs), self) 1749 finally: 1750 if len(kwargs) > 0: File
/databricks/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py:1355,
in JavaMember.call(self, *args) 1349 command =
proto.CALL_COMMAND_NAME +\ 1350 self.command_header +\ 1351
args_command +\ 1352 proto.END_COMMAND_PART 1354 answer =
self.gateway_client.send_command(command)
-> 1355 return_value = get_return_value( 1356 answer, self.gateway_client, self.target_id, self.name) 1358 for temp_arg
in temp_args: 1359 if hasattr(temp_arg, "_detach"): File
/databricks/spark/python/pyspark/errors/exceptions/captured.py:230, in
capture_sql_exception..deco(*a, **kw)
226 converted = convert_exception(e.java_exception)
227 if not isinstance(converted, UnknownException):
228 # Hide where the exception came from that shows a non-Pythonic
229 # JVM exception message.
--> 230 raise converted from None
231 else:
232 raise

解决方案

可以通过以下方式强制集群使用Apache Iceberg官方运行时,覆盖Databricks的Uniform Reader:

1. 添加Spark配置参数

在集群的Spark配置中新增以下两项:

  • spark.sql.sources.providers.iceberg:org.apache.iceberg.spark.source.IcebergSource
  • spark.databricks.delta.uniform.enabled:false

参数作用:

  • 第一项明确指定iceberg对应的数据源类为Apache Iceberg官方实现,解决多数据源冲突问题
  • 第二项禁用Databricks的Uniform功能,彻底关闭Delta与Iceberg的兼容交互层

2. 直接指定数据源全类名(临时方案)

如果不想修改集群配置,可在CREATE TABLE语句中替换USING iceberg为全类名:

CREATE TABLE IF NOT EXISTS spark_catalog.my_database.my_iceberg_table (
    col1 INT,
    col2 STRING
)
USING org.apache.iceberg.spark.source.IcebergSource

3. 验证配置有效性

修改配置后重启集群,执行建表操作即可避免冲突。若仍有问题,可检查:

  • Apache Iceberg runtime版本是否与Databricks Runtime兼容(14.3 LTS对应Spark 3.5,Iceberg 1.5.2兼容)
  • 确保spark.databricks.sql.iceberg.handle-tables保持为false,阻止Databricks接管Iceberg表处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 16:24:51