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.NativeAzureFileSystemspark.sql.catalog.spark_catalog:org.apache.iceberg.spark.SparkCatalogspark.hadoop.fs.abfss.impl:org.apache.hadoop.fs.azurebfs.AzureBlobFileSystemspark.databricks.sql.iceberg.handle-tables:falsespark.sql.extensions:org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensionsspark.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.4org.apache.hadoop:hadoop-azure:3.4.0org.apache.hadoop:hadoop-client-runtime:3.4.0org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.5.2pandas==2.2.2pyarrow==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.IcebergSourcespark.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

