Daft × Apache Gravitino:端到端集成测试实战指南(gvfs:// 文件集与 Catalog 目录操作)
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
本篇指南以 Daft 仓库中的 Gravitino 集成测试套件(tests/integration/gravitino/)为骨架,完整讲解如何在本地用 Docker Compose 拉起 Apache Gravitino + MySQL + MinIO 三件套,通过gvfs://协议读写 Gravitino fileset(本地file://与 S3 存储两种后端),并通过 Daft 的 Catalog API 完成目录、Schema、表的元数据操作。读完本文,你将掌握一套可复现的 Gravitino 连接器验证流程,并理解其底层 REST 调用链与 PyArrow 文件系统适配实现。
一、背景:为什么需要 Gravitino 集成测试
Apache Gravitino 是一个开源的多租户元数据中心,为各类数据源与存储系统提供统一元数据管理。在 Daft 中,Gravitino 连接器(位于 daft/catalog/__gravitino/)让用户能够:
- 通过
Catalog.from_gravitino()接入 Daft 的 Catalog 体系,统一列举 catalog/schema/table; - 读取 Iceberg、Hive/Parquet 等多格式表;
- 通过
gvfs://协议直接读写 Gravitino fileset(一种管理原始文件的元数据实体,底层可落在 S3、GCS、Azure Blob 等存储上)。
集成测试(tests/integration/gravitino/)的作用,正是以真实运行的 Gravitino 服务为依赖,端到端验证这套连接器的正确性——包括元数据 CRUD、文件读写、glob 发现、错误处理与清理逻辑。
从实现层面看,该测试覆盖了连接器的两条核心链路:
- 元数据链路:通过 Gravitino REST API(
Accept: application/vnd.gravitino.v1+json)创建/删除 metalake、catalog、schema、fileset、table,见 test_utils.py; - 数据链路:通过 Daft 的
IOConfig(含 Gravitino 配置)把gvfs://路径解析成真实存储位置,见 daft/io/gravitino_filesystem.py。
二、环境准备:一键拉起 Gravitino + MySQL + MinIO
测试依赖一个本地 Gravitino 服务。仓库在 tests/integration/gravitino/docker-compose/docker-compose.yml 中提供了完整的编排文件,包含三个服务:
| 服务 | 镜像 | 端口(默认) | 用途 |
|---|---|---|---|
gravitino | apache/gravitino:1.0.1 | 8090 | Gravitino 服务器(HTTP 端口由GRAVITINO_HTTP_PORT=8090指定) |
mysql | mysql:8.0(linux/amd64) | 3306 | 关系型 Catalog 测试的后端存储(root/root) |
gravitino-minio | quay.io/minio/minio | 9001(映射容器 9000) | S3 兼容对象存储,用于 S3 fileset 测试 |
启动命令(源自 README):
cd tests/integration/gravitino/docker-compose docker compose up -d2.1 端口与凭据的可配置性
docker-compose 中的所有端口均支持通过环境变量覆盖:GRAVITINO_PORT(默认 8090)、MYSQL_PORT(默认 3306)、MINIO_PORT(默认 9001,映射到容器内 9000)。MinIO 的默认访问凭据为minioadmin/minioadmin,与测试中 conftest.py 里的gravitino_minio_io_configfixture 一一对应:
return daft.io.IOConfig( s3=daft.io.S3Config( endpoint_url="http://127.0.0.1:9001", key_id="minioadmin", access_key="minioadmin", region_name="us-east-1", use_ssl=False, ) )注意:测试进程跑在宿主机上,而 Gravitino 跑在容器里,因此两者看到的 MinIO 地址不同——容器内是http://daft-gravitino-minio:9000,宿主机是http://127.0.0.1:9001。测试在创建 catalog 时会先用容器内地址,随后通过update_catalog(REST PUT 的setProperty操作)把s3-endpoint改写成宿主机地址,这正是 test_gravitino_fileset_s3.py 中反复出现update_catalog(... {"@type": "setProperty", "property": "s3-endpoint", "value": "http://127.0.0.1:9001"})的原因。
2.2 entrypoint:为 Gravitino 开启 S3 文件系统支持
默认的 Gravitino 镜像只支持file://与hdfs://存储方案。为了让 S3 fileset 测试可运行,gravitino-entrypoint.sh 在启动前会向 Gravitino 的 fileset 配置文件注入一条关键属性:
FILESET_CONF="/root/gravitino/catalogs/fileset/conf/fileset.conf" if ! grep -q "gravitino.bypass.fs.s3a.path.style.access=true" "$FILESET_CONF" 2>/dev/null; then echo "gravitino.bypass.fs.s3a.path.style.access=true" >> "$FILESET_CONF" fi exec /bin/bash /root/gravitino/bin/start-gravitino.shgravitino.bypass.fs.s3a.path.style.access=true使 Gravitino 内部使用的 Hadoop S3A 文件系统以path-style 访问模式(http://host:9000/bucket/key)访问 MinIO,而不是默认的 virtual-hosted 模式——后者要求 DNS 能解析bucket.host形式的主机名,对本地 MinIO 不适用。脚本在追加配置后仍通过原始start-gravitino.sh启动服务。
三、导出连接配置并运行测试
启动服务后,按 README 的说明导出连接设置(若使用仓库自带 compose 文件,默认值已匹配,无需修改):
export GRAVITINO_ENDPOINT=${GRAVITINO_ENDPOINT:-http://127.0.0.1:8090} export GRAVITINO_METALAKE=${GRAVITINO_METALAKE:-metalake_demo}运行全部 Gravitino 集成测试:
DAFT_RUNNER=native pytest tests/integration/gravitino -v -m integration这里有两个关键点:
DAFT_RUNNER=native指定使用 Daft 的本地原生执行引擎;-m integration只选中打了@pytest.mark.integration()标记的用例。所有 Gravitino 用例都带有该标记(见各测试文件),因为它们必须连接真实服务,无法在无依赖的单元测试环境中运行。
3.1 连接参数与认证环境变量
conftest.py 中以 session 级 fixture 定义了全部连接参数,均可通过环境变量覆盖:
| 环境变量 | 默认值 | 说明 |
|---|---|---|
GRAVITINO_ENDPOINT | http://127.0.0.1:8090 | Gravitino REST 服务地址 |
GRAVITINO_METALAKE | metalake_demo | 使用的 metalake 名称 |
GRAVITINO_AUTH_TYPE | simple | 认证方式(simple/oauth2) |
GRAVITINO_USERNAME | admin(simple 模式下) | 用户名 |
GRAVITINO_PASSWORD | 无 | 密码 |
GRAVITINO_TOKEN | 无 | OAuth2 bearer token |
GRAVITINO_TEST_FILE | 内部示例文件 gvfs 路径 | 指向一个已存在文件的gvfs://URL |
GRAVITINO_TEST_DIR | 内部示例目录 gvfs 路径 | 指向一个已存在 fileset 的gvfs://目录 |
local_gravitino_clientfixture 会把这些参数组装成GravitinoClient(endpoint, metalake_name, auth_type, username, password, token),并据此构造 Daft 侧的两个重要对象:
Catalog.from_gravitino(gravitino_endpoint, gravitino_metalake, ...) # 元数据操作 local_gravitino_client.to_io_config() # gvfs:// 数据读写3.2 健康检查:等待 Gravitino 就绪
由于容器启动需要时间,conftest.py中的_wait_for_gravitino会轮询 Gravitino 的版本接口GET /api/version(携带Accept: application/vnd.gravitino.v1+json),默认最多等待 120 秒、每 3 秒重试一次,直到服务响应或超时抛错。类似地,MySQL 测试通过 TCP 连接 3306 端口的方式等待就绪(test_gravitino_table.py 中的_wait_for_mysql)。
四、测试文件总览
按 README 的划分,套件包含三个测试文件:
| 文件 | 覆盖内容 |
|---|---|
| test_gravitino_fileset.py | 基于本地file://存储的 fileset 读写 |
| test_gravitino_fileset_s3.py | 基于 MinIO(S3 兼容)存储的 fileset 读写 |
| test_gravitino_table.py | Catalog/Table 元数据操作(含 MySQL 关系表) |
4.1 测试基础设施:REST API 封装
test_utils.py 是整套测试的基础设施,它直接调用 Gravitino REST API 完成元数据生命周期管理,全部通过GravitinoClient._session发送请求:
ensure_metalake:GET /api/metalakes/{name}检查,若 404 则POST /api/metalakes创建;create_catalog/update_catalog:POST /api/metalakes/{m}/catalogs创建(类型如FILESET、RELATIONAL),PUT /api/metalakes/{m}/catalogs/{c}执行变更(如setProperty);create_schema/delete_schema:Schema 的创建与级联删除(?cascade=true);create_fileset/delete_fileset:fileset 的创建与删除。创建时通过storageLocations声明存储位置(见下文 4.2);delete_catalog:先PATCH将 catalog 的inUse置为false,再DELETE,规避 Gravitino 对使用中 catalog 的删除限制。
4.2 gvfs 路径与存储位置解析
test_gravitino_fileset.py中的_resolve_storage_uri展示了 Daft 解析gvfs://的核心逻辑:
parsed = urlparse(gvfs_path) # 必须为 gvfs://fileset/... catalog, schema, fileset, *rest = segments # 取前三个路径段 fileset_obj = client.load_fileset(f"{catalog}.{schema}.{fileset}") storage_uri = fileset_obj.fileset_info.storage_location.rstrip("/") if rest: storage_uri = f"{storage_uri}/{'/'.join(rest)}"即:gvfs 路径被拆成catalog/schema/fileset/剩余路径四部分,前三段用于在 Gravitino 中加载 fileset 元数据、取出真实存储位置,剩余路径则拼接到存储 URI 之后。本地测试创建的 fileset 存储 URI 是tmp_path的file://URI;S3 测试则使用s3a://<bucket>/<prefix>/。写入 fileset 元数据时,仓库同时设置properties["location"]与storageLocations({"default": storage_uri}),以兼容 Gravitino 1.0+ 的多存储位置格式。
4.3 本地 file:// fileset 测试
test_gravitino_fileset.py 的prepared_filesetfixture 完整演示了"建数据 → 建元数据 → 测功能 → 全量清理"的闭环:
- 用
daft.from_pydict({"id": [1,2,3], "value": [...]}).write_parquet()在临时目录写入 sample.parquet; - 通过
ensure_metalake+create_catalog+create_schema+create_fileset创建唯一的 catalog/schema/fileset; - 暴露
gvfs_root = gvfs://fileset/{catalog}/{schema}/{fileset}; finally中按 fileset → schema → catalog 顺序删除元数据,并shutil.rmtree清理本地数据。
基于该 fixture 的四个用例:
test_read_fileset_over_gvfs:daft.read_parquet(gvfs_file, io_config=gravitino_io_config)直接读取 fileset 内的 parquet,排序后断言与写入数据一致;test_list_files_via_glob:用glob_path_with_stats(f"{gvfs_root}/**/*.parquet", FileFormat.Parquet, io_config)验证 glob 列出的是真实存在的 parquet 文件(注意write_parquet会生成目录结构,因此文件前缀包含sample.parquet/);test_from_glob_path_reads_files:daft.from_glob_path(glob_pattern, io_config=...)返回文件路径 DataFrame;test_delete_file_via_gvfs_path:先经_resolve_storage_uri把 gvfs 路径还原成本地路径删除,随后断言 glob 结果为空、且直接读该文件会抛异常——即"数据删除后元数据感知一致"。
4.4 S3 fileset 测试:从 MinIO 到外部 Gravitino
test_gravitino_fileset_s3.py 是覆盖面最广的文件,它使用s3fs直接操作 MinIO 进行 bucket 创建与清理(s3_bucketfixture),并创建带 S3 文件系统配置的 catalog:
create_catalog( client, metalake, catalog_name, properties={ "filesystem-providers": "s3", "s3-endpoint": "http://daft-gravitino-minio:9000", # 容器内地址 "s3-access-key-id": "minioadmin", "s3-secret-access-key": "minioadmin", }, ) # 随后 update_catalog 将 s3-endpoint 改为 http://127.0.0.1:9001(宿主机地址)按 README 的总结,S3 套件覆盖:
- 基础读取:
test_read_s3_fileset_over_gvfs通过gvfs://读 S3 上的 parquet; - glob 与文件列举:
test_list_s3_fileset_via_glob、test_from_glob_path_s3_reads_files; - 分区数据处理:
test_s3_fileset_partitioned_data用partition_cols=["category"]写分区 parquet,再以daft.read_parquet(glob, ..., hive_partitioning=True)读回全部 6 行并校验分区列; - 错误处理:
test_s3_fileset_missing_file断言读不存在的文件抛异常;test_s3_fileset_empty_glob断言对空目录 glob 收集后为空 DataFrame; - 写入能力:
test_write_parquet_to_gvfs/test_write_csv_to_gvfs/test_write_json_to_gvfs分别向gvfs://路径写 parquet/csv/json,并用 s3fs 和直读s3://路径双重校验; - 完整回环:
test_write_and_read_roundtrip_gvfs经gvfs://写入再经gvfs://读回。
这些写入用例依赖于 daft/io/gravitino_filesystem.py 提供的 PyArrow 文件系统适配(GravitinoFSHandler、输出流实现等),它让 PyArrow 的 parquet 写入能够把字节流定向到gvfs://目标。
4.4.1 对接外部 Gravitino 部署
README 还给出了"不依赖本地 docker-compose、直连已有 S3 fileset 部署"的用法——前提是对方 Gravitino 已配置好 S3 fileset:
# 指向一个已存在的 S3-backed fileset export GRAVITINO_TEST_DIR="gvfs://fileset/<catalog>/<schema>/<fileset>/" export GRAVITINO_TEST_FILE="gvfs://fileset/<catalog>/<schema>/<fileset>/file.parquet" export GRAVITINO_ENDPOINT="http://your-gravitino-server:8090" # 只运行外部 fileset 用例 DAFT_RUNNER=native pytest tests/integration/gravitino/test_gravitino_fileset_s3.py::test_read_existing_s3_fileset -v -m integration对应地,test_read_existing_s3_fileset与test_read_specific_s3_file带有skipif条件:仅当GRAVITINO_TEST_DIR/GRAVITINO_TEST_FILE以gvfs://开头时才执行,否则跳过。这保证了在没有外部部署时套件不会报错。
4.5 Catalog/Table 元数据测试(含 MySQL)
test_gravitino_table.py 验证 Daft Catalog API 与 Gravitino 的对接:
test_catalog_from_gravitino:Catalog.from_gravitino(endpoint, metalake)返回非空 Catalog,且其name为gravitino_{metalake}(对应 daft/catalog/__gravitino/_catalog.py 中GravitinoCatalog.name的实现);test_catalog_has_table_false:不存在的表返回False;test_catalog_list_tables_returns_identifiers:list_tables()返回标识符列表;test_catalog_get_table_not_found:获取不存在的表抛出 Daft 的NotFoundError。
mysql_gravitino_catalogfixture 则更进一步,通过 Gravitino REST API 创建了一个jdbc-mysql类型的 catalog(jdbc-url: jdbc:mysql://mysql:3306、root/root),并在两个 schema 下分别建了users、orders、products三张表(integer/varchar/decimal 类型)。test_gravitino_mysql_integration最终验证:
catalog.list_tables()能列出{catalog}.{schema}.{table}形式的三级标识符;catalog.has_table(...)对真实存在的表返回True、对不存在的表返回False。
五、从测试看连接器实现要点
5.1 Catalog 适配层
GravitinoCatalog(daft/catalog/__gravitino/_catalog.py)实现了 Daft 的Catalog抽象:
_get_table:调用self._inner.load_table(str(ident)),GravitinoTableNotFoundError被转换为 Daft 统一的NotFoundError;_list_namespaces/_list_tables:支持按catalog、catalog.schema粒度的 pattern 过滤,无 pattern 时遍历所有 catalog 与 namespace 聚合返回;- 文件顶部明确标注:这些是内部 API,请使用
load_gravitino()或Catalog.from_gravitino()(见 daft/catalog/__gravitino/init.py 的导出)。
5.2 gvfs 数据面实现
daft/io/gravitino_filesystem.py是 gvfs 协议的数据面核心:GravitinoFSHandler提供 FSSpec 风格的接口(ls、open、mkdir等),并显式声明"流式直读gvfs://尚未实现,读取请用daft.read_parquet()等高层 API",而输出流则借助io_put把写入字节定向到gvfs://路径。这正是测试中"读用daft.read_parquet+IOConfig、写用df.write_*"模式的底层原因。
5.3 清理策略:测试不留下任何垃圾
两套 fileset 测试与 MySQL 测试都在finally中执行清理:先删 fileset,再删 schema(可级联),再删 catalog,最后删除本地/S3 上的实际数据(shutil.rmtree或 s3fs 递归删除 bucket)。README 也明确说明:"The tests clean up both the server-side metadata and storage when finished."——这是该套件可反复运行、无状态残留的关键设计。
六、可选配置与运行前提
6.1 用真实数据验证 gvfs IO
README 指出:可选地配置GRAVITINO_TEST_FILE与GRAVITINO_TEST_DIR,指向你 Gravitino 部署中真实存在的 fileset,即可让 gvfs IO 测试读取具体数据(而不是默认的示例路径)。默认值(gvfs://fileset/s3_fileset_catalog3/test_schema/test_fileset/...)仅当本地测试创建过同名 fileset 时才有效。
6.2 运行前提汇总
- Docker + Docker Compose(用于本地三件套),或一个可访问的 Gravitino 部署(0.9.0+,见 docs/connectors/gravitino.md 的 Requirements);
- 安装了带 Gravitino 支持的 Daft:
pip install "daft[gravitino]"; - Python 依赖
requests(conftest 与 test_utils 直接使用)与s3fs(S3 测试使用); - 连接器 API 处于 Beta 阶段,接口可能随版本演进而变化(docs/connectors/gravitino.md 中有明确警告)。
6.3 已知限制
从测试与文档可以确认的当前边界(见 docs/connectors/gravitino.md 的 Limitations 一节):
- 凭据签发(credential vending)尚未实现;
- gvfs 写入目前只支持 S3-backed fileset(其他存储后端仍在规划中);
- 连接器直接调用 Gravitino REST API,而非 Gravitino 官方 Python client;
GravitinoCatalog的create_table/drop_table/create_namespace等方法当前抛出NotImplementedError。
七、快速验证清单
cd tests/integration/gravitino/docker-compose && docker compose up -d启动三件套;- 确认 entrypoint 已向
fileset.conf写入gravitino.bypass.fs.s3a.path.style.access=true(容器日志或docker exec查看); - 导出
GRAVITINO_ENDPOINT/GRAVITINO_METALAKE(默认即可); - 运行
DAFT_RUNNER=native pytest tests/integration/gravitino -v -m integration; - 观察通过项:本地 fileset 读写、S3 fileset 读写/glob/分区/错误处理/写入回环、MySQL 表的元数据列举。
这套测试既是连接器的回归保障,也是一份"如何在 Daft 中使用 Gravitino"的可运行范例:从Catalog.from_gravitino()到gvfs://fileset/<catalog>/<schema>/<fileset>/...的读写,从 MinIO path-style 配置到 REST API 的元数据生命周期管理,都可以直接迁移到你的生产集成方案中。
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考