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

Iceberg结合Hive Metastore在Spark中无法生成独立目录问题咨询

Iceberg自定义Catalog未在Hive Metastore中创建的问题

问题现象

按照Iceberg官方文档配置Spark连接Hive Metastore后,出现以下问题:

  • 配置名为hive_catalog的Iceberg Catalog,成功执行创建namespace、表、写入数据等操作,但查看Hive元数据CTLGS表,仅存在默认的hive catalog,新创建的namespace和表均关联到该默认catalog。
  • 配置另一个名为other_catalog的Catalog并尝试创建同名namespace时,提示已存在,且新namespace仍关联到hive catalog,无法通过不同自定义Catalog隔离同名namespace。

测试验证

第一段测试PySpark代码

import os
from pyspark.sql import SparkSession

deps = [
    "org.apache.iceberg:iceberg-spark-runtime-3.3_2.12:1.2.1",
    "org.apache.iceberg:iceberg-aws:1.2.1",
    "software.amazon.awssdk:bundle:2.17.257",
    "software.amazon.awssdk:url-connection-client:2.17.257"
]
os.environ["PYSPARK_SUBMIT_ARGS"] = f"--packages {','.join(deps)} pyspark-shell"
os.environ["AWS_ACCESS_KEY_ID"] = "minioadmin"
os.environ["AWS_SECRET_ACCESS_KEY"] = "minioadmin"
os.environ["AWS_REGION"] = "eu-east-1"


catalog = "hive_catalog"
spark = SparkSession.\
    builder.\
    appName("Iceberg Reader").\
    config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions").\
    config(f"spark.sql.catalog.{catalog}", "org.apache.iceberg.spark.SparkCatalog").\
    config(f"spark.sql.catalog.{catalog}.type", "hive").\
    config(f"spark.sql.catalog.{catalog}.uri", "thrift://localhost:9083").\
    config(f"spark.sql.catalog.{catalog}.io-impl", "org.apache.iceberg.aws.s3.S3FileIO") .\
    config(f"spark.sql.catalog.{catalog}.s3.endpoint", "http://localhost:9000").\
    config(f"spark.sql.catalog.{catalog}.warehouse", "s3a://lakehouse").\
    config("hive.metastore.uris", "thrift://localhost:9083").\
    enableHiveSupport().\
    getOrCreate()

# Raises error
spark.sql("CREATE NAMESPACE wrong_catalog.new_db;")

# Correct creation of namespace
spark.sql(f"CREATE NAMESPACE {catalog}.new_db;")

# Create table
spark.sql(f"CREATE TABLE {catalog}.new_db.new_table (col1 INT, col2 STRING);")

# Insert data
spark.sql(f"INSERT INTO {catalog}.new_db.new_table VALUES (1, 'first'), (2, 'second');")

# Read data
spark.sql(f"SELECT * FROM {catalog}.new_db.new_table;").show()
#|col1|  col2|
#+----+------+
#|   1| first|
#|   2|second|
#+----+------+

# Read metadata
spark.sql(f"SELECT * FROM {catalog}.new_db.new_table.files;").show()
#+-------+--------------------+-----------+-------+------------+------------------+------------------+----------------+-----------------+----------------+--------------------+--------------------+------------+-------------+------------+-------------+--------------------+
#|content|           file_path|file_format|spec_id|record_count|file_size_in_bytes|      column_sizes|    value_counts|null_value_counts|nan_value_counts|        lower_bounds|        upper_bounds|key_metadata|split_offsets|equality_ids|sort_order_id|    readable_metrics|
#+-------+--------------------+-----------+-------+------------+------------------+------------------+----------------+-----------------+----------------+--------------------+--------------------+------------+-------------+------------+-------------+--------------------+
#|      0|s3a://lakehouse/n...|    PARQUET|      0|           1|               652|{1 -> 47, 2 -> 51}|{1 -> 1, 2 -> 1}| {1 -> 0, 2 -> 0}|              {}|{1 -> ���, 2 -> ...|{1 -> ���, 2 -> ...|        null|          [4]|        null|            0|{{47, 1, 0, null,...|
#|      0|s3a://lakehouse/n...|    PARQUET|      0|           1|               660|{1 -> 47, 2 -> 53}|{1 -> 1, 2 -> 1}| {1 -> 0, 2 -> 0}|              {}|{1 -> ���, 2 -> ...|{1 -> ���, 2 -> ...|        null|          [4]|        null|            0|{{47, 1, 0, null,...|
#+-------+--------------------+-----------+-------+------------+------------------+------------------+----------------+-----------------+----------------+--------------------+--------------------+------------+-------------+------------+-------------+--------------------+

Hive元数据CTLGS表输出

|CTLG_ID|NAME|DESC                    |LOCATION_URI    |
|-------|----|------------------------|----------------|
|1      |hive|Default catalog for Hive|s3a://lakehouse/|

