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

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_IDNAMEDESCLOCATION_URI
1hiveHive默认Catalogfile:/user/hive/warehouse
3fintechFintech业务Catalogs3a://fintech
2telcoTelco业务Catalogs3a://telco
4plantsPlants业务Catalogs3a://plants

DBS表内容:

DB_IDDESCDB_LOCATION_URINAMEOWNER_NAMEOWNER_TYPECTLG_NAME
1Hive默认数据库file:/user/hive/warehousedefaultpublicROLEhive
6手动创建的Trino schemas3a://fintech/digipaysdigipayscarl.coxUSERfintech
7手动创建的Trino schemas3a://telco/worldonlineworldonlinecarl.coxUSERtelco
8手动创建的Trino schemas3a://plants/plastixmeters_buscarl.coxUSERplants
10同事通过Spark创建的数据库s3a://corporate/sparkk.dbsparkkjohn.digweedUSERhive
  • 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()

关键修正点

  1. 添加hive.catalog-name配置:每个Catalog都需要指定该参数,明确关联HMS中对应的Catalog记录,这是Trino配置hive.metastore.thrift.catalog-name在Spark中的对应参数。
  2. 移除错误的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 17:54:53