ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

DataHub Iceberg Catalog 实战指南:以 DataHub 为统一目录服务管理 Iceberg 表

2026/9/17 1:43:03 拓冰建站 浏览量
DataHub Iceberg Catalog 实战指南:以 DataHub 为统一目录服务管理 Iceberg 表 DataHub Iceberg Catalog 实战指南以 DataHub 为统一目录服务管理 Iceberg 表【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本文基于本仓库 docs/iceberg-catalog.md 编写配套示例代码位于 metadata-ingestion/examples/iceberg。文章面向希望用 DataHub 替代或代理自有 Iceberg REST Catalog、打通「目录管理 权限管控 元数据发现」的开发者与平台工程师读完即可完成从仓库开通、建表、读写到公开读配置的完整闭环。DataHub Iceberg Catalog 将 DataHub 的元数据服务GMS暴露为一个符合 Apache Iceberg REST Catalog 协议的目录端点让 PyIceberg、Spark 等计算引擎可以像使用原生 Iceberg Catalog 一样通过 DataHub 创建、读取和管理 Iceberg 表。核心价值在于元数据天然双向一致Iceberg 中建的表会立刻以 Dataset 形式出现在 DataHub 中供人类与搜索发现访问受 DataHub 权限模型管控vended credentials 机制按策略授予临时 S3 凭证从而把「数据目录」与「数据权限」统一到同一套体系内。注意该功能当前处于开放Beta阶段对应 Open Source DataHub 1.0.0 版本。若遇问题或建议请联系 DataHub 维护方。适用场景与概念映射文档明确了四类典型用途通过 DataHub 创建并管理 Iceberg 表保持 DataHub 与 Iceberg 两侧元数据一致将 Iceberg 表元数据暴露在 DataHub 中便于数据发现借助 DataHub 权限模型为 Iceberg 表提供安全访问控制。为了让两类模型对齐文档给出了如下的概念映射表这是理解后续所有配置尤其是权限与资源类型的基础Iceberg 概念DataHub 概念说明Overall platformdataPlatformicebergplatformWarehousedataPlatformInstance存储凭据等仓库级信息Namespacecontainer子类型为 Namespace 的层级容器Tables, ViewsdatasetDataset URN 为 UUID重命名后保持不变Table/View NamesplatformResource指向 Dataset随重命名操作变化其中值得注意的设计是Dataset 的 URN 使用与表名解耦的 UUID而表名/视图名通过platformResource关联因此表重命名不会导致 URN 变化血缘与权限的锚点保持稳定。前置条件开始前需要准备以下四类前提本地已安装并运行 DataHubdatahubCLI 已安装且配置指向该实例CLI 的默认图客户端认证方式见 metadata-ingestion/src/datahub/cli 下的 client 实现配置好 AWS 凭据及相应权限设置以下环境变量DH_ICEBERG_CLIENT_IDyour_client_id DH_ICEBERG_CLIENT_SECRETyour_client_secret DH_ICEBERG_AWS_ROLEarn:aws:iam::123456789012:role/your-role-name # Example format DH_ICEBERG_DATA_ROOTs3://your-bucket/path其中DH_ICEBERG_CLIENT_ID即AWS_ACCESS_KEY_IDDH_ICEBERG_CLIENT_SECRET即AWS_SECRET_ACCESS_KEY。若使用 PyIceberg需要按 PyIceberg 支持的任一方式将本地 DataHub 配置为目录。例如创建~/.pyiceberg.yamlcatalog: local_datahub: uri: http://localhost:8080/iceberg warehouse: arctic_warehouse版本说明本文 Python 代码片段基于仓库 metadata-ingestion/examples/iceberg 目录中的脚本需要安装pyiceberg[duckdb] 0.8.1该目录的 requirements.txt 还要求pyarrow 19.0.0Spark 示例的实测版本为3.5.3_2.12。必需的 AWS 权限DH_ICEBERG_AWS_ROLE指定的角色必须对DH_ICEBERG_DATA_ROOT指向的 S3 位置具备读写权限。权限应只授予 Iceberg 表实际存储的 S3 bucket 与路径前缀。此外角色还必须配置信任策略允许使用DH_ICEBERG_CLIENT_ID/DH_ICEBERG_CLIENT_SECRET提供的 AWS 凭据来 AssumeRole——这正是下文「Vended Credentials」机制能工作的前提。搭建步骤1. 开通 WarehouseProvision a Warehouse在 DataHub 中创建一个 Iceberg warehouse。两种方式任选其一方式 ACLIdatahub iceberg create -w arctic_warehouse -d $DH_ICEBERG_DATA_ROOT -i $DH_ICEBERG_CLIENT_ID --client_secret $DH_ICEBERG_CLIENT_SECRET --region us-east-1 --role $DH_ICEBERG_AWS_ROLE方式 BPythonpyiceberg# File: provision_warehouse.py import os from constants import warehouse # Assert that env variables are present assert os.environ.get(DH_ICEBERG_CLIENT_ID), ( DH_ICEBERG_CLIENT_ID variable is not present ) assert os.environ.get(DH_ICEBERG_CLIENT_SECRET), ( DH_ICEBERG_CLIENT_SECRET variable is not present ) assert os.environ.get(DH_ICEBERG_AWS_ROLE), ( DH_ICEBERG_AWS_ROLE variable is not present ) assert os.environ.get(DH_ICEBERG_DATA_ROOT), ( DH_ICEBERG_DATA_ROOT variable is not present ) assert os.environ.get(DH_ICEBERG_DATA_ROOT, ).startswith(s3://) os.system( fdatahub iceberg create --warehouse {warehouse} --data_root $DH_ICEBERG_DATA_ROOT/{warehouse} --client_id $DH_ICEBERG_CLIENT_ID --client_secret $DH_ICEBERG_CLIENT_SECRET --region us-east-1 --role $DH_ICEBERG_AWS_ROLE )仓库中的 provision_warehouse.py 与上例一致可作为可直接运行的模板。从 CLI 实现iceberg_cli.py可以看到create命令的完整行为以urn:li:dataPlatformInstance:(iceberg,warehouse)作为仓库 URN若已存在则报错退出依次执行validate_warehouse(data_root)、validate_creds(client_id, client_secret, region)、validate_role(role, storage_client, duration_seconds)三段校验保证 S3 路径、凭据与角色在创建时即有效将clientId与clientSecret写入 DataHub Secret生成clientIdUrn/clientSecretUrn仓库 aspect 记录dataRoot、region、role、env与可选的tempCredentialExpirationSeconds。此外还支持可选参数-d/--description描述、-e/--env环境默认值见DEFAULT_FABRIC_TYPE、-x/--duration_seconds临时凭据有效期。开通后必须授予权限确保你的 DataHub 用户对资源类型Data Platform Instance拥有以下 4 项特权随 Iceberg 支持引入DATA_MANAGE_VIEWS_PRIVILEGEDATA_MANAGE_TABLES_PRIVILEGEDATA_MANAGE_NAMESPACES_PRIVILEGEDATA_LIST_ENTITIES_PRIVILEGE可在 DataHub UI 的 Policies 区块中授予。2. 创建表创建表前必须先创建 Namespace。下面以「滑雪场指标表」为例展示两种引擎方式 Aspark-sql启动前需要定义$GMS_HOST、$GMS_PORT、$WAREHOUSE与$USER_PATDataHub 个人访问令牌用于连接目录。DataHub Cloud 的 Iceberg Catalog URL 为https://your-instance.acryl.io/gms/iceberg/本地运行则将GMS_HOST设为localhost、GMS_PORT设为8080。示例中WAREHOUSE为arctic_warehouse。spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.6.1,org.apache.iceberg:iceberg-aws-bundle:1.6.1 \ --conf spark.sql.extensionsorg.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \ --conf spark.sql.catalog.spark_catalogorg.apache.iceberg.spark.SparkSessionCatalog \ --conf spark.sql.catalog.localorg.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.local.typerest \ --conf spark.sql.catalog.local.urihttp://$GMS_HOST:$GMS_PORT/iceberg/ \ --conf spark.sql.catalog.local.warehouse$WAREHOUSE \ --conf spark.sql.catalog.local.token$USER_PAT \ --conf spark.sql.catalog.local.rest-metrics-reporting-enabledfalse \ --conf spark.sql.catalog.local.header.X-Iceberg-Access-Delegationvended-credentials \ --conf spark.sql.defaultCataloglocal然后执行 SQLCREATE NAMESPACE alpine_db; CREATE TABLE alpine_db.ski_resorts ( resort_id BIGINT NOT NULL COMMENT Unique identifier for each ski resort, resort_name STRING NOT NULL COMMENT Official name of the ski resort, daily_snowfall BIGINT COMMENT Amount of new snow in inches during the last 24 hours, conditions STRING COMMENT Current snow conditions description, last_updated TIMESTAMP COMMENT Timestamp of when the snow report was last updated );方式 BPythonpyicebergfrom constants import namespace, table_name, warehouse from pyiceberg.schema import Schema from pyiceberg.types import LongType, NestedField, StringType, TimestampType from pyiceberg.catalog import load_catalog from datahub.ingestion.graph.client import get_default_graph # Get DataHub graph client for authentication graph get_default_graph() # Define schema with documentation schema Schema( NestedField( field_id1, nameresort_id, field_typeLongType(), requiredTrue, docUnique identifier for each ski resort ), NestedField( field_id2, nameresort_name, field_typeStringType(), requiredTrue, docOfficial name of the ski resort ), NestedField( field_id3, namedaily_snowfall, field_typeLongType(), requiredFalse, docAmount of new snow in inches during the last 24 hours ), NestedField( field_id4, nameconditions, field_typeStringType(), requiredFalse, docCurrent snow conditions description ), NestedField( field_id5, namelast_updated, field_typeTimestampType(), requiredFalse, docTimestamp of when the snow report was last updated ) ) # Load catalog and create table catalog load_catalog(local_datahub, warehousewarehouse, tokengraph.config.token) catalog.create_namespace(namespace) catalog.create_table(f{namespace}.{table_name}, schema)仓库中的 create_table.py 是更完整的版本使用constants.py中的warehouse arctic_warehouse、namespace alpine_db、table_name resort_metrics并对「命名空间已存在」「表已存在」两种异常做了容错处理建表后立刻写入样本数据并读回验证。关于表名与示例差异的说明Spark 示例建表名为alpine_db.ski_resortsPython 示例基于constants.py使用alpine_db.resort_metrics且文档第 4 节读取示例中 SQL 查询的是alpine_db.resort_metrics。实际使用时请保持建表名与后续读写语句一致两者都是同一套目录机制的演示。3. 写入数据spark-sqlINSERT INTO alpine_db.ski_resorts (resort_id, resort_name, daily_snowfall, conditions, last_updated) VALUES (1, Snowpeak Resort, 12, Powder, CURRENT_TIMESTAMP()), (2, Alpine Valley, 8, Packed, CURRENT_TIMESTAMP()), (3, Glacier Heights, 15, Fresh Powder, CURRENT_TIMESTAMP());Pythonpyiceberg PyArrow写入时必须精确匹配 Iceberg 表的 schemafrom constants import namespace, table_name, warehouse from pyiceberg.catalog import load_catalog from datahub.ingestion.graph.client import get_default_graph import pyarrow as pa from datetime import datetime graph get_default_graph() catalog load_catalog(local_datahub, warehousewarehouse, tokengraph.config.token) # Create PyArrow schema to match Iceberg schema pa_schema pa.schema([ (resort_id, pa.int64(), False), # False means not nullable (resort_name, pa.string(), False), (daily_snowfall, pa.int64(), True), (conditions, pa.string(), True), (last_updated, pa.timestamp(us), True), ]) # Create sample data sample_data pa.Table.from_pydict( { resort_id: [1, 2, 3], resort_name: [Snowpeak Resort, Alpine Valley, Glacier Heights], daily_snowfall: [12, 8, 15], conditions: [Powder, Packed, Fresh Powder], last_updated: [ datetime.now(), datetime.now(), datetime.now() ] }, schemapa_schema ) # Write to table table catalog.load_table(f{namespace}.{table_name}) table.overwrite(sample_data) # Refresh table to see changes table.refresh()注意 PyArrow 类型与 Iceberg 类型的对应BIGINT↔pa.int64()、STRING↔pa.string()、TIMESTAMP↔pa.timestamp(us)元组第三个布尔值表示是否可空与建表时的required一致。4. 读取数据spark-sqlSELECT * from alpine_db.resort_metrics;Pythonpyiceberg DuckDB借助pyiceberg[duckdb]将扫描结果落到内存 DuckDB 中查询from pyiceberg.catalog import load_catalog from constants import namespace, table_name, warehouse from datahub.ingestion.graph.client import get_default_graph # Get DataHub graph client for authentication graph get_default_graph() catalog load_catalog(local_datahub, warehousewarehouse, tokengraph.config.token) table catalog.load_table(f{namespace}.{table_name}) con table.scan().to_duckdb(table_nametable_name) # Query the data print(\nResort Metrics Data:) print(- * 50) for row in con.execute(fSELECT * FROM {table_name}).fetchall(): print(row)仓库中 read_table.py 是上述读取逻辑的独立精简版drop_table.py 则演示了catalog.drop_table(f{namespace}.{table_name})的删表操作。参考信息Trust Policy 示例创建 AWS 角色时需配置信任策略允许你的凭据所对应的 IAM 用户/角色 AssumeRole{ Version: 2012-10-17, Statement: [ { Effect: Allow, Principal: { AWS: arn:aws:iam::123456789012:user/iceberg-user // The IAM user or role associated with your credentials }, Action: sts:AssumeRole } ] }Vended Credentials 与会话有效期创建 warehouse 时vended credentials颁发给计算引擎的临时 S3 凭证的会话有效期可通过datahub icebergCLI 的--duration_seconds参数配置。未指定时默认 3600 秒且最小不能低于 900 秒15 分钟。CLI 默认值定义在 iceberg_cli.py 的DEFAULT_CREDS_EXPIRY_DURATION_SECONDS。对 AWS 而言创建角色时还可以设置MaxSessionDurationAWS 支持 1 到 12 小时的范围——它会作为--duration_seconds配置上限。即最终有效期为「CLI 参数」与「AWS 角色上限」二者的较小者。可通过datahub iceberg create --help查看全部参数及默认值创建后也可用datahub iceberg update修改注意update 目前只能更新凭据与 role不能更新 region参见 iceberg_cli.py 附近实现。Namespace 与配置要求Namespace 必须先创建然后才能在其中建表。Spark SQL 默认配置为避免每条 SQL 都写全 catalog 名与 namespace 前缀可在 Spark 配置中设置默认值--conf spark.sql.defaultCatalogdefault-catalog-name \ --conf spark.sql.catalog.local.default-namespacedefault-namespace公开只读访问 Iceberg 表如需让特定表无需认证即可读取public read-only access按以下三步启用S3 侧确保DH_ICEBERG_DATA_ROOT所在 S3 文件夹配置了公共读策略使该前缀下的文件无需 AWS 凭据即可读取。GMS 侧以如下环境变量启动 GMS二者默认值同时定义在 application.yaml 的icebergCatalog键下ENABLE_PUBLIC_READtrue启用该能力不设置时默认falsePUBLICLY_READABLE_TAG指定表示「允许公开访问」的 Tag 名称不设置时默认为PUBLICLY_READABLE。也可以直接在metadata-service/configuration/src/main/resources/application.yaml的icebergCatalog键下配置这些默认值。打标签GMS 以公开读能力启动后对每个希望免认证访问的 Dataset 打上上述 Tag。访问公开表时spark-sql 需做如下特殊配置spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.6.1,org.apache.iceberg:iceberg-aws-bundle:1.6.1\ --conf spark.sql.extensionsorg.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \ --conf spark.sql.catalog.spark_catalogorg.apache.iceberg.spark.SparkSessionCatalog \ --conf spark.sql.catalog.localorg.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.local.typerest \ --conf spark.sql.catalog.local.urihttp://${GMS_HOST}:${GMS_PORT}/public-iceberg/ \ --conf spark.sql.catalog.local.warehousearctic_warehouse \ --conf spark.sql.catalog.local.header.X-Iceberg-Access-Delegationfalse \ --conf spark.sql.catalog.local.rest-metrics-reporting-enabledfalse \ --conf spark.sql.catalog.local.client.regionus-east-1 \ --conf spark.sql.catalog.local.client.credentials-providersoftware.amazon.awssdk.auth.credentials.AnonymousCredentialsProvider \ --conf spark.sql.defaultCataloglocal \ --conf spark.sql.catalog.local.default-namespacealpine_db与认证访问的关键差异如下REST Catalog URI 变为gms_host:port/**public-iceberg/**对应服务端公开端点见 PublicIcebergApiController.javaDATA ROOT 所在 S3 区域通过spark.sql.catalog.local.client.region指定X-Iceberg-Access-Delegation头设置为false而非vended-credentialscredentials-provider设为software.amazon.awssdk.auth.credentials.AnonymousCredentialsProvider匿名凭据不提供Personal Access Token。需要注意在未认证会话中尝试访问未打公开 Tag的表会以NoSuchTableException失败而不是授权失败错误——这是为了避免向匿名用户泄露表是否存在的信息。DataHub Iceberg CLI 管理命令除create外CLI实现见 iceberg_cli.py还提供以下管理命令列出所有 Warehousedatahub iceberg list更新 Warehouse 配置datahub iceberg update \ -w $WAREHOUSE_NAME \ -d $DH_ICEBERG_DATA_ROOT \ -i $DH_ICEBERG_CLIENT_ID \ --client_secret $DH_ICEBERG_CLIENT_ID \ --region $DH_ICEBERG_REGION \ --role DH_ICEBERG_AWS_ROLE删除 Warehousedatahub iceberg delete -w $WAREHOUSE_NAME注意该命令会删除该 warehouse 下 Namespace 对应的 Containers、以及 Tables/Views 对应的 Datasets但不会删除data_root中任何底层数据文件——它只清理目录条目。删除执行成功后 CLI 会输出删除了多少 datasets 与 namespaces见 iceberg_cli.py 附近。更多选项与帮助datahub iceberg [command] --help从其他 Iceberg Catalog 迁移从其他目录迁移时可使用system.register_table注册现有 Iceberg 表从而在不移动底层数据的前提下通过 DataHub 开始管理这些表call system.register_table(myTable, s3://my-s3-bucket/my-data-root/myNamespace/myTable/metadata/00000-f9dbba67-df4f-4742-9ba5-123aa2bb4076.metadata.json); -- Read from newly registered table select * from myTable;注册时需提供表元数据文件metadata.json的完整 S3 路径。安全与权限模型DataHub 通过策略Policy为 Iceberg 操作提供细粒度访问控制。随 Iceberg 支持引入的权限如下表操作所需特权资源类型CREATE 或 DROP namespacesManage NamespacesData Platform InstanceCREATE、ALTER 或 DROP tablesManage TablesData Platform InstanceCREATE、ALTER 或 DROP viewsManage ViewsData Platform Instance从表或视图 SELECTRead Only contenteditable="false">【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考