Hive元数据DBS和TBLS表输出

|DB_ID|DESC                 |DB_LOCATION_URI          |NAME   |OWNER_NAME    |OWNER_TYPE|CTLG_NAME|
|-----|---------------------|-------------------------|-------|--------------|----------|---------|
|1    |Default Hive database|s3a://lakehouse/         |default|public        |ROLE      |hive     |
|2    |                     |s3a://lakehouse/new_db.db|new_db |thijsvandepoll|USER      |hive     |


|TBL_ID|CREATE_TIME  |DB_ID|LAST_ACCESS_TIME|OWNER         |OWNER_TYPE|RETENTION    |SD_ID|TBL_NAME |TBL_TYPE      |VIEW_EXPANDED_TEXT|VIEW_ORIGINAL_TEXT|IS_REWRITE_ENABLED|
|------|-------------|-----|----------------|--------------|----------|-------------|-----|---------|--------------|------------------|------------------|------------------|
|1     |1.683.707.647|2    |80.467          |thijsvandepoll|USER      |2.147.483.647|1    |new_table|EXTERNAL_TABLE|                  |                  |0                 |

第二段测试PySpark代码

import os
from pyspark.sql import SparkSession

deps = [
    "org.apache.iceberg:iceberg-spark-runtime-3.3_2.12:1.2.1",
    "org.apache.iceberg:iceberg-aws:1.2.1",
    "software.amazon.awssdk:bundle:2.17.257",
    "software.amazon.awssdk:url-connection-client:2.17.257"
]
os.environ["PYSPARK_SUBMIT_ARGS"] = f"--packages {','.join(deps)} pyspark-shell"
os.environ["AWS_ACCESS_KEY_ID"] = "minioadmin"
os.environ["AWS_SECRET_ACCESS_KEY"] = "minioadmin"
os.environ["AWS_REGION"] = "eu-east-1"

catalog = "other_catalog"
spark = SparkSession.\
    builder.\
    appName("Iceberg Reader").\
    config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions").\
    config(f"spark.sql.catalog.{catalog}", "org.apache.iceberg.spark.SparkCatalog").\
    config(f"spark.sql.catalog.{catalog}.type", "hive").\
    config(f"spark.sql.catalog.{catalog}.uri", "thrift://localhost:9083").\
    config(f"spark.sql.catalog.{catalog}.io-impl", "org.apache.iceberg.aws.s3.S3FileIO") .\
    config(f"spark.sql.catalog.{catalog}.s3.endpoint", "http://localhost:9000").\
    config(f"spark.sql.catalog.{catalog}.warehouse", "s3a://lakehouse").\
    config("hive.metastore.uris", "thrift://localhost:9083").\
    enableHiveSupport().\
    getOrCreate()

# Error that catalog already exists
spark.sql(f"CREATE NAMESPACE {catalog}.new_db;")
# pyspark.sql.utils.AnalysisException: Namespace 'new_db' already exists

# Create another namespace
spark.sql(f"CREATE NAMESPACE {catalog}.other_db;")

# Try to access data from other catalog using current catalog
spark.sql("SELECT * FROM {catalog}.new_db.new_table;").show()
#|col1|  col2|
#+----+------+
#|   1| first|
#|   2|second|
#+----+------+

Hive元数据DBS表输出

|DB_ID|DESC                 |DB_LOCATION_URI            |NAME    |OWNER_NAME    |OWNER_TYPE|CTLG_NAME|
|-----|---------------------|---------------------------|--------|--------------|----------|---------|
|1    |Default Hive database|s3a://lakehouse/           |default |public        |ROLE      |hive     |
|2    |                     |s3a://lakehouse/new_db.db  |new_db  |thijsvandepoll|USER      |hive     |
|3    |                     |s3a://lakehouse/other_db.db|other_db|thijsvandepoll|USER      |hive     |

问题解答

这是预期行为,原因如下:
Iceberg的Hive Catalog实现是基于Hive Metastore的现有默认catalog(即hive catalog)来存储元数据,不会在Hive Metastore中创建新的catalog条目。你配置的自定义catalog名称(如hive_catalog、other_catalog)只是Spark层面的逻辑别名,用于区分不同的配置参数,但底层的namespace和表元数据仍然会关联到Hive的默认catalog。

如果需要实现catalog级别的隔离,可以采用以下两种方案:

  • 使用独立的Hive Metastore实例:为每个Iceberg catalog配置不同的Hive Metastore服务,实现完全的元数据隔离。
  • 配置不同的warehouse路径:为每个自定义catalog设置不同的spark.sql.catalog.<catalog-name>.warehouse路径,通过物理存储路径来隔离不同catalog的数据和元数据,即使Hive元数据仍在同一个默认catalog下,也能实现逻辑上的隔离。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 12:57:08