使用 DataHub Source 在 DataHub 实例间迁移元数据:从数据库与 Kafka 的完整数据迁移指南
使用 DataHub Source 在 DataHub 实例间迁移元数据从数据库与 Kafka 的完整数据迁移指南【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读DataHub Source 是 datahub 项目 metadata-ingestion 模块内置的一个元数据迁移型连接器它的核心场景是将一个 DataHub 实例中的元数据完整迁移到另一个 DataHub 实例。它同时从源实例的数据库versioned aspects、Kafkatimeseries aspects 的 MCL 日志以及可选的 DataHub API 三个位置拉取数据并借助有状态摄入stateful ingestion的检查点机制实现断点续跑。读完本文你将掌握该连接器的概念映射、完整配置方法、迁移原理、检查点行为以及性能与排除项的最佳实践。一、连接器概览与概念映射DataHub Source 是一个面向 DataHub 元数据的集成模块。它的官方定位是DataHub 实用工具utility或元数据专注型集成在 metadata-ingestion 仓库的 docs/sources/datahub 目录下提供完整文档、预置 recipe 示例与集成说明。该连接器覆盖与迁移相关的元数据实体和运维对象并捕获有状态删除检测stateful deletion detection——即软删除实体在目标实例中会被过滤或标记而不是被当作普通数据复制。由于该连接器负责的是 DataHub 到 DataHub 的元数据搬运概念映射本质上是一张源概念 → DataHub 概念的对照表说明迁移后的元数据在目标 DataHub 中如何落位源概念Source ConceptDataHub 概念说明NotesPlatform/account/project scopePlatform Instance、Container在平台上下文中组织资产Core technical asset如 table/view/topic/fileDataset主要被摄入的技术资产Schema fields / columnsSchemaField支持 schema 提取时包含Ownership and collaboration principalsCorpUser、CorpGroup由支持所有权和身份元数据的模块发出Dependencies and processing relationshipsLineage edges支持并启用血缘提取时可用官方文档同时说明具体概念的映射仍处于待完善pending状态上表展示的是 DataHub 中的通用概念映射方式因此实际迁移时建议以目标实例中最终生成的实体类型为准进行验证。二、工作原理数据从哪来、按什么顺序迁移DataHub Source 的数据拉取有两个主要来源一个可选来源在 datahub_pre.md 中描述如下DataHub 数据库读取 versioned aspects版本化方面即每个实体可多次更新、带版本号的历史状态DataHub Kafka读取 MCL LogMetadata Change Log 中的 timeseries aspects时序方面DataHub API可选通过 Graph 客户端按 URN 拉取版本化方面。执行顺序遵循明确的优先级先从数据库读完全部数据再从 Kafka 摄入时序数据。为了防止该源无限期运行它不会摄入 datahub_source 摄入作业启动之后产生的新数据——这个截止时间点stop_time会记录在报告中对应源码 report.py 中的stop_time字段。数据库和 Kafka 的数据都按时间顺序读取数据库侧按createdon时间戳排序Kafka 侧按每个分区的 offset 排序。从源码 datahub_source.py 的get_workunits_internal可以看到完整的执行流程读取上次运行的检查点状态stateful_ingestion_handler.get_last_run_state()若配置了pull_from_datahub_api先从 API 拉取版本化方面若配置了database_connection创建DataHubDatabaseReader并按createdon增量拉取版本化方面随后提交数据库进度若配置了kafka_connection先查询软删除 URN 列表除非include_soft_deleted_entitiestrue再从 Kafka 拉取 timeseries aspects并提交 offset 进度。每条产出的元数据都会经过urn_pattern过滤不符合 allow/deny 规则的 URN 会被跳过。2.1 数据库读取的索引要求为保证从数据库按createdon增量读取的效率请务必确认createdon列已建立索引。新建的数据库默认会带有名为timeIndex的索引但旧库可能需要手动创建CREATE INDEX timeIndex ON metadata_aspect_v2 (createdon);注意如果没有该索引源可能会运行得极慢并对数据库产生显著负载。从源码 datahub_database_reader.py 可以看到数据库读取基于metadata_aspect_v2表可通过database_table_name配置默认表名与批大小常量定义在 config.pyDEFAULT_DATABASE_TABLE_NAME metadata_aspect_v2 DEFAULT_KAFKA_TOPIC_NAME MetadataChangeLog_Timeseries_v1 DEFAULT_DATABASE_BATCH_SIZE 10_000读取器使用 SQLAlchemy 连接仅支持 PostgreSQL / MySQL / MariaDB 的服务端游标并实现了混合分页策略按createdon时间戳分页与 offset 分页之间自动切换内存占用更优。它还内建了VersionOrderer当开启include_all_versions时同一createdon时间戳下 version 0 的行会被排到最后从而保证目标实例先写入历史版本、最后写入最新版本避免出现最新版本被旧版本覆盖的问题。2.2 Kafka 读取的实现细节Kafka 读取由 datahub_kafka_reader.py 实现使用 Confluent Kafka 的DeserializingConsumer配合 Schema Registry 的 Avro 反序列化消费MetadataChangeLog_Timeseries_v1主题。要点如下group.id采用datahub_source-{pipeline_name}前缀auto.offset.resetearliest、enable.auto.commitfalsecheckpoint 由有状态摄入自行管理通过on_assign回调把上次记录的每个分区 offset 恢复到消费者实现精确续跑遇到created时间超过stop_time的 MCL 即停止消费默认跳过exclude_aspects中的方面名并丢弃 DELETE 类型的变更与软删除实体的时序方面统计在报告中。三、前置条件在运行摄入之前需要满足以下条件见 datahub_pre.md确认到源实例的网络连通性拥有有效的认证凭据具备本模块所需的元数据 API 读取权限直接访问源 DataHub 实例的数据库、Kafka broker 与 Kafka Schema Registry。也就是说这个连接器需要的是底层基础设施直连权限而不是仅靠 GMS API 就能完成的普通源。四、完整 Recipe 配置详解预置的 recipe 示例位于 datahub_recipe.yml以下是带注释的完整版本可直接作为迁移配置的起点pipeline_name: datahub_source_1 datahub_api: server: http://localhost:8080 # 从 localhost:8080 的 DataHub 实例迁移数据 token: token source: type: datahub config: include_all_versions: false database_connection: scheme: mysqlpymysql # 或 Postgres 使用 postgresqlpsycopg2 host_port: database_host:database_port username: username password: password database: database kafka_connection: bootstrap: boostrap_url:9092 schema_registry_url: schema_registry_url:8081 stateful_ingestion: enabled: true ignore_old_state: false urn_pattern: deny: # 忽略所有匹配该正则的 datahub 元数据 urn - ^denied.urn.* allow: # 只摄入匹配该正则的 datahub 元数据 urn - ^allowed.urn.* flags: set_system_metadata: false # 是否复制系统元数据 # 这里写入一个 DataHub 实例 # 你也可以使用其他 sink例如把数据写入文件 sink: type: datahub-rest config: server: destination_gms_url token: token提示示例中flags.set_system_metadata由配置类中的preserve_system_metadata字段默认 true接管recipe 里的 flags 写法是早期约定的兼容形式。若希望完全复制源端系统元数据可保持默认若目标端不希望保留源端系统元数据可将其关闭。4.1 配置项速查表结合 config.py 中DataHubSourceConfig的字段定义各配置项含义如下配置项默认值说明database_connectionNoneSQLAlchemy 连接配置提供后才会读取 versioned aspectskafka_connectionNoneKafka 消费连接配置提供后才会读取 timeseries aspectspull_from_datahub_apifalse是否通过 DataHub API 拉取版本化方面隐藏文档项include_all_versionsfalse是否包含每个 aspect 的所有历史版本关闭时只取最新版本version 0include_soft_deleted_entitiestrue是否包含已被软删除的实体exclude_aspects{datahubIngestionRunSummary, datahubIngestionCheckpoint, testResults}要排除的方面名若要整体排除实体类型请改用urn_pattern.denydatabase_query_batch_size10000每次从数据库取回的记录数database_table_namemetadata_aspect_v2存放所有 versioned aspects 的数据库表名kafka_topic_nameMetadataChangeLog_Timeseries_v1存放 timeseries MCL 的 Kafka 主题名stateful_ingestionenabledTrue有状态摄入配置本源默认开启commit_state_interval1000每处理多少条记录提交一次检查点commit_with_parse_errorsfalse出现解析错误时是否仍更新 createdon 时间戳与 Kafka offseturn_patterndeny 内置敏感实体模式URN 过滤规则allow/deny 正则drop_duplicate_schema_fieldsfalse是否丢弃schemaMetadata中重复的 schema 字段路径query_timeoutNone每次查询的超时秒数preserve_system_metadatatrue是否从源端复制系统元数据配置校验规则源码check_ingesting_data如果database_connection、kafka_connection、pull_from_datahub_api三者均未配置配置校验会直接报错Your current config will not ingest any data提示至少配置其中一个理想情况下两者都配。数据库方言限制database_connection.scheme若为 MySQL必须是mysqlpymysql否则校验报错validate_mysql_scheme。PostgreSQL 则使用postgresqlpsycopg2。默认 URN 拒绝模式DEFAULT_URN_DENY_PATTERNS会自动排除以下环境相关实体防止把加密凭据复制到目标端、或创建损坏实体DEFAULT_URN_DENY_PATTERNS [ urn:li:dataHubIngestionSource:.*, urn:li:dataHubSecret:.*, urn:li:globalSettings:.*, urn:li:dataHubExecutionRequest:.*, ]如果你自定义了urn_pattern源码中会记录_urn_pattern_was_set标记并给出警告务必在自定义规则中保留上述默认排除项。五、有状态摄入与检查点机制该连接器默认开启有状态摄入stateful_ingestion.enabledtrue检查点按数据库createdon时间戳 Kafka 分区 offset两个维度记录状态模型见 state.py 中的DataHubIngestionState包含database_createdon_ts与kafka_offsets: Dict[partition, offset]。行为要点见 datahub_post.md首次运行从数据库最早数据与 Kafka 最早 offset 开始读取周期性提交每处理commit_state_interval默认 1000条记录就保存一次检查点上次createdon时间戳与各分区 offset中断后重启不会丢失太多进度重启会重复少量数据由于检查点是每 N 条提交一次新运行会从最近一次检查点重新摄入少量已处理数据属于预期行为错误处理摄入过程中若遇到错误例如网络错误导致某个 aspect 无法发出源会继续运行但停止提交检查点除非commit_with_parse_errorstrue。这样重新运行时可以补摄入之前遗漏的数据但代价是该错误之后的所有数据都会被重新摄入强制全量重放设置不同的pipeline_name或将stateful_ingestion.ignore_old_state: true即可从头重放全部数据。检查点提交逻辑见 datahub_source.py 的_commit_progress只有在没有解析错误或配置了commit_with_parse_errors时才会更新createdon时间戳与 Kafka offset 并提交检查点。六、排除项配置迁移时建议过滤的 URN迁移实例时通常需要排除一些包含实例专属元数据的 URN 类型如设置、角色、策略、摄入源、摄入运行记录等否则会把源实例的敏感或无用配置一并复制。官方文档建议从以下配置起步source: config: urn_pattern: # 摄入时忽略/包含的 URN 模式 deny: # 忽略所有匹配该正则的 datahub 元数据 urn - ^urn:li:role.* # 仅在不想摄入角色时排除 - ^urn:li:dataHubRole.* # 仅在不想摄入角色时排除 - ^urn:li:dataHubPolicy.* # 仅在不想摄入策略时排除 - ^urn:li:dataHubIngestionSource.* # 仅在不想摄入摄入源时排除 - ^urn:li:dataHubSecret.* - ^urn:li:dataHubExecutionRequest.* - ^urn:li:dataHubAccessToken.* - ^urn:li:dataHubUpgrade.* - ^urn:li:inviteToken.* - ^urn:li:globalSettings.* - ^urn:li:dataHubStepState.*此外exclude_aspects默认已排除datahubIngestionRunSummary、datahubIngestionCheckpoint、testResults三类方面。七、限制与注意事项7.1 有状态摄入的限制只能拉取 Kafka 保留期内的 timeseries aspectsKafka 默认保留 90 天更早的时序数据无法通过本源获取不检测硬删除hard timeseries deletions例如通过 CLI 的 datahub delete 命令删除的数据不会反映到目标实例删除后源端数据在目标实例中仍然存在同一 createdon 时间戳的大量 aspect当某一时间戳下的 aspect 数量很大时有状态摄入无法在该时间戳内部保存部分检查点后续运行时该时间戳的全部 aspect 都会被重新摄入。7.2 一般限制模块行为受限于源平台暴露的 API、权限与元数据范围官方文档统一表述具体以能力表中的说明为准。八、性能与规模化建议针对大规模迁移官方文档给出如下建议确保metadata_aspect_v2.createdon已建索引timeIndex这是数据库读取性能的基石在目标端启用异步摄入async ingestion避免同步写入成为瓶颈使用独立消费者mae-consumer 与 mce-consumer迁移大量数据时考虑扩展消费者副本数增加 GMS pod 数量提升冗余度与对节点驱逐的韧性提高 Elasticsearch 线程数迁移大量数据时可通过ELASTICSEARCH_THREAD_COUNT环境变量增加 ES 的线程数。此外源码层面 datahub_database_reader.py 已内置服务端游标 流式读取与database_query_batch_size分批机制能有效控制迁移过程中的内存占用。九、常见问题排查Troubleshooting如果摄入失败官方文档建议按以下顺序排查详见 datahub_post.md先校验凭据、权限、连通性、范围过滤配置是否正确再检查摄入日志中的源特定错误根据错误信息调整配置后重试。结合源码可补充两个典型场景配置中未提供任何数据源数据库/Kafka/API 均为空时配置校验会直接抛出ValueError日志会提示 Your current config will not ingest any data——优先检查 recipe 中是否填入了database_connection或kafka_connection数据库方言为 MySQL 但 scheme 不是mysqlpymysql时同样会在校验阶段报错请按 datahub_recipe.yml 中的写法修正。十、小结DataHub Source 是 DataHub 生态中以 DataHub 为源的迁移型连接器它通过数据库直连 Kafka 消费 可选 API 三条通道把源实例的 versioned aspects 与 timeseries aspects 按时间顺序完整搬运到目标实例并依靠createdon/Kafka offset 检查点实现可靠续跑。在规划迁移时重点把握三点建好createdon索引保证数据库读取性能、按需配置urn_pattern与exclude_aspects过滤实例专属元数据、利用pipeline_name与ignore_old_state控制全量/增量重放。相关实现细节可进一步阅读 datahub_source.py、config.py 与 datahub_recipe.yml。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考