Catalog 多了,统一管理就成了新问题
前四篇分别演示了 Spark、Kafka Connect、Flink、Trino 接入 OSS Tables。四个引擎、四种配置写法,但连的是同一个东西:OSS Tables 对外暴露的那一个 Iceberg REST Catalog 端点。表、Schema、快照都在同一份元数据里,各引擎看到的命名空间完全一致,新建一张表也不需要去每个引擎那里确认一遍。换句话说,引擎数量增长本身并不会让元数据碎片化,这正是 REST Catalog 这套标准协议的价值所在。
真正的麻烦来自另一个维度:一家企业里的 Catalog 往往不止 OSS Tables 一个。存量的自建 Hive Metastore 还在跑,兄弟团队有自己的 Iceberg REST Catalog,业务侧有 MySQL、PostgreSQL 这类 JDBC 库表,还有大量非表形态的训练数据和日志文件集。它们是彼此独立的元数据源,各有各的命名空间、权限模型和接入方式。于是,当分析师想知道“公司里到底有哪些数据”,得跑好几个系统分别查,引擎要做跨源查询,得把这堆 Catalog 逐个配进来,再接入一个新引擎,这份配置还得重抄一遍——成本是“引擎数 × Catalog 数”。
Apache Gravitino 要解的就是这个乘法问题。作为生态里的元数据中枢,它是一个开源的统一元数据管理平台,定位是“数据湖的元数据联邦层”:把多个异构 Catalog 挂在同一个 Metalake 之下,不管底层是 Iceberg、Hive、JDBC 还是 Fileset,在 Gravitino 里都是一致的命名空间、统一的元数据视图和权限治理。
OSS Tables 兼容 Iceberg REST Catalog 协议,因此可以直接作为 Gravitino 的一个 Catalog Backend 挂载进来;而 Gravitino 自带的 Iceberg REST 服务又会把整个 Metalake 重新以标准 REST Catalog 协议暴露出去。最直接的效果是:Spark、Flink、Trino 只连 Gravitino 这一个入口,就能在同一份 SQL 里同时访问 OSS Tables 和 Hive、JDBC 上的存量数据,而端点、warehouse ARN、签名方式这些接入细节收敛到 Gravitino 一侧统一维护。
前提条件
已部署 Apache Gravitino 服务。
已创建 OSS Tables 的 Table Bucket。如未创建,请参见OSS Tables。
已安装 Spark 3.5 及以上版本,用于通过 Gravitino 访问 OSS Tables 数据。
步骤一:配置环境变量
Gravitino Iceberg REST 服务和后续访问 OSS Tables 的 Spark 客户端均需要凭证。环境变量名使用 AWS_ 前缀,是因为底层 Iceberg SigV4 签名模块和 S3FileIO 复用 AWS SDK 标准凭证链;实际填入的是阿里云账号的凭证。
配置 Gravitino 服务环境变量
在启动 Gravitino Iceberg REST 服务前设置以下环境变量。配置后,创建 Catalog 时无需填写rest.access-key-id和rest.secret-access-key。
export AWS_ACCESS_KEY_ID=<阿里云AccessKey ID> export AWS_SECRET_ACCESS_KEY=<阿里云AccessKey Secret> export AWS_REGION=<地域,例如cn-hangzhou> # 可选,使用STS临时凭证时配置 export AWS_SESSION_TOKEN=<阿里云STS TOKEN>配置 Spark 环境变量
在启动 Spark Driver 和 Executor 前设置以下环境变量。配置后,无需填写spark.sql.catalog.oss_cata.s3.access-key-id和spark.sql.catalog.oss_cata.s3.secret-access-key。
export AWS_ACCESS_KEY_ID=<阿里云AccessKey ID> export AWS_SECRET_ACCESS_KEY=<阿里云AccessKey Secret> export AWS_REGION=<地域,例如cn-hangzhou> # 可选,使用STS临时凭证时配置 export AWS_SESSION_TOKEN=<阿里云STS TOKEN> export AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED说明:生产环境推荐使用环境变量或 STS 临时凭证,避免在配置文件中明文写入长期 AccessKey。分布式部署时,需确保 Gravitino 服务进程以及 Spark Driver、Executor 均能读取各自所需的环境变量。
步骤二:准备依赖包
由于目前 Iceberg 社区版本的 OSSFileIO 性能欠佳,建议客户端配置 S3 FileIO。为避免 Gravitino 返回的 io-impl 覆盖客户端设置,需要修改 Gravitino 源码:编辑iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/CatalogWrapperForREST.java,删除catalogPropertiesToClientKeys中的IcebergConstants.IO_IMPL一行:
private static final Set<String> catalogPropertiesToClientKeys = ImmutableSet.of( // IcebergConstants.IO_IMPL, // 删除这一行 IcebergConstants.AWS_S3_REGION, IcebergConstants.ICEBERG_S3_ENDPOINT, IcebergConstants.ICEBERG_OSS_ENDPOINT,此外,还需将以下依赖包部署到 Gravitino 的libs目录中:
依赖包 | 版本 | 用途 |
gravitino-iceberg-aws-bundle-xxx.jar | 与 Gravitino 版本匹配 | 加载 S3 相关依赖。建议在 Gravitino 源码中执行 |
步骤三:创建 OSS Tables 相关资源
在配置 Gravitino Catalog 前,需先准备好 OSS Tables 的 Table Bucket,并获取其接入端点信息。
创建 Table Bucket
若尚未创建 Table Bucket,可通过 ossutil 或 AWS CLI 创建,详细步骤请参见OSS Tables。创建完成后,Table Bucket 的 ARN 格式为acs:osstables:{region}:{accountId}:bucket/{bucketName},该 ARN 将作为步骤三中 Gravitino Catalog 的warehouse参数值。
获取接入端点
OSS Tables 提供 Iceberg REST Catalog 端点,Gravitino 通过该端点连接并管理 OSS Tables 的表元数据;该端点将作为步骤三中 Gravitino Catalog 的uri参数值。端点格式如下:
内网:
https://{region}-internal.oss-tables.aliyuncs.com/iceberg外网:
https://{region}.oss-tables.aliyuncs.com/iceberg
说明:底层数据文件通过 S3FileIO 访问 OSS 数据面(OSS 端点格式为https://oss-{region}-internal.aliyuncs.com(内网)或https://oss-{region}.aliyuncs.com(外网))。该端点由 OSS Tables REST Catalog 返回并由客户端 S3FileIO 使用,无需在 Gravitino Catalog 中单独指定(相关依赖与源码调整见步骤一)。
步骤四:创建基于 OSS Tables 的 Gravitino REST Catalog
以下示例中,Metalake 名称均为my_metalake,Gravitino Iceberg REST 服务绑定的也是该 Metalake。
配置 Gravitino Iceberg REST 服务
在 Gravitino 配置文件中配置 Iceberg REST Catalog 的 Metalake:
# ===== Iceberg REST Auxiliary Service ===== gravitino.auxService.names = iceberg-rest # 按需更改,上述依赖包可放在该目录下 gravitino.iceberg-rest.classpath = iceberg-rest-server/libs, iceberg-rest-server/conf gravitino.iceberg-rest.host = 0.0.0.0 gravitino.iceberg-rest.httpPort = 9001 gravitino.iceberg-rest.catalog-config-provider = dynamic-config-provider gravitino.iceberg-rest.gravitino-metalake = my_metalake创建 Metalake
curl -X POST http://localhost:8090/api/metalakes \ -H "Content-Type: application/json" \ -d '{ "name": "my_metalake", "comment": "Main metalake", "properties": {} }'创建 Catalog
创建基于 OSS Tables 的 Catalog:
curl -X POST http://localhost:8090/api/metalakes/my_metalake/catalogs \ -H "Content-Type: application/json" \ -d '{ "name": "osstables_catalog", "type": "RELATIONAL", "provider": "lakehouse-iceberg", "comment": "OSS Tables via Iceberg REST", "properties": { "catalog-backend": "custom", "catalog-backend-impl": "org.apache.iceberg.rest.RESTCatalog", "uri": "https://cn-hangzhou-internal.oss-tables.aliyuncs.com/iceberg", "warehouse": "acs:osstables:cn-hangzhou:{accountId}:bucket/your-table-bucket-name", "rest.auth.type": "sigv4", "rest.signing-region": "cn-hangzhou", "rest.signing-name": "osstables" } }'说明:配置中的 AccessKey、Endpoint、Region 请根据实际情况修改。其中uri为 OSS Tables 的 Iceberg REST Catalog 端点,内网格式为https://{region}-internal.oss-tables.aliyuncs.com/iceberg;warehouse为 Table Bucket 的 ARN,格式为acs:osstables:{region}:{accountId}:bucket/{bucketName}。
步骤五:通过 Gravitino 访问 OSS Tables
完成上述操作后,Gravitino 上已构建名为osstables_catalog的 Catalog。在 Spark 中配置以下参数,即可通过 Gravitino 访问 OSS Tables:
spark.sql.catalog.oss_cata=org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.oss_cata.type=rest spark.sql.catalog.oss_cata.io-impl=org.apache.iceberg.aws.s3.S3FileIO spark.sql.catalog.oss_cata.uri=http://127.0.0.1:9001/iceberg/ spark.sql.catalog.oss_cata.warehouse=osstables_catalog spark.sql.catalog.oss_cata.s3.region=<地域,例如cn-hangzhou> spark.sql.catalog.oss_cata.s3.endpoint=https://oss-{region}-internal.aliyuncs.com说明:必须显式配置spark.sql.catalog.oss_cata.io-impl=org.apache.iceberg.aws.s3.S3FileIO。由于步骤一已移除 Gravitino 对io-impl的透传,客户端需自行指定使用 S3FileIO,否则将回退到 HadoopFileIO 并因无法识别oss://路径而报错。同时需将s3.endpoint指向 OSS 数据面端点,并在启动 Spark 前设置环境变量AWS_REQUEST_CHECKSUM_CALCULATION=WHEN_REQUIRED(关闭 OSS 暂不支持的分块上传编码)。
配置完成后,即可使用标准 SQL 操作 OSS Tables 中的数据。以下示例中 Catalog 名称均为oss_cata。
管理Namespace
Namespace(命名空间)用于对表进行逻辑分组,作用相当于数据库。
-- 查看现有Namespace SHOW NAMESPACES IN oss_cata; -- 创建Namespace CREATE NAMESPACE oss_cata.my_namespace; -- 删除Namespace(需先删除其中所有Table) DROP NAMESPACE oss_cata.my_namespace;建表与表管理
-- 创建非分区表 CREATE TABLE oss_cata.my_namespace.users ( id BIGINT NOT NULL COMMENT '用户ID', name STRING COMMENT '用户名', email STRING COMMENT '邮箱', created_at TIMESTAMP COMMENT '创建时间' ) USING iceberg; -- 创建分区表(按天分区) CREATE TABLE oss_cata.my_namespace.events ( id BIGINT NOT NULL, event_type STRING, data STRING, ts TIMESTAMP ) USING iceberg PARTITIONED BY (days(ts)); -- 查看Namespace中的所有Table SHOW TABLES IN oss_cata.my_namespace; -- 查看表结构 DESCRIBE TABLE oss_cata.my_namespace.users; -- 删除表(OSS Tables 要求必须带 PURGE 关键字,否则会报错:OSS Tables only supports dropping tables with purge enabled) DROP TABLE oss_cata.my_namespace.users PURGE;数据写入与查询
-- 插入数据 INSERT INTO oss_cata.my_namespace.users VALUES (1, '张三', 'zhangsan@example.com', TIMESTAMP '2024-01-15 10:30:00'), (2, '李四', 'lisi@example.com', TIMESTAMP '2024-01-16 14:20:00'), (3, '王五', 'wangwu@example.com', TIMESTAMP '2024-01-17 09:15:00'); -- 全表查询 SELECT * FROM oss_cata.my_namespace.users; -- 条件查询 SELECT * FROM oss_cata.my_namespace.users WHERE id = 2; -- 聚合查询 SELECT COUNT(*) AS total FROM oss_cata.my_namespace.users; -- 分组聚合 SELECT name, COUNT(*) AS cnt FROM oss_cata.my_namespace.users GROUP BY name; -- 更新数据 UPDATE oss_cata.my_namespace.users SET name = '赵六' WHERE id = 3; -- 删除数据 DELETE FROM oss_cata.my_namespace.users WHERE id = 1; -- 查询验证 SELECT * FROM oss_cata.my_namespace.users ORDER BY id;分区表操作
-- 插入分区数据 INSERT INTO oss_cata.my_namespace.events VALUES (1, 'click', '{"page": "home"}', TIMESTAMP '2024-01-15 10:30:00'), (2, 'view', '{"page": "product"}', TIMESTAMP '2024-01-15 11:00:00'), (3, 'click', '{"page": "detail"}', TIMESTAMP '2024-01-16 09:00:00'); -- 分区裁剪查询(仅扫描匹配分区) SELECT * FROM oss_cata.my_namespace.events WHERE ts >= TIMESTAMP '2024-01-15 00:00:00' AND ts < TIMESTAMP '2024-01-16 00:00:00'; -- 聚合统计 SELECT event_type, COUNT(*) AS cnt FROM oss_cata.my_namespace.events GROUP BY event_type;时间旅行查询
Iceberg 支持时间旅行(Time Travel)查询,可读取历史某个时间点的数据快照。
-- 查看快照历史 SELECT snapshot_id, committed_at, operation FROM oss_cata.my_namespace.users.snapshots; -- 基于快照ID查询历史数据 SELECT * FROM oss_cata.my_namespace.users VERSION AS OF <snapshot_id>; -- 查询指定时间点的数据 SELECT * FROM oss_cata.my_namespace.users TIMESTAMP AS OF TIMESTAMP '2024-01-16 00:00:00'; -- 查看数据文件分布 SELECT * FROM oss_cata.my_namespace.users.files;