Daft × Apache Gravitino:端到端集成测试实战指南(gvfs:// 文件集与 Catalog 目录操作)

发布时间:2026/9/17 21:29:06
Daft × Apache Gravitino:端到端集成测试实战指南(gvfs:// 文件集与 Catalog 目录操作)
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 APIAccept: application/vnd.gravitino.v1json创建/删除 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 中提供了完整的编排文件包含三个服务服务镜像端口默认用途gravitinoapache/gravitino:1.0.18090Gravitino 服务器HTTP 端口由GRAVITINO_HTTP_PORT8090指定mysqlmysql:8.0linux/amd643306关系型 Catalog 测试的后端存储root/rootgravitino-minioquay.io/minio/minio9001映射容器 9000S3 兼容对象存储用于 S3 fileset 测试启动命令源自 READMEcd 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( s3daft.io.S3Config( endpoint_urlhttp://127.0.0.1:9001, key_idminioadmin, access_keyminioadmin, region_nameus-east-1, use_sslFalse, ) )注意测试进程跑在宿主机上而 Gravitino 跑在容器里因此两者看到的 MinIO 地址不同——容器内是http://daft-gravitino-minio:9000宿主机是http://127.0.0.1:9001。测试在创建 catalog 时会先用容器内地址随后通过update_catalogREST 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.accesstrue $FILESET_CONF 2/dev/null; then echo gravitino.bypass.fs.s3a.path.style.accesstrue $FILESET_CONF fi exec /bin/bash /root/gravitino/bin/start-gravitino.shgravitino.bypass.fs.s3a.path.style.accesstrue使 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_RUNNERnative pytest tests/integration/gravitino -v -m integration这里有两个关键点DAFT_RUNNERnative指定使用 Daft 的本地原生执行引擎-m integration只选中打了pytest.mark.integration()标记的用例。所有 Gravitino 用例都带有该标记见各测试文件因为它们必须连接真实服务无法在无依赖的单元测试环境中运行。3.1 连接参数与认证环境变量conftest.py 中以 session 级 fixture 定义了全部连接参数均可通过环境变量覆盖环境变量默认值说明GRAVITINO_ENDPOINThttp://127.0.0.1:8090Gravitino REST 服务地址GRAVITINO_METALAKEmetalake_demo使用的 metalake 名称GRAVITINO_AUTH_TYPEsimple认证方式simple/oauth2GRAVITINO_USERNAMEadminsimple 模式下用户名GRAVITINO_PASSWORD无密码GRAVITINO_TOKEN无OAuth2 bearer tokenGRAVITINO_TEST_FILE内部示例文件 gvfs 路径指向一个已存在文件的gvfs://URLGRAVITINO_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.v1json默认最多等待 120 秒、每 3 秒重试一次直到服务响应或超时抛错。类似地MySQL 测试通过 TCP 连接 3306 端口的方式等待就绪test_gravitino_table.py 中的_wait_for_mysql。四、测试文件总览按 README 的划分套件包含三个测试文件文件覆盖内容test_gravitino_fileset.py基于本地file://存储的 fileset 读写test_gravitino_fileset_s3.py基于 MinIOS3 兼容存储的 fileset 读写test_gravitino_table.pyCatalog/Table 元数据操作含 MySQL 关系表4.1 测试基础设施REST API 封装test_utils.py 是整套测试的基础设施它直接调用 Gravitino REST API 完成元数据生命周期管理全部通过GravitinoClient._session发送请求ensure_metalakeGET /api/metalakes/{name}检查若 404 则POST /api/metalakes创建create_catalog/update_catalogPOST /api/metalakes/{m}/catalogs创建类型如FILESET、RELATIONALPUT /api/metalakes/{m}/catalogs/{c}执行变更如setPropertycreate_schema/delete_schemaSchema 的创建与级联删除?cascadetruecreate_fileset/delete_filesetfileset 的创建与删除。创建时通过storageLocations声明存储位置见下文 4.2delete_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://URIS3 测试则使用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_metalakecreate_catalogcreate_schemacreate_fileset创建唯一的 catalog/schema/fileset暴露gvfs_root gvfs://fileset/{catalog}/{schema}/{fileset}finally中按 fileset → schema → catalog 顺序删除元数据并shutil.rmtree清理本地数据。基于该 fixture 的四个用例test_read_fileset_over_gvfsdaft.read_parquet(gvfs_file, io_configgravitino_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_filesdaft.from_glob_path(glob_pattern, io_config...)返回文件路径 DataFrametest_delete_file_via_gvfs_path先经_resolve_storage_uri把 gvfs 路径还原成本地路径删除随后断言 glob 结果为空、且直接读该文件会抛异常——即数据删除后元数据感知一致。4.4 S3 fileset 测试从 MinIO 到外部 Gravitinotest_gravitino_fileset_s3.py 是覆盖面最广的文件它使用s3fs直接操作 MinIO 进行 bucket 创建与清理s3_bucketfixture并创建带 S3 文件系统配置的 catalogcreate_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 上的 parquetglob 与文件列举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_partitioningTrue)读回全部 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_DIRgvfs://fileset/catalog/schema/fileset/ export GRAVITINO_TEST_FILEgvfs://fileset/catalog/schema/fileset/file.parquet export GRAVITINO_ENDPOINThttp://your-gravitino-server:8090 # 只运行外部 fileset 用例 DAFT_RUNNERnative 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 元数据测试含 MySQLtest_gravitino_table.py 验证 Daft Catalog API 与 Gravitino 的对接test_catalog_from_gravitinoCatalog.from_gravitino(endpoint, metalake)返回非空 Catalog且其name为gravitino_{metalake}对应 daft/catalog/__gravitino/_catalog.py 中GravitinoCatalog.name的实现test_catalog_has_table_false不存在的表返回Falsetest_catalog_list_tables_returns_identifierslist_tables()返回标识符列表test_catalog_get_table_not_found获取不存在的表抛出 Daft 的NotFoundError。mysql_gravitino_catalogfixture 则更进一步通过 Gravitino REST API 创建了一个jdbc-mysql类型的 catalogjdbc-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 适配层GravitinoCatalogdaft/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_parquetIOConfig、写用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 IOREADME 指出可选地配置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 支持的 Daftpip install daft[gravitino]Python 依赖requestsconftest 与 test_utils 直接使用与s3fsS3 测试使用连接器 API 处于 Beta 阶段接口可能随版本演进而变化docs/connectors/gravitino.md 中有明确警告。6.3 已知限制从测试与文档可以确认的当前边界见 docs/connectors/gravitino.md 的 Limitations 一节凭据签发credential vending尚未实现gvfs 写入目前只支持 S3-backed fileset其他存储后端仍在规划中连接器直接调用 Gravitino REST API而非 Gravitino 官方 Python clientGravitinoCatalog的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.accesstrue容器日志或docker exec查看导出GRAVITINO_ENDPOINT/GRAVITINO_METALAKE默认即可运行DAFT_RUNNERnative 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),仅供参考