news 2026/9/18 17:50:51

DataHub 实例间元数据迁移:datahub 源连接器(datahub_pre)原理与实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
DataHub 实例间元数据迁移:datahub 源连接器(datahub_pre)原理与实战指南

DataHub 实例间元数据迁移:datahub 源连接器(datahub_pre)原理与实战指南

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

导读

本篇技术指南聚焦于 DataHub 元数据迁移场景中的datahub源连接器(对应文档 metadata-ingestion/docs/sources/datahub/datahub_pre.md),讲解如何将一个 DataHub 实例中的元数据(版本化 Aspect 与时序 Aspect)完整迁移到另一个 DataHub 实例。读者将掌握该连接器的数据读取顺序、关键配置项、数据库索引要求、Kafka 消费策略与有状态增量机制,并能直接写出一份可运行、可断点续传的迁移 Recipe。


一、连接器概述:在 DataHub 实例之间"搬运"元数据

datahub源连接器(platform: "datahub",支持状态为 GA,见 datahub_source.py)的核心用途只有一个:将元数据从一个 DataHub 实例迁移到另一个 DataHub 实例。它并非从某个外部系统(数据库、数仓)采集元数据,而是把源 DataHub 已持久化的元数据原样读出,再以 MetadataChangeProposalWrapper(MCP)的形式作为 WorkUnit 交给目标端 sink 写入。

该连接器从两个位置拉取数据:

数据来源数据类型说明
DataHub 数据库版本化 Aspect(Versioned Aspects)存储在源实例 GMS 后端数据库的metadata_aspect_v2表中
DataHub Kafka时序 Aspect(Timeseries Aspects)通过 Kafka 主题中的 Metadata Change Log(MCL) 事件读取

读取顺序是固定的:先完整读取数据库中的版本化 Aspect,再消费 Kafka 中的时序 MCL。这一顺序在 get_workunits_internal 中体现得十分清楚:_get_database_workunits执行完毕后才会进入_get_kafka_workunits

二、防止无限运行的 stop_time 机制

文档明确指出一个关键设计:为防止该源"永远跑下去",连接器不会消费 ingestion job 启动之后才产生的数据。每次运行开始时,get_workunits_internal会立即记录一个stop_time

self.report.stop_time = datetime.now(tz=timezone.utc) logger.info(f"Ingesting DataHub metadata up until {self.report.stop_time}")

该时间点被同时用于:

  • 数据库侧:作为查询时间范围的上界(createdon < stop_time);
  • Kafka 侧:消费时若遇到 MCL 的审计时间戳(audit stamp)超过stop_time,则立即停止读取(见 datahub_kafka_reader.py)。

stop_time本身会被写入运行报告(report.py),方便在 UI 或日志中确认本次迁移的时间边界。

三、有序读取:数据库按 createdon、Kafka 按 offset

为了保证迁移的一致性,两路数据都按**时间顺序(chronological order)**读取:

  • 数据库:按createdon时间戳排序读取;
  • Kafka:按每个分区(partition)的 offset 顺序读取,断点续传时按{partition -> offset}映射恢复(见 state.py)。

数据库侧的实际 SQL 查询(datahub_database_reader.py)以ORDER BY createdon, urn, aspect, version排序,并通过"时间分页 + offset 分页"混合策略分批拉取:每当一批数据全部落在同一createdon时改用 offset 递增,否则以最新createdon作为下一批起点。值得注意的是,文档与代码都提示这种分页在跨批边界处可能返回重复行(多个行共享同一createdon时),这是有状态去重/幂等设计下被接受的预期行为。

关于 createdon 索引的硬性要求

文档特别强调:要正确、高效地读取数据库,必须保证metadata_aspect_v2表的createdon建有索引。新建的数据库默认会带一个名为timeIndex的索引,但历史数据库可能需要手工创建:

CREATE INDEX timeIndex ON metadata_aspect_v2 (createdon);

文档以醒目警告提示:若缺少该索引,连接器可能运行极慢,并对数据库造成显著的负载(每次时间范围查询都会触发全表扫描)。

四、运行前置条件(Prerequisites)

在启动迁移前,需要满足以下条件:

  1. 网络连通性:能够访问源 DataHub 实例所在的环境;
  2. 有效的认证凭据:具备读取元数据 API 的权限;
  3. 只读权限:拥有本模块所需元数据 API 的读取权限。

特别地,该连接器必须直连源 DataHub 实例的以下三部分基础设施(而非仅通过 GMS HTTP API):

  • 数据库(存储版本化 Aspect)
  • Kafka broker(存储时序 MCL)
  • Kafka Schema Registry(用于反序列化 Avro 格式的 MCL)

