DataHub Metadata File 接入指南:file Source 的配置、文件格式与状态化删除检测详解
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
导读
Metadata File(元数据文件)是 DataHub 元数据摄取体系中的核心回放机制:它可以作为存储与湖仓场景下元数据的载体,也可以作为通用的"元数据交换文件"用于解耦元数据采集与元数据入库。本文基于 DataHub 官方文档 metadata-file 源集成说明,结合仓库源码深入讲解file类型 Source 的完整配置、MCE/MCP 文件格式、读取模式、连接测试与状态化删除检测,帮助读者在acryl-datahub环境中独立完成"元数据文件生成 → 回放 → 校验"的端到端工作流。
一、Metadata File 集成概览
根据官方文档,Metadata File 是一个存储与湖仓(storage and lakehouse)平台,而 DataHub 与 Metadata File 的集成覆盖了文件/湖仓元数据实体,例如数据集(datasets)、路径(paths)与容器(containers),并且支持状态化删除检测(stateful deletion detection)。
在仓库实现层面,这一集成对应metadata-ingestion中的通用文件回放 Source——type: file。其源码位于 metadata-ingestion/src/datahub/ingestion/source/file.py,类名GenericFileSource,其类注释明确说明:
This plugin pulls metadata from a previously generated file. The metadata file sink can produce such files, and a number of samples are included in the examples/mce_files directory.
也就是说,该模块用于将先前生成好的元数据文件(MCE / MCP JSON)回放到 DataHub,既可作为生产环境摄取工作流的一环,也是调试元数据管线的利器。配合 Metadata File Sink 使用时,可以形成"从任意数据源采集 → 落盘为元数据文件 → 再回放进入 DataHub"的解耦架构。
二、概念映射:Source 概念与 DataHub 概念
官方文档给出了集成层面向的通用概念映射关系(特定场景下的精确映射仍在完善中),完整继承如下:
| Source 概念 | DataHub 概念 | 说明 |
|---|---|---|
| Platform/account/project scope | Platform Instance, Container | 在平台上下文内组织资产。 |
| Core technical asset(如表/视图/主题/文件) | Dataset | 主要摄入的技术资产。 |
| Schema fields / columns | SchemaField | 在支持 Schema 抽取时包含。 |
| Ownership and collaboration principals | CorpUser, CorpGroup | 由支持所有权与身份元数据的模块发出。 |
| Dependencies and processing relationships | Lineage edges | 在支持且启用血缘抽取时可用。 |
这张表的核心含义是:无论元数据来源于哪个平台,DataHub 最终都统一收敛到 Dataset、SchemaField、CorpUser/CorpGroup、Container、Platform Instance 等标准实体与关系上,从而保证不同数据源在 DataHub 中的一致视图。例如仓库自带的示例文件 single_mce.json 中就包含一个CorpUserSnapshot,其urn为urn:li:corpuser:harshal,携带CorpUserInfo方面数据,正是"所有权主体 → CorpUser"映射的具体落地。
三、快速开始:最小可用 Recipe
官方在 file_recipe.yml 中给出了最小配置骨架:
source: type: file config: # Coordinates path: ./path/to/mce/file.json sink: # sink configs其中type: file指定使用GenericFileSource,path指向待回放的元数据文件。若要组成完整的"采集 → 落盘 → 回放"链路,可参照 Metadata File Sink 的 Quickstart recipe,由任意 Source 采集后写入文件:
source: # source configs sink: type: file config: filename: ./path/to/mce/file.json注意:Sink 侧早期配置字段为
filename;而 Source 侧的filename已被标记为废弃(deprecated),统一使用path,源码中通过pydantic_renamed_field("filename", "path", print_warning=False)做了自动兼容迁移(见 file.py)。
四、核心配置项详解
fileSource 的全部配置集中在FileSourceConfig(继承自StatefulIngestionConfigBase),定义于 metadata-ingestion/src/datahub/ingestion/source/file.py。各字段含义如下:
| 字段 | 必填 | 默认值 | 说明 |
|---|---|---|---|
path | ✅ | 无 | 待摄入的文件或目录路径,也可以是远程文件 URL。若指向目录,则处理该目录下所有扩展名匹配file_extension(默认.json)的文件。 |
file_extension | 否 | .json | 指向目录时用于过滤待处理文件的扩展名;特殊值*表示处理目录下所有文件(不区分扩展名)。配置时会自动补全前导点(add_leading_dot_to_extension校验器)。 |
read_mode | 否 | AUTO | 文件读取模式,枚举值为STREAM/BATCH/AUTO,详见第五节。 |
aspect | 否 | 无 | 若设置,只读取该 aspect 对应的元数据(例如只回放ownership方面),用于定向回放。 |
count_all_before_starting | 否 | true | 启用时在开始前先统计文件总记录数,用于精确估算完成时间;若启动阶段耗时过长可关闭。 |
stateful_ingestion | 否 | 无 | 状态化摄取配置,类型为StatefulStaleMetadataRemovalConfig,用于启用删除检测,详见第七节。 |
filename(废弃) | 否 | 无 | 已废弃,自动迁移为path。 |
从源码实现看,path的解析并不局限于本地文件系统:get_filenames()会先通过get_path_schema(path_str)解析路径协议,再从fs_registry获取对应的文件系统实现(见 file.py)。这意味着只要注册表中存在对应协议,即可直接摄入远程文件(如对象存储等),这是该模块面向存储/湖仓场景的重要支撑。
aspect过滤逻辑在get_workunits_internal()中实现:当对象为 MCP/MCPW 且aspectName与配置不一致时直接跳过(见 file.py),适合只回放某类方面的增量场景。
五、支持的文件格式:MCE 与 MCP
GenericFileSource读取的是 DataHub 标准的元数据事件 JSON 文件,支持两种主流对象,判定逻辑位于_from_obj_for_file()(见 file.py):
- MCE(MetadataChangeEvent):JSON 对象中包含
proposedSnapshot键,对应实体快照,例如single_mce.json中的CorpUserSnapshot; - MCP / MCPW(MetadataChangeProposal / Wrapper):JSON 对象中包含
aspect键,对应单方面更新; - UsageAggregationClass:包含
bucket键,属于已废弃的用量聚合格式,读取到时会打印 warning 并丢弃(Dropping deprecated UsageAggregationClass)。
每个对象读取后会执行item.validate()校验,校验失败会抛出ValueError(f"Failed to parse: {obj}")。
文件内部结构支持两种组织方式(_iterate_file_batch中的兼容逻辑,见 file.py):
- JSON 数组:顶层为对象列表,逐条回放,如 mce_list.json 中包含两个
CorpUserSnapshot; - 单个 JSON 对象:兼容旧版单对象格式,直接回放。
仓库 examples/mce_files 目录提供了大量可直接试用的样例,除用户快照外还包括容器(test_containers.json)、域(test_domains.json)、结构化属性(test_structured_properties.json)、标签策略(associate_tags_policies_test.json)等,是验证fileSource 行为的最佳素材。
六、读取模式与性能:AUTO / STREAM / BATCH
read_mode控制文件的读取策略,源码中枚举定义于FileReadMode(file.py):
| 模式 | 行为 |
|---|---|
BATCH | 一次性json.load整个文件后逐条 yield,适合中小文件。 |
STREAM | 使用ijson增量解析流式读取,内存占用低,适合超大文件。 |
AUTO | 自动决策:文件大小小于 100MB(_minsize_for_streaming_mode_in_bytes = 100 * 1000 * 1000,见 file.py)时走 BATCH,否则走 STREAM。 |
决策逻辑位于_iterate_file()(file.py):AUTO 模式按文件字节数切换模式,STREAM 模式读取顶层item。流式模式下若开启count_all_before_starting,会先完整扫描一遍统计元素总数(用于进度百分比),随后fp.seek(0)回到文件头开始正式读取(file.py)。
配合这一机制,FileSourceReport(继承自StaleEntityRemovalSourceReport)会记录每个文件的字节数、元素数、已读字节/元素数、解析/计数/反序列化耗时,并实时计算完成百分比与预计剩余时间(compute_stats,见 file.py)。在大文件批量回放场景下,这些统计信息是监控进度与定位性能瓶颈的直接依据。
七、状态化摄取与删除检测
官方文档明确该集成支持状态化删除检测。在实现上,GenericFileSource继承自StatefulIngestionSourceBase,其允许启用的 Workunit 处理器为(file.py):
AutoWorkunitsReporterProcessor:自动汇报 Workunit 进度;AutoStaleEntityRemovalProcessor:自动陈旧实体删除处理器,即删除检测的核心;- 当配置了
stateful_ingestion时还会插入AutoStatusAspectProcessor,用于维护实体的 status 方面。
使用方式是在 recipe 的source.config下增加:
source: type: file config: path: ./path/to/mce/file.json stateful_ingestion: remove_stale_metadata: true启用后,每次回放会以本次文件内容为基准,对比状态化检查点,将上次存在但本次文件中已不存在的实体识别为"陈旧"并触发删除,从而保证 DataHub 中的元数据与源文件保持一致,避免残留过时资产。这是元数据文件回放链路中维持数据一致性的关键能力,与第五节中的aspect定向回放组合使用,可实现细粒度的增量维护。
八、连接测试(Test Connection)
fileSource 声明了SourceCapability.TEST_CONNECTION能力且"默认启用"(@capability(SourceCapability.TEST_CONNECTION, "Enabled by default")),支持状态为 GA(@support_status(SupportStatus.GA))。
其实现为静态方法test_connection(file.py),校验逻辑包括:
- 路径是否存在(不存在则报告
doesn't appear to be a valid file or directory); - 是否为目录且是否可读(
os.access(path, os.R_OK)); - 若为目录,还需具备执行权限(
os.X_OK)。
因此在运行摄取前,可以借助datahub ingest的--test-connection或 UI 中的测试连接功能,快速验证path的可读性,这与官方前置条件文档 file_pre.md 中"确保对源元数据 API 具备读取权限"的要求一致。
九、故障排查
官方文档 file_post.md 给出了通用的排查顺序:先校验凭据、权限、连通性与范围过滤,再查看摄取日志中的源相关错误并调整配置。
结合源码,fileSource 最常见的失败点与排查方向如下:
- 反序列化失败:
iterate_generic_file捕获每条记录的反序列化异常,通过self.report.failure(message="Failed to deserialize metadata", context=f"{file_status.path}-{i}", ...)记录具体文件与记录序号(file.py),排查时应优先检查该文件的 JSON 结构是否符合 MCE/MCP 格式(缺少proposedSnapshot/aspect键会直接抛Unknown object type); - 路径/权限问题:参考第八节,
test_connection会给出具体的不可读、无执行权限提示; - 启动过慢:
count_all_before_starting对超大文件会多一次全量扫描,若启动时间不可接受可将其设为false; - 扩展名未匹配:指向目录时,确保文件扩展名与
file_extension(默认.json)一致,或显式设为*。
十、生成与回放:构建解耦的元数据工作流
综合上述内容,fileSource 的典型生产用法是构建解耦式摄取管线:
- 任意数据源(数据库、消息队列、湖仓等)通过各自的 Source 采集元数据;
- 使用 Metadata File Sink(
sink.type: file)将结果落盘为 MCE/MCP JSON 文件,实现"采集"与"入库"的解耦,也便于人工检查与调试中间产物; - 在目标环境(如生产 DataHub 实例)使用
type: file的 Source 回放该文件,配合stateful_ingestion做删除检测; - 通过
aspect定向回放实现精细化更新,借助test_connection与FileSourceReport的进度统计验证任务健康度。
该模式在跨网络隔离环境迁移元数据、批量重放历史元数据、以及 CI 中做摄取回归测试等场景下尤为实用。仓库中 examples/mce_files 下的样例文件即可作为回放练习与单元验证的直接素材。
延伸阅读
- 源集成文档目录:metadata-ingestion/docs/sources/metadata-file(含 README.md、file_pre.md、file_post.md 与 file_recipe.yml)
- Source 核心实现:metadata-ingestion/src/datahub/ingestion/source/file.py
- 配套 Sink 文档:metadata-ingestion/sink_docs/metadata-file.md
- 示例元数据文件:metadata-ingestion/examples/mce_files
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考