DataHub MLflow 数据源接入指南:模型注册、实验、运行与血缘的元数据采集实战
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
导读
本文围绕 DataHub 的 MLflow 元数据接入模块(metadata-ingestion中的mlflowsource)展开,系统讲解如何将 MLflow 平台中的 Registered Model、Model Version、Experiment、Run 以及 run-to-model 血缘采集到 DataHub,并覆盖 Tags 与有状态删除检测等能力。读完本文,你将掌握 MLflow 与 DataHub 实体之间的完整映射关系、全部配置项的含义与取值、认证与数据集血缘配置方法,并能直接编写可运行的 ingestion recipe 完成一次生产级采集。
MLflow 接入模块概述
MLflow 是一个机器学习平台,DataHub 为它提供了官方的元数据接入模块(source 类型为mlflow)。该模块覆盖以下核心对象:
- Registered Model(注册模型):模型在 MLflow Model Registry 中的容器,用于承载同一模型的多个版本;
- Model Version(模型版本):模型的每一次具体迭代,拥有独立的 artifacts 与元数据;
- Experiment(实验):组织相关 Run 的逻辑分组,可记录参数、指标与 artifacts;
- Run(运行):一次具体执行,包含执行细节、参数、指标以及与模型的血缘关系;
- run-to-model 血缘:Run 到模型版本的产出关系;
- Tags:将 Model Stage 映射为 DataHub Tag;
- 有状态删除检测(stateful deletion detection):基于
StatefulIngestionSourceBase实现,用于清理已删除实体的陈旧元数据。
模块的声明信息可以在 mlflow.py 中看到:平台名platform = "mlflow",支持状态为BETA(SupportStatus.BETA),并声明了 Descriptions、Containers(MLflow Experiment)与 Tags 三项能力。
两个重要边界说明
原文档明确指出以下两点限制,采集前需要知晓:
- MLflow 特性不会以 DataHub
MlFeature实体采集:MLflow 的 tracking API 不记录 feature-to-column 的溯源信息,连接器没有可供读取的数据,因此该 aspect 无法填充; - 模型不会在模型级别关联其训练数据集(
mlModelTrainingData不会被填充),但 Run 级别的数据集血缘会被采集——详见下文概念映射表中的Dataset Input行。
概念映射:MLflow 对象如何落到 DataHub 实体
这是整个连接器设计的核心。下表完整对应原文档的 Concept Mapping 表,并补充了来自源码的 URN 构造细节:
| MLflow 源概念 | DataHub 实体 | 说明 |
|---|---|---|
| Registered Model | MlModelGroup | Model Group 的名称与 Registered Model 的名称相同(如my_mlflow_model)。Registered Model 在 MLflow 中是同一模型多个版本的容器。对应 URN 通过make_ml_model_group_urn生成(见 mlflow.py) |
| Model Version | MlModel | 模型名称为{registered_model_name}{model_name_separator}{model_version}(例如 Registered Model 为my_mlflow_model、Version 为 1 时,模型名是my_mlflow_model_1,后续为my_mlflow_model_2……)。每个 Model Version 代表模型的一次特定迭代,拥有自己的 artifacts 与元数据。分隔符由配置项model_name_separator控制,默认_(见 mlflow.py) |
| Experiment | Container | MLflow 中的每个 Experiment 映射为 DataHub 的一个 Container,其 subtype 为MLFLOW_EXPERIMENT,用于组织相关 Run 并记录参数、指标与 artifacts(见 mlflow.py) |
| Run | DataProcessInstance | 捕获 Run 的执行细节、参数、指标以及与模型的血缘。Run 同时带有MLFLOW_TRAINING_RUNsubtype、MLTrainingRunPropertiesClass(超参数 + 训练指标 + artifacts URI)与DataProcessInstanceRunEventClass(运行状态、耗时)(见 mlflow.py) |
| Model Stage | Tag | Model Stage 与 Tag 的映射关系:Production →mlflow_production,Staging →mlflow_staging,Archived →mlflow_archived,None →mlflow_none。Model Stage 表示每个版本所处的部署状态。Tag 名称由mlflow_+ stage 小写构成,且带颜色标记(见 mlflow.py 与_make_stage_tag_name) |
| Dataset Input | DataProcessInstanceInput | 通过mlflow.log_input()记录,将 Run 与其训练数据集关联。需要启用materialize_dataset_inputs才会同时创建被引用的数据集实体(见 mlflow.py) |
Stage 到 Tag 的源码级细节
源码中内置了四种 Stage 的描述与十六进制颜色(mlflow.py):
| Stage | 生成的 Tag | 描述 | 颜色 |
|---|---|---|---|
| Production | mlflow_production | Production Stage for an ML model in MLflow Model Registry | #308613 |
| Staging | mlflow_staging | Staging Stage for an ML model in MLflow Model Registry | #FACB66 |
| Archived | mlflow_archived | Archived Stage for an ML model in MLflow Model Registry | #5D7283 |
| None | mlflow_none | None Stage for an ML model in MLflow Model Registry | #F2F4F5 |
单元测试test_stages(test_mlflow_source.py)验证了恰好生成 4 个 stage tag,且命名符合mlflow_+ stage 小写的规则。
版本集(Version Set)与别名
从源码看,每个 Registered Model 还会生成一个VersionSetUrn(由 platform 与模型名经datahub_guid计算),每个 Model Version 通过VersionPropertiesClass关联到该版本集,并携带sortId(版本号 10 位补零)与aliases(来自model_version.aliases)。这为 DataHub 的模型版本管理与别名查询提供了基础(mlflow.py)。
环境要求与前置条件
来自 mlflow_pre.md 与 mlflow_post.md 的官方要求:
- 网络连通性:确保 ingestion 执行环境能够访问 MLflow 服务;
- 认证凭据:提供有效的认证信息(用户名/密码);
- 只读权限:拥有该模块所需元数据 API 的读取权限;
- 版本兼容性:需要MLflow 服务端 1.28.0 或更高版本。如果使用更早的版本,Experiments 与 Runs 的采集将被跳过。
版本兼容性的源码依据位于_traverse_mlflow_search_func(mlflow.py):当 MLflow 返回ENDPOINT_NOT_FOUND错误码时,连接器会记录一条 warning,提示“请升级到 1.28.0 或更高版本以确保兼容性,跳过 Experiments 与 Runs 的采集”,并继续其余采集流程。
配置文件(Recipe)与全部配置项
最小可用的 recipe 如下(源自 mlflow_recipe.yml):
source: type: mlflow config: # Coordinates tracking_uri: tracking_uri sink: # sink configs完整配置项定义在MLflowConfig(mlflow.py),逐项说明如下:
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
tracking_uri | str | None | MLflow Tracking Server URI。未设置时使用 MLflow 默认 tracking_uri(本地mlruns/目录或MLFLOW_TRACKING_URI环境变量) |
registry_uri | str | None | Model Registry Server URI。未设置时使用 MLflow 默认 registry_uri(取 tracking_uri 的值或MLFLOW_REGISTRY_URI环境变量) |
model_name_separator | str | _ | 分隔模型名与版本号的字符串,如model_1或model-1 |
base_external_url | str | None | 构造指向 MLflow 的外部 URL 时使用的基础 URL。未设置时,若 tracking_uri 是 HTTP URL 则使用 tracking_uri;两者都不可用时,不生成外部 URL |
materialize_dataset_inputs | bool | False | 是否物化(创建)每个 Run 的数据集输入实体 |
source_mapping_to_platform | dict | None | 将 MLflow 数据集 source type 映射为 DataHub platform 的映射表 |
username | str | None | MLflow 认证用户名 |
password | TransparentSecretStr | None | MLflow 认证密码(敏感信息自动脱敏) |
stateful_ingestion | StatefulStaleMetadataRemovalConfig | None | 有状态陈旧元数据删除配置 |
认证配置
可以在 recipe 中直接配置username与password(mlflow_post.md 中的示例):
source: type: mlflow config: tracking_uri: "http://127.0.0.1:5000" username: <username> password: <password>源码中的认证逻辑(_configure_client,见 mlflow.py)值得注意:
- 如果
username与password只设置了其中一个,会直接抛出ValueError("Both username and password must be set together"); - 两者都设置时,会分别写入环境变量
MLFLOW_TRACKING_USERNAME与MLFLOW_TRACKING_PASSWORD,再构造MlflowClient(tracking_uri, registry_uri)。
数据集血缘与平台映射
通过source_mapping_to_platform可以将不同 MLflow 引擎(source type)的数据集关联到指定的 DataHub 平台(mlflow_post.md 示例):
source_mapping_to_platform: huggingface: snowflake # Maps Hugging Face datasets to Snowflake platform http: s3 # Maps HTTP data sources to s3 platform默认行为:仅按平台与名称链接到已存在的数据集,不创建新数据集。要自动创建数据集,需要启用materialize_dataset_inputs:
materialize_dataset_inputs: true # Creates new datasets if they don't exist两个配置项可以独立组合使用:
# Only map to existing datasets materialize_dataset_inputs: false source_mapping_to_platform: huggingface: snowflake # Maps Hugging Face datasets to Snowflake platform pytorch: snowflake # Maps PyTorch datasets to Snowflake platform # Create new datasets and map platforms materialize_dataset_inputs: true source_mapping_to_platform: huggingface: snowflake pytorch: snowflake注意:原文档中存在
materlize_dataset_inputs的拼写,实际配置键名为materialize_dataset_inputs(见 mlflow.py),配置时请以源码为准。
平台解析优先级
_get_dataset_platform_from_source_type(mlflow.py)定义了 source type 到平台的解析优先级:
- 用户映射:
source_mapping_to_platform中显式配置的映射; - 内置映射:例如
gs→gcs; - 直接平台匹配:source type 本身是 DataHub 已知的有效平台名(通过
KNOWN_VALID_PLATFORM_NAMES校验)则直接使用。
若解析不到平台:当materialize_dataset_inputs=false时,仅创建 mlflow 平台下的数据集引用(reference),不建立 upstream;当materialize_dataset_inputs=true时,会记录 failure 报告,提示 “No mapping dataPlatform found for dataset input source type. Please addmaterialize_dataset_inputs.source_mapping_to_platformin config.”(见 mlflow.py)。
数据集输入的处理分支
_get_dataset_input_workunits(mlflow.py)对每个 Run 的 dataset input 按 source type 分三类处理:
- local / code 类型:始终创建
mlflow平台下的本地数据集实体(含 schema 与自定义属性),不涉及外部平台; - 托管数据集 + 物化开启:创建目标平台(如 snowflake)的 hosted dataset,同时创建 mlflow 平台下的数据集引用,并为其添加指向 hosted dataset 的
UpstreamLineageClass(type=COPY); - 托管数据集 + 物化关闭:只创建数据集引用,若平台可解析则尝试通过 upstream 链接到已存在的目标数据集(需要 DataHub graph 可查询)。
最终所有数据集引用 URN 会以DataProcessInstanceInputClass的inputEdges挂到 Run 对应的DataProcessInstance上,完成 run-to-dataset 血缘。
数据集 schema 解析
数据集 schema 支持 MLflowmlflow_colspec格式:当dataset.schema是合法 JSON 且包含mlflow_colspec键时,会提取(name, type)列表写入数据集 schema;JSON 解析失败时记录 warning 并跳过 schema,仅把原始 schema 字符串放入自定义属性(见 mlflow.py)。
内部工作流程:一次采集发生什么
从get_workunits_internal(mlflow.py)可以看到采集分三阶段:
_get_tags_workunits():先为 Model Registry 的四种 Stage 生成 Tag 实体;_get_experiment_workunits():遍历所有 Experiment,每个 Experiment 生成 Container;再遍历其下所有 Run,为每个 Run 生成DataProcessInstance相关 workunit,并处理数据集输入血缘;_get_ml_model_workunits():遍历所有 Registered Model,生成MlModelGroup;再遍历其下所有 Model Version,为每个版本生成MlModel(属性、版本属性、Stage Tag 关联)。
数据获取统一走_traverse_mlflow_search_func(mlflow.py),它会自动处理search_experiments、search_runs、search_registered_models、search_model_versions返回的PagedList分页游标,直到没有下一页为止。单元测试test_traverse_mlflow_search_func(test_mlflow_source.py)用一个三页的假分页函数验证了该遍历逻辑。
Run 到模型的产出血缘
在_get_run_workunits中,连接器通过search_model_versions(filter_string=f"run_id = '{run_id}'")反查该 Run 产出的 Model Version,并为其构造DataProcessInstanceOutputClass,将模型版本 URN 挂到outputEdges(mlflow.py)。这样就形成了Run → DataProcessInstance → MlModel的完整产出链路。
无 Run 的 Model Version 处理
并非所有 Model Version 都有关联的 Run(例如手工注册的模型)。_get_mlflow_run(mlflow.py)在model_version.run_id为空时返回None;此时_get_ml_model_properties_workunit会生成hyperParams=None、trainingMetrics=None、trainingJobs=[]的模型属性(mlflow.py)。单元测试test_model_without_run(test_mlflow_source.py)验证了该行为。
外部 URL 构造
- Model Version 外部 URL:优先使用
base_external_url,否则当tracking_uri以http开头时回退到 tracking_uri,格式为{base}/#/models/{name}/versions/{version}(mlflow.py)。单元测试test_make_external_link_*系列(test_mlflow_source.py)分别验证了本地 URI 不生成 URL、远程 URI 自动生成、配置优先于 tracking_uri 三种场景; - Run 外部 URL:仅当 tracking_uri 为 HTTP 时生成
{base}/#/experiments/{experiment_id}/runs/{run_id}(mlflow.py)。
用户信息与运行事件
每个 Run 还会生成一个platformResource用户实体(以run.info.user_id为主键),并作为DataProcessInstancePropertiesClass.created.actor;运行结束时(存在end_time)会产出DataProcessInstanceRunEventClass,其中:
- 状态映射(
_convert_run_result_type,见 mlflow.py):FINISHED→SUCCESS,FAILED→FAILURE,其他 →SKIPPED; - 附带
durationMillis(end_time − start_time)与nativeResultType=mlflow。
验证与测试:如何确认采集正确
仓库提供了两级测试来验证连接器行为:
- 单元测试(metadata-ingestion/tests/unit/test_mlflow_source.py):覆盖 Stage Tag 生成、
model_name_separator对 URN 的影响、无 Run 模型、分页遍历、外部 URL 构造,以及materialize_dataset_inputs开/关与平台映射组合下的数据集血缘行为; - 集成测试(metadata-ingestion/tests/integration/mlflow/test_mlflow_source.py):使用真实
MlflowClient在临时mlruns/目录中创建实验、Run(含参数与指标)、Registered Model 与 Model Version,并将 Stage 切换到 Archived,然后以 file sink 运行完整 Pipeline,与黄金文件 mlflow_mcps_golden.json 对比产出的 MCP(Metadata Change Proposal)序列。
集成测试的 recipe 结构(来自 test_mlflow_source.py)也可以作为本地验证模板:
source: type: mlflow config: tracking_uri: /path/to/mlruns sink: type: file config: filename: mlflow_source_mcps.json局限性说明
原文档指出,模块行为受源平台 API、权限与暴露的元数据限制。结合源码,具体表现为:
- 不采集 MLflow 的 feature 级信息(无
MlFeature实体、不填充mlModelTrainingData); - 当服务端版本低于 1.28.0 时,Experiments 与 Runs 采集被跳过(仅保留模型注册相关采集);
materialize_dataset_inputs=true但平台无法解析时,该数据集输入会被跳过并记录 failure;- 本地文件型 tracking_uri 不会生成外部 URL。
故障排查建议
按 mlflow_post.md 的官方建议,采集失败时按以下顺序排查:
- 验证凭据:
username/password必须成对出现,否则连接器启动即报错; - 验证权限与网络连通性:确保对 MLflow Tracking/Registry API 有读取权限且网络可达;
- 检查服务端版本:低于 1.28.0 时留意日志中
MLflow API Endpoint Not Found for Experiments的 warning; - 审查采集日志中的 source 级错误信息,结合错误上下文调整配置(例如补充
source_mapping_to_platform映射)。
通过以上配置与原理说明,即可将 MLflow 的模型资产、实验与运行历史完整接入 DataHub,形成可搜索、可追溯的机器学习元数据图谱。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考