备注:仓库源码还提供了一条可选的pull_from_datahub_api路径(见下文),此时才要求配置datahub_api(即 GMS REST 连接)而非直连数据库。

五、Recipe 配置详解:从连接串到增量策略

datahub源的配置类为DataHubSourceConfig(config.py),继承了StatefulIngestionConfigBase。一个典型的迁移 Recipe 结构如下:

source: type: datahub config: # --- 数据来源一:数据库(版本化 Aspect) --- database_connection: scheme: mysql+pymysql # MySQL 必须使用 mysql+pymysql host: <source-db-host> port: 3306 username: <db-user> password: <db-password> db: datahub # 通常为 datahub 库 # --- 数据来源二:Kafka(时序 Aspect) --- kafka_connection: bootstrap: <source-kafka-broker:9092> schema_registry_url: http://<source-schema-registry:8081> # 可选:consumer_config 中追加 SASL/SSL 等认证参数 # --- 迁移行为控制 --- include_all_versions: false include_soft_deleted_entities: true exclude_aspects: - datahubIngestionRunSummary - datahubIngestionCheckpoint - testResults stateful_ingestion: enabled: true sink: type: datahub-rest config: server: http://<target-datahub-gms:8080> token: <target-gms-token>

5.1 核心配置参数

以下参数均来自DataHubSourceConfig(默认值以源码为准):

参数默认值作用
database_connectionNone数据库连接配置(SQLAlchemyConnectionConfig),提供版本化 Aspect 数据源
kafka_connectionNoneKafka 连接配置(KafkaConsumerConnectionConfig),提供时序 MCL 数据源
include_all_versionsFalse是否包含每个 Aspect 的全部历史版本;关闭时只取最新版本(version = 0
include_soft_deleted_entitiesTrue是否包含被软删除(soft deleted)的实体
exclude_aspects{"datahubIngestionRunSummary", "datahubIngestionCheckpoint", "testResults"}需要排除的 Aspect 名称集合
database_query_batch_size10000数据库单次查询拉取的行数
database_table_namemetadata_aspect_v2存储版本化 Aspect 的数据库表名
kafka_topic_nameMetadataChangeLog_Timeseries_v1存储时序 MCL 的 Kafka 主题名
stateful_ingestionenabled=True有状态增量摄入(该源默认开启,与多数源不同)
commit_state_interval1000每处理多少条记录提交一次检查点
commit_with_parse_errorsFalse出现解析错误时是否仍推进 createdon/offset 检查点
pull_from_datahub_apiFalse(隐藏参数)改为通过 DataHub API 拉取版本化 Aspect
max_workers5 * cpu_count(隐藏参数)DataHub API 拉取时的线程数
urn_patterndeny 环境专属 URN(见下)URN 过滤模式
drop_duplicate_schema_fieldsFalse是否丢弃schemaMetadata中的重复字段路径(源库存在重复、目标端有服务端去重时适用)
query_timeoutNone数据库单次查询超时(秒)
preserve_system_metadataTrue是否复制源系统的 systemMetadata

5.2 三个值得注意的设计点

(1)至少配置一个数据来源,否则拒绝运行。配置类通过模型校验器强制要求:database_connectionkafka_connectionpull_from_datahub_api三者至少提供一个,否则抛出ValueError("Your current config will not ingest any data...")(config.py)。文档推荐两者都配("ideally both"),以保证版本化与时序 Aspect 都能迁移。

(2)MySQL 连接串有硬性约束。当 scheme 包含mysql时,必须是mysql+pymysql,否则校验失败(config.py)。

(3)默认排除环境专属 URN。urn_pattern默认 deny 以下四类 URN(config.py):

urn:li:dataHubIngestionSource:.* urn:li:dataHubSecret:.* urn:li:globalSettings:.* urn:li:dataHubExecutionRequest:.*

这是为了防止把加密凭据(Secret)与易产生脏实体的环境配置复制到目标实例。源码甚至会在用户显式自定义urn_pattern时给出 warning,提醒保留这些默认 deny 规则(datahub_source.py)。

5.3 关于 exclude_aspects 的警告

exclude_aspects仅适用于"想摄入实体但剔除某些 Aspect"的场景;若要整体排除某类实体,应使用urn_pattern.deny。文档与配置注释同时警告:排除 key aspect 而保留其他 Aspect 可能产生无效实体(config.py)。

六、源码级原理:三路读取器的工作方式

DataHubSource在运行时按需实例化三个读取器,对应三种数据获取路径:

6.1 DataHubDatabaseReader:版本化 Aspect 的主路径

datahub_database_reader.py 负责直连源实例数据库:

  • 通过 SQLAlchemy 创建引擎,使用get_sql_alchemy_url()拼接连接串;
  • 查询逻辑上对metadata_aspect_v2自关联,把statusAspect(version = 0)中的removed字段提取出来,用于软删除过滤。该 JSON 提取表达式按方言区分:PostgreSQL 用((metadata::json)->>'removed')::boolean,其他(如 MySQL)用JSON_EXTRACT(metadata, '$.removed')(datahub_database_reader.py);
  • 支持 PostgreSQL/MySQL/MariaDB 的服务端游标流式读取stream_results=True+yield_per=batch_size),并可按query_timeout设置statement_timeout(PG)或max_execution_time(MySQL);
  • include_all_versions=True时,启用VersionOrderer:同一createdon时间戳下版本 0 的 Aspect 被延后到最后输出,保证"最新版本后写、旧版本先写"的稳定顺序(datahub_database_reader.py);
  • 解析时通过ASPECT_MAP将 JSON 元数据还原为对应的 Aspect 类对象,preserve_system_metadata=True时会复制 systemMetadata(并剥离isNoOp标记);解析失败计入num_database_parse_errorsdatabase_parse_errors明细报告。

