Iceberg结合Hive Metastore在Spark中无法生成独立目录问题咨询
Iceberg自定义Catalog未在Hive Metastore中创建的问题
问题现象
按照Iceberg官方文档配置Spark连接Hive Metastore后,出现以下问题:
- 配置名为
hive_catalog的Iceberg Catalog,成功执行创建namespace、表、写入数据等操作,但查看Hive元数据CTLGS表,仅存在默认的hivecatalog,新创建的namespace和表均关联到该默认catalog。 - 配置另一个名为
other_catalog的Catalog并尝试创建同名namespace时,提示已存在,且新namespace仍关联到hivecatalog,无法通过不同自定义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
相关产品推荐
相关产品推荐

