DataHub MLflow 数据源接入指南:架构映射、认证配置与数据集血缘实践
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
MLflow 是机器学习生命周期管理平台,本文围绕 DataHub 元数据摄取管道中内置的mlflow连接器(mlflow.py),系统讲解其摄取能力、MLflow 与 DataHub 实体之间的概念映射、认证与数据集血缘配置,并结合源码与单元测试说明底层实现机制。读完本文,你将掌握如何编写可运行的 MLflow 摄取 Recipe、如何用source_mapping_to_platform与materialize_dataset_inputs控制数据集血缘行为,以及如何排障常见摄取问题。
连接器能力总览
MLflow 连接器用于将 MLflow 中的元数据摄取到 DataHub,面向生产环境摄取工作流。模块整体能力声明位于源码的装饰器定义中(mlflow.py):
| 能力 | 说明 | 是否需要额外配置 |
|---|---|---|
| 描述(DESCRIPTIONS) | 提取 MLflow Registered Model 与 Model Version 的描述信息 | 否 |
| 容器(CONTAINERS) | 将 MLflow Experiments 提取为容器(subtype 为 MLFLOW_EXPERIMENT) | 否 |
| 标签(TAGS) | 为 MLflow Model Registry 的 Stage 生成对应标签 | 否 |
| 数据集血缘(Dataset Lineage) | 将 Run 的 dataset inputs 映射为 DataHub 数据集并建立血缘 | 是(见下文"数据集血缘"章节) |
| 有状态删除(Stateful Deletion) | 通过stateful_ingestion配置启用陈旧实体清理 | 否(可选) |
从源码看,MLflowSource继承自StatefulIngestionSourceBase,因此天然支持有状态摄取;其摄取入口get_workunits_internal(mlflow.py)依次产出三类工作单元:Stage 标签、Experiments(含 Runs 与数据集输入)、ML Models(Registered Model 与 Model Version)。
版本兼容性
连接器要求 MLflow 服务器版本1.28.0 或更高。如果使用更早的版本,Experiments 和 Runs 的摄取将被跳过。
该约束同样体现在源码中:_traverse_mlflow_search_func遍历 MLflow 的search_experiments/search_runs等分页接口时,若捕获到ENDPOINT_NOT_FOUND错误,会在报告中记录警告并跳过实验与运行的摄取(mlflow.py)。这也是文档建议升级到 1.28.0 的底层原因——较旧版本缺少按分页 token 遍历的搜索端点。
摄取内容边界
MLflow 集成覆盖:Registered Models、Model Versions、Experiments、Runs,以及 Run 到 Model 的血缘,同时捕获标签和有状态删除检测。
需要特别说明的两个边界:
- MLflow 特性不会以 DataHub
MlFeature实体摄入:MLflow 的 tracking API 不记录 feature 到列的出处信息,连接器没有可读取的数据来源。 - 模型不会在模型级别链接其训练数据集(
mlModelTrainingData不会被填充),但 Run 级别的数据集血缘会被捕获(详见下方概念映射的 Dataset Input 行)。
概念映射:MLflow 实体如何落到 DataHub
MLflow 中的实体与 DataHub 元数据模型并不是一一对应,连接器遵循以下映射规则(来源:mlflow README):
| MLflow 源概念 | DataHub 目标概念 | 说明 |
|---|---|---|
| Registered Model | MlModelGroup | Model Group 的名称与 Registered Model 名称相同(如my_mlflow_model)。Registered Model 在 MLflow 中充当同一模型多个版本的容器。 |
| 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 代表模型的一次具体迭代,带自己的产物与元数据。 |
| Experiment | Container | MLflow 中每个 Experiment 映射为 DataHub 的 Container。Experiment 组织相关 Runs,是模型开发迭代的逻辑分组,可追踪参数、指标与产物。 |
| Run | DataProcessInstance | 捕获 Run 的执行细节、参数、指标以及与模型的血缘。 |
| Model Stage | Tag | Stage 与标签的映射关系:Production →mlflow_production、Staging →mlflow_staging、Archived →mlflow_archived、None →mlflow_none。Model Stage 标识每个版本的部署状态。 |
| Dataset Input | DataProcessInstanceInput | 通过mlflow.log_input()记录,将 Run 链接到其训练数据集。需要开启materialize_dataset_inputs才会同时创建被引用的数据集实体。 |
这些映射在源码中均有对应实现:
- Stage 标签在
_get_tags_workunits中创建,标签名称由_make_stage_tag_name统一生成为mlflow_{stage_name.lower()}(mlflow.py),且四个 Stage 均带有预置的颜色与描述(Production 为绿色#308613、Staging 为黄色#FACB66、Archived 为灰色#5D7283、None 为浅灰#F2F4F5)。 - Experiment 通过
Container与ExperimentKey(平台 + experiment 名称)生成容器实体,并把mlflow.note.content解析为容器描述、artifacts_location写入自定义属性(mlflow.py)。 - Model Version 的命名由
model_name_separator配置决定(默认_),在_make_ml_model_urn中拼接(mlflow.py);同时每个 Model Version 还会通过VersionPropertiesClass关联到VersionSet,并写入别名(aliases)与排序 ID(mlflow.py)。 - Run 被建模为
DataProcessInstance,其参数、指标分别转为MLHyperParamClass与MLMetricClass,运行状态(FINISHED/FAILED 等)映射为 SUCCESS/FAILURE/SKIPPED(mlflow.py)。
前置条件
开始摄取前,请确认:
- 到 MLflow 源(tracking server / registry server)的网络连通性;
- 有效的认证凭据;
- 本模块所需元数据 API 的读取权限。
安装与最小 Recipe
mlflow连接器随 DataHub 元数据摄取框架提供,通过metadata-ingestion包使用。最小 Recipe 如下(来源:mlflow_recipe.yml):
source: type: mlflow config: # Coordinates tracking_uri: tracking_uri sink: # sink configstracking_uri用于指定 MLflow Tracking Server 地址(如http://127.0.0.1:5000)。根据源码MLflowConfig的定义(mlflow.py),若未设置,连接器将回退到 MLflow 默认行为(本地mlruns/目录或MLFLOW_TRACKING_URI环境变量)。
完整配置参数说明
以下是MLflowConfig支持的全部配置项(来自源码字段定义与注释):
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
tracking_uri | str | None(回退 MLflow 默认值) | Tracking Server URI;未设置时使用本地mlruns/目录或MLFLOW_TRACKING_URI环境变量 |
registry_uri | str | None(回退默认值) | Registry Server URI;未设置时使用tracking_uri或MLFLOW_REGISTRY_URI环境变量 |
model_name_separator | str | _ | 模型名称与其版本号之间的分隔符(如model_1或model-1) |
base_external_url | str | None | 构造指向 MLflow UI 的外部 URL 时使用的基础 URL;未设置时若tracking_uri为 HTTP URL 则使用之;两者均未设置则不生成外部 URL |
materialize_dataset_inputs | bool | False | 是否为每个 Run 物化(创建)数据集输入实体 |
source_mapping_to_platform | dict | None | 将 MLflow 数据集 source type 映射到 DataHub 平台的映射表 |
username | str | None | MLflow 认证用户名 |
password | TransparentSecretStr | None | MLflow 认证密码(摄入日志中会被脱敏处理) |
stateful_ingestion | StatefulStaleMetadataRemovalConfig | None | 有状态摄取配置,用于陈旧实体删除 |
另外,MLflowConfig继承自EnvConfigMixin,这意味着标准的环境相关配置(如env,默认PROD)同样可用,并会体现在生成的实体 URN 中。
认证配置
连接器支持通过username与password两个配置项向 MLflow 服务器进行认证:
source: type: mlflow config: tracking_uri: "http://127.0.0.1:5000" username: <username> password: <password>底层实现中,_configure_client会先校验username与password必须成对出现——只设置其中一个会直接抛出ValueError(mlflow.py)。校验通过后,连接器将凭据写入MLFLOW_TRACKING_USERNAME与MLFLOW_TRACKING_PASSWORD环境变量,再构造MlflowClient(tracking_uri=..., registry_uri=...)。由于password字段使用TransparentSecretStr类型,摄入日志中不会明文打印密码。
数据集血缘配置
MLflow Run 可以通过mlflow.log_input()记录训练数据集(Dataset Input)。连接器支持将不同 MLflow 引擎产生的数据集映射到指定的 DataHub 平台,这是通过source_mapping_to_platform配置项实现的。
平台映射规则
source_mapping_to_platform: huggingface: snowflake # Maps Hugging Face datasets to Snowflake platform http: s3 # Maps HTTP data sources to s3 platform平台解析的完整优先级(来自_get_dataset_platform_from_source_type,mlflow.py):
- 用户自定义映射:
source_mapping_to_platform中配置的 source type → 平台映射优先使用; - 内置映射:
gs会被转换为gcs; - 直接平台匹配:若 source type 本身是 DataHub 已知平台名称(如
snowflake、s3、bigquery),则直接使用。
其中"已知平台名称"来自 data_platforms.py 中的KNOWN_VALID_PLATFORM_NAMES列表,包含bigquery、cassandra、databricks、delta-lake、dbt、feast、file、gcs、hdfs、hive、mssql、mysql、oracle、postgres、redshift、s3、sagemaker、snowflake、streamlit。注意该列表注释说明它并不完整,仅用于血缘生成时的自动平台映射与手动覆盖,不应作为 URN 校验依据。
数据集物化(materialize)行为
默认行为:仅按平台和名称链接到已存在的数据集,不会创建新数据集。
若希望自动创建数据集实体,需启用materialize_dataset_inputs:
materlize_dataset_inputs: true # Creates new datasets if they don't exist两个配置项可以独立组合使用:
# Only map to existing datasets materlize_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 materlize_dataset_inputs: true source_mapping_to_platform: huggingface: snowflake pytorch: snowflake注意:
materlize_dataset_inputs是官方文档与 Recipe 中使用的键名(拼写如此,非materialize),配置时请保持与文档一致。源码中的 Python 字段名为materialize_dataset_inputs,两者对应同一配置。
底层血缘实现
_get_dataset_input_workunits(mlflow.py)按以下分支处理每个 dataset input:
- 本地/代码数据集(source type 为
local或code):直接在mlflow平台下创建数据集实体,不涉及外部平台映射; - 托管数据集(其他 source type):
- 若
materialize_dataset_inputs开启:先在映射平台下创建托管数据集实体;若 source type 找不到任何映射平台,会在报告中记录失败(failure),提示"请配置materialize_dataset_inputs.source_mapping_to_platform"并跳过该数据集; - 随后在
mlflow平台下创建一个数据集引用实体,并通过UpstreamLineageClass(type 为COPY)指向外部平台的数据集——即使未开启物化,只要平台可解析,引用与上游血缘依然会建立(与"默认仅链接已有数据集"的行为一致); - 最后将所有数据集引用作为
DataProcessInstanceInput的inputEdges附加到 Run 上,形成 Run → 训练数据集的输入血缘。
- 若
此外,连接器会尝试解析数据集 schema:优先读取mlflow_colspec格式的字段名与类型;若 schema 不是该格式则原样放入自定义属性schema,JSON 解析失败时会在报告中记录警告(mlflow.py)。
测试用例印证
单元测试 test_mlflow_source.py 覆盖了上述全部分支:
test_materialization_disabled_with_supported_platform:关闭物化且平台可解析时,仅创建引用并建立 COPY 上游;test_materialization_disabled_with_unsupported_platform:关闭物化且平台不可解析时,引用实体不带上游;test_materialization_enabled_with_supported_platform:开启物化时创建托管数据集实体;test_materialization_enabled_with_unsupported_platform:开启物化但无映射平台时记录失败并跳过;test_materialization_enabled_with_custom_mapping:通过source_mapping_to_platform将不支持的 source type(如unsupported_platform)映射为snowflake后成功物化。
外部链接与运行细节
连接器会为以下实体生成指向 MLflow UI 的外部 URL:
- Model Version:
{base_url}/#/models/{model_name}/versions/{version}。基础 URL 优先使用base_external_url,其次在tracking_uri为 HTTP URL 时使用tracking_uri;两者皆无则不生成(mlflow.py)。单元测试test_make_external_link_local/test_make_external_link_remote/test_make_external_link_remote_via_config验证了这三种情形。 - Run:
{tracking_uri}/#/experiments/{experiment_id}/runs/{run_id},仅在tracking_uri以 http 开头时生成(mlflow.py)。
此外,Run 的元数据还包含:run.info.run_name(缺失时回退 run_id)作为展示名称、user_id(缺失时回退为mlflow)作为执行者、artifact_uri作为输出 URL,以及start_time/end_time计算出的运行时长(durationMillis)(mlflow.py)。
限制
模块行为受平台暴露的源 API、权限与元数据约束,请参考能力说明中标注为不支持或有条件的特性。当前已知限制包括:
- MLflow tracking API 不暴露 feature 到列的出处,因此无法摄取
MlFeature实体; - 模型级别不填充
mlModelTrainingData,训练数据集血缘仅在 Run 级别体现; - 若 Model Version 没有关联的 Run(
model_version.run_id为空),则该模型的超参数与训练指标不可用(源码注释明确说明,mlflow.py); - Experiments 与 Runs 的摄取依赖 MLflow 1.28.0+ 的搜索 API,旧版本会被跳过。
故障排查
若摄取失败,请按以下顺序排查:
- 校验凭据:确认
username/password有效且成对配置(只配其中一个会触发ValueError); - 校验权限:确认服务账号对 tracking/registry API 具备读取权限;
- 校验连通性:确认
tracking_uri/registry_uri可从执行摄取的主机访问; - 校验范围过滤:确认没有误配置导致实体被过滤;
- 查看摄入日志:连接器会在
SourceReport中记录 source-specific 错误与警告(如 API 端点缺失、schema 解析失败、物化平台缺失等),根据日志中的具体错误调整配置后重试。
例如,若日志出现 "MLflow API Endpoint Not Found for Experiments" 警告,说明 MLflow 版本低于 1.28.0,需升级服务器或接受 Experiments/Runs 被跳过的行为;若出现 "Unable to materialize dataset inputs" 失败,说明开启物化后未给对应 source type 配置source_mapping_to_platform映射。
参考资料
- 连接器官方文档:mlflow_pre.md、mlflow_post.md、mlflow_recipe.yml、mlflow README
- 核心实现:mlflow.py
- 平台列表:data_platforms.py
- 单元测试:test_mlflow_source.py
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考