另外,get_all_aspects会分两轮查询:先拉取urn:li:structuredProperty:*(结构化属性),等待structured_properties_template_cache_invalidation_interval(默认 1 秒)让目标端模板缓存失效后,再拉取其余 Aspect,避免结构化属性模板未注册导致后续 Aspect 引用失败(datahub_database_reader.py)。

6.2 DataHubKafkaReader:时序 Aspect 的消费路径

datahub_kafka_reader.py 使用 Confluent Kafka Python 客户端:

  • DeserializingConsumer+AvroDeserializer消费主题,反序列化依赖Schema Registry(这就是前置条件要求直连 schema registry 的原因);
  • 关键消费参数:auto.offset.reset=earliest(从头开始)、enable.auto.commit=False(手动管理进度,配合有状态检查点);
  • consumer group 固定为datahub_source-{pipeline_name}(datahub_kafka_reader.py),on_assign回调会把各分区 offset 恢复到上次检查点记录的{partition -> offset},未记录的分区从OFFSET_BEGINNING开始;
  • 逐条轮询(poll(10))解析为MetadataChangeLogClass,解析失败计入num_kafka_parse_errors;遇到created.time > stop_time的 MCL 立即停止(时间边界机制);命中exclude_aspects的 MCL 跳过并计入num_kafka_excluded_aspects

6.3 DataHubApiReader:可选的 API 拉取路径

pull_from_datahub_api=True时,连接器改用 DataHub Graph API 拉取版本化 Aspect(datahub_api_reader.py):

  • 通过ctx.graph(即 Recipe 中的datahub_api配置)调用get_urns_by_filter枚举实体 URN,并依据include_soft_deleted_entities选择过滤已软删除实体;
  • 使用ThreadPoolExecutor(线程数 =max_workers)并发为每个 URN 调用get_entity_semityped获取其全部 Aspect;
  • 若未配置datahub_api,该路径会直接记录 failure。

七、有状态增量:断点续传与幂等设计

与大多数源不同,datahub源将stateful_ingestion默认开启。状态由StatefulDataHubIngestionHandler(state.py)管理,检查点结构为:

class DataHubIngestionState(CheckpointStateBase): database_createdon_ts: NonNegativeInt = 0 # 数据库侧进度(毫秒时间戳) kafka_offsets: Dict[int, NonNegativeInt] # Kafka 各分区已消费的 offset

运行逻辑(datahub_source.py):

  1. 从上次检查点恢复from_createdon(数据库)与from_offsets(Kafka);
  2. 数据库每处理一条记录就更新database_createdon_ts,每commit_state_interval(默认 1000)条提交一次;
  3. Kafka 每消费一条 MCL 就按offset + 1更新对应分区进度;
  4. 只有当没有解析错误(或显式设置commit_with_parse_errors=True)时才推进检查点,避免"带病提交"导致数据永久丢失。

这保证了迁移任务中途失败后可安全重跑:数据库从上次createdon续读,Kafka 从上次 offset 续消费,天然幂等。

八、运行报告:如何验证迁移结果

连接器报告(report.py)提供以下关键指标,可用于核对迁移完整性:

