Spark Iceberg配置:如何通过单个Hive Metastore访问多Catalog?
问题解决:PySpark 3.5.1 识别Hive Metastore多Catalog及对应数据库
现有环境
- 涉及行业:fintech、telco、plants
- Hive Metastore(基于PostgreSQL)以Docker部署,地址为
thrift://hms-host:9083。PostgreSQL的CTLGS表已手动插入4条Catalog记录:
| CTLG_ID | NAME | DESC | LOCATION_URI |
|---|---|---|---|
| 1 | hive | Hive默认Catalog | file:/user/hive/warehouse |
| 3 | fintech | Fintech业务Catalog | s3a://fintech |
| 2 | telco | Telco业务Catalog | s3a://telco |
| 4 | plants | Plants业务Catalog | s3a://plants |
DBS表内容:
| DB_ID | DESC | DB_LOCATION_URI | NAME | OWNER_NAME | OWNER_TYPE | CTLG_NAME |
|---|---|---|---|---|---|---|
| 1 | Hive默认数据库 | file:/user/hive/warehouse | default | public | ROLE | hive |
| 6 | 手动创建的Trino schema | s3a://fintech/digipays | digipays | carl.cox | USER | fintech |
| 7 | 手动创建的Trino schema | s3a://telco/worldonline | worldonline | carl.cox | USER | telco |
| 8 | 手动创建的Trino schema | s3a://plants/plastix | meters_bus | carl.cox | USER | plants |
| 10 | 同事通过Spark创建的数据库 | s3a://corporate/sparkk.db | sparkk | john.digweed | USER | hive |
- Minio集群端点为
https://my-minio-host:9000,具备access和secret密钥,包含3个存储桶:s3a://fintech、s3a://telco、s3a://plants
目标
编写正确配置的PySpark 3.5.1代码,识别并展示Hive Metastore中的fintech、telco、plants Catalog及对应数据库。
当前问题
现有PySpark代码仅返回hive Catalog下的default和sparkk数据库,无法获取worldonline等预期库。代码及输出如下:
import pyspark from pyspark.sql import SparkSession CATALOG_CLASS = "org.apache.iceberg.spark.SparkCatalog" CATALOG_TYPE = "hive" HMS_URI = "thrift://hms-host:9083" AWS_ENDPOINT = "https://my-minio-host:9000" AWS_ACCESS_KEY = "123" AWS_SECRET_KEY = "12345678" spark = SparkSession.builder \ .appName("IcebergHiveMetastoreIntegration") \ .config("spark.sql.catalog.fintech", CATALOG_CLASS) \ .config("spark.sql.catalog.telco", CATALOG_CLASS) \ .config("spark.sql.catalog.plants", CATALOG_CLASS) \ .config("spark.sql.catalog.fintech.type", CATALOG_TYPE) \ .config("spark.sql.catalog.telco.type", CATALOG_TYPE) \ .config("spark.sql.catalog.plants.type", CATALOG_TYPE) \ .config("spark.sql.catalog.fintech.uri", HMS_URI) \ .config("spark.sql.catalog.telco.uri", HMS_URI) \ .config("spark.sql.catalog.plants.uri", HMS_URI) \ .config("spark.sql.catalog.fintech.default-namespace", "fintech") \ .config("spark.sql.catalog.telco.default-namespace", "telco") \ .config("spark.sql.catalog.plants.default-namespace", "plants") \ .config("spark.hadoop.fs.s3a.endpoint", AWS_ENDPOINT) \ .config("spark.hadoop.fs.s3a.access.key", AWS_ACCESS_KEY) \ .config("spark.hadoop.fs.s3a.secret.key", AWS_SECRET_KEY) \ .config("spark.hadoop.fs.s3a.path.style.access", "true") \ .getOrCreate() spark.sql("USE telco") spark.sql("SHOW CATALOGS").show() spark.sql("SHOW NAMESPACES").show() spark.sql("SHOW DATABASES").show()
输出:
+-------------+ | catalog| +-------------+ | telco| |spark_catalog| +-------------+ +---------+ |namespace| +---------+ | default| | sparkk| +---------+ +---------+ |namespace| +---------+ | default| | sparkk| +---------+
备注
Trino配合Iceberg连接器环境中,已通过配置hive.metastore.thrift.catalog-name=fintech等参数实现预期效果,需在Spark中复现该行为。
解决方案
问题核心是Iceberg的Hive类型Catalog需要明确指定绑定到HMS中的哪个Catalog,对应Trino的hive.metastore.thrift.catalog-name参数,在Spark中配置项为spark.sql.catalog.<catalog-name>.hive.catalog-name。同时,原代码中错误设置了default-namespace为Catalog名称,需调整或移除。
修正后的代码如下:
import pyspark from pyspark.sql import SparkSession CATALOG_CLASS = "org.apache.iceberg.spark.SparkCatalog" CATALOG_TYPE = "hive" HMS_URI = "thrift://hms-host:9083" AWS_ENDPOINT = "https://my-minio-host:9000" AWS_ACCESS_KEY = "123" AWS_SECRET_KEY = "12345678" spark = SparkSession.builder \ .appName("IcebergHiveMetastoreIntegration") \ # 配置各个Catalog实现类 .config("spark.sql.catalog.fintech", CATALOG_CLASS) \ .config("spark.sql.catalog.telco", CATALOG_CLASS) \ .config("spark.sql.catalog.plants", CATALOG_CLASS) \ # 指定Catalog类型为Hive .config("spark.sql.catalog.fintech.type", CATALOG_TYPE) \ .config("spark.sql.catalog.telco.type", CATALOG_TYPE) \ .config("spark.sql.catalog.plants.type", CATALOG_TYPE) \ # 绑定到HMS中对应的Catalog记录 .config("spark.sql.catalog.fintech.hive.catalog-name", "fintech") \ .config("spark.sql.catalog.telco.hive.catalog-name", "telco") \ .config("spark.sql.catalog.plants.hive.catalog-name", "plants") \ # HMS thrift地址 .config("spark.sql.catalog.fintech.uri", HMS_URI) \ .config("spark.sql.catalog.telco.uri", HMS_URI) \ .config("spark.sql.catalog.plants.uri", HMS_URI) \ # Minio存储配置 .config("spark.hadoop.fs.s3a.endpoint", AWS_ENDPOINT) \ .config("spark.hadoop.fs.s3a.access.key", AWS_ACCESS_KEY) \ .config("spark.hadoop.fs.s3a.secret.key", AWS_SECRET_KEY) \ .config("spark.hadoop.fs.s3a.path.style.access", "true") \ .getOrCreate() # 测试telco Catalog spark.sql("USE telco") print("=== Telco Catalog 信息 ===") spark.sql("SHOW CATALOGS").show() spark.sql("SHOW NAMESPACES").show() spark.sql("SHOW DATABASES").show() # 测试fintech Catalog spark.sql("USE fintech") print("\n=== Fintech Catalog 信息 ===") spark.sql("SHOW NAMESPACES").show() # 测试plants Catalog spark.sql("USE plants") print("\n=== Plants Catalog 信息 ===") spark.sql("SHOW NAMESPACES").show()
关键修正点
- 添加
hive.catalog-name配置:每个Catalog都需要指定该参数,明确关联HMS中对应的Catalog记录,这是Trino配置hive.metastore.thrift.catalog-name在Spark中的对应参数。 - 移除错误的
default-namespace配置:原代码将default-namespace设为Catalog名称,但HMS中对应Catalog下的数据库并非该名称(比如telco Catalog下的数据库是worldonline),移除后Spark会自动加载对应Catalog下的所有数据库,也可按需设置为目标数据库名称(如spark.sql.catalog.telco.default-namespace=worldonline)。
执行修正后的代码后,即可正确获取各Catalog下的对应数据库:telco Catalog显示worldonline,fintech显示digipays,plants显示meters_bus。
内容的提问来源于stack exchange,提问作者deeplay
相关产品推荐
相关产品推荐