指标含义
stop_time本次运行的读取截止时间(与文档描述的防止无限运行机制对应)
num_database_aspects_ingested从数据库摄入的版本化 Aspect 数量
num_database_parse_errors数据库侧解析失败条数(按 error → aspect → urn 记录明细)
num_kafka_aspects_ingested从 Kafka 摄入的时序 Aspect 数量
num_kafka_parse_errorsKafka 侧反序列化失败条数
num_kafka_excluded_aspectsexclude_aspects跳过的 MCL 数
num_timeseries_deletions_dropped被丢弃的时序 DELETE 变更数
num_timeseries_soft_deleted_aspects_dropped因实体软删除而被丢弃的时序 Aspect 数

其中两个"丢弃"指标对应 datahub_source.py 的过滤逻辑:ChangeTypeClass.DELETE的时序变更被跳过;当include_soft_deleted_entities=False时,先从数据库查出软删除 URN 列表,再从 Kafka 流中过滤掉这些实体的时序 Aspect。

九、迁移实操建议与注意事项

  1. 先建索引再迁移:确认源库存在timeIndexcreatedon索引),否则先执行CREATE INDEX timeIndex ON metadata_aspect_v2 (createdon);
  2. 保持默认 urn_pattern:保留对 Ingestion Source / Secret / Settings 等环境专属 URN 的排除,避免复制加密凭据;
  3. 优先"双连接"配置:同时配置database_connectionkafka_connection,才能完整迁移版本化与时序两类 Aspect;若只配一个,另一类会被跳过(日志中会出现 "Skipping ingestion of versioned aspects..." 等提示);
  4. 注意 MySQL 方言约束:MySQL 场景必须写scheme: mysql+pymysql
  5. 保留默认 exclude_aspectsdatahubIngestionRunSummarydatahubIngestionCheckpointtestResults属于运行期状态,迁移它们没有业务价值,默认排除是合理选择;
  6. 利用有状态增量:由于默认开启 stateful ingestion,迁移失败后直接重跑同一 Recipe 即可续传,无需清空目标端;
  7. 核对运行报告:迁移完成后对照上文表格中的各项计数,确认无大量解析错误、且stop_time覆盖了预期的数据时间范围。

十、小结

datahub源连接器(datahub_pre)是 DataHub 官方提供、状态为 GA 的实例间迁移工具,通过"数据库(版本化 Aspect)+ Kafka MCL(时序 Aspect)"双通道、按createdon/offset 有序读取,配合默认开启的有状态检查点与stop_time边界机制,实现了可增量、可续传、可核对的元数据迁移。部署前重点检查:源库createdon索引、三端网络连通性(数据库 / Kafka / Schema Registry)、以及 Recipe 中数据来源与过滤规则的配置。


延伸阅读

  • 时序 Aspect 依赖的 MCL 事件模型:docs/what/mxe.md
  • 源连接器入口与主流程:datahub_source.py
  • 完整配置参数与默认值:config.py
  • 数据库读取实现:datahub_database_reader.py
  • Kafka 消费实现:datahub_kafka_reader.py
  • 有状态检查点实现:state.py
  • 运行报告指标:report.py

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/18 17:49:59

MySQL 修改密码全攻略:8.0/5.7、忘记密码与认证插件

MySQL 密码这个东西&#xff0c;平时放在那里谁也不会多看一眼&#xff0c;直到某天你打开 Navicat 或者敲下mysql -u root -p&#xff0c;回车三次都提示ERROR 1045 (28000): Access denied for user rootlocalhost&#xff0c;那一刻脑子里是空的。尤其是刚装完 MySQL 的朋友…

作者头像 李华
网站建设 2026/9/18 17:48:42

Windows计划任务排查:应急响应中快速定位可疑任务与持久化后门

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/18 17:47:08

LeetCode 212:Trie + DFS 高效解决多单词搜索与剪枝优化

1. LeetCode 212 到底在考什么这道题其实把 LeetCode 里两个高频考点缝合在了一起&#xff1a;DFS&#xff08;深度优先搜索&#xff09;和经典数据结构Trie&#xff08;前缀树&#xff09;。很多人在做到第 79 题"单词搜索"的时候觉得挺顺&#xff0c;一个二维棋盘 …

作者头像 李华
网站建设 2026/9/18 17:44:51

Ubuntu 20.04软件安装与软件中心打不开修复指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/18 17:42:00

喷雾燃烧机理与数值模拟实战解析

1. 喷雾燃烧基础与工程应用喷雾燃烧技术是现代动力装置的核心技术之一&#xff0c;从航空发动机到工业锅炉都离不开这项关键技术。作为一名长期从事燃烧仿真研究的工程师&#xff0c;我见证了许多项目因为对喷雾燃烧机理理解不足而导致性能不达标的情况。本文将系统梳理喷雾燃烧…

作者头像 李华
网站建设 2026/9/18 17:40:31

工具调用偶发循环?Sonnet 5 用 TaoToken 先核对 Base URL

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华