ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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 已持久化的元数据原样读出再以 MetadataChangeProposalWrapperMCP的形式作为 WorkUnit 交给目标端 sink 写入。该连接器从两个位置拉取数据数据来源数据类型说明DataHub 数据库版本化 AspectVersioned Aspects存储在源实例 GMS 后端数据库的metadata_aspect_v2表中DataHub Kafka时序 AspectTimeseries Aspects通过 Kafka 主题中的 Metadata Change LogMCL 事件读取读取顺序是固定的先完整读取数据库中的版本化 Aspect再消费 Kafka 中的时序 MCL。这一顺序在 get_workunits_internal 中体现得十分清楚_get_database_workunits执行完毕后才会进入_get_kafka_workunits。二、防止无限运行的 stop_time 机制文档明确指出一个关键设计为防止该源永远跑下去连接器不会消费 ingestion job 启动之后才产生的数据。每次运行开始时get_workunits_internal会立即记录一个stop_timeself.report.stop_time datetime.now(tztimezone.utc) logger.info(fIngesting DataHub metadata up until {self.report.stop_time})该时间点被同时用于数据库侧作为查询时间范围的上界createdon stop_timeKafka 侧消费时若遇到 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在启动迁移前需要满足以下条件网络连通性能够访问源 DataHub 实例所在的环境有效的认证凭据具备读取元数据 API 的权限只读权限拥有本模块所需元数据 API 的读取权限。特别地该连接器必须直连源 DataHub 实例的以下三部分基础设施而非仅通过 GMS HTTP API数据库存储版本化 AspectKafka broker存储时序 MCLKafka Schema Registry用于反序列化 Avro 格式的 MCL备注仓库源码还提供了一条可选的pull_from_datahub_api路径见下文此时才要求配置datahub_api即 GMS REST 连接而非直连数据库。五、Recipe 配置详解从连接串到增量策略datahub源的配置类为DataHubSourceConfigconfig.py继承了StatefulIngestionConfigBase。一个典型的迁移 Recipe 结构如下source: type: datahub config: # --- 数据来源一数据库版本化 Aspect --- database_connection: scheme: mysqlpymysql # MySQL 必须使用 mysqlpymysql 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-token5.1 核心配置参数以下参数均来自DataHubSourceConfig默认值以源码为准参数默认值作用database_connectionNone数据库连接配置SQLAlchemyConnectionConfig提供版本化 Aspect 数据源kafka_connectionNoneKafka 连接配置KafkaConsumerConnectionConfig提供时序 MCL 数据源include_all_versionsFalse是否包含每个 Aspect 的全部历史版本关闭时只取最新版本version 0include_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_ingestionenabledTrue有状态增量摄入该源默认开启与多数源不同commit_state_interval1000每处理多少条记录提交一次检查点commit_with_parse_errorsFalse出现解析错误时是否仍推进 createdon/offset 检查点pull_from_datahub_apiFalse隐藏参数改为通过 DataHub API 拉取版本化 Aspectmax_workers5 * cpu_count隐藏参数DataHub API 拉取时的线程数urn_patterndeny 环境专属 URN见下URN 过滤模式drop_duplicate_schema_fieldsFalse是否丢弃schemaMetadata中的重复字段路径源库存在重复、目标端有服务端去重时适用query_timeoutNone数据库单次查询超时秒preserve_system_metadataTrue是否复制源系统的 systemMetadata5.2 三个值得注意的设计点1至少配置一个数据来源否则拒绝运行。配置类通过模型校验器强制要求database_connection、kafka_connection、pull_from_datahub_api三者至少提供一个否则抛出ValueError(Your current config will not ingest any data...)config.py。文档推荐两者都配ideally both以保证版本化与时序 Aspect 都能迁移。2MySQL 连接串有硬性约束。当 scheme 包含mysql时必须是mysqlpymysql否则校验失败config.py。3默认排除环境专属 URN。urn_pattern默认 deny 以下四类 URNconfig.pyurn: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做自关联把statusAspectversion 0中的removed字段提取出来用于软删除过滤。该 JSON 提取表达式按方言区分PostgreSQL 用((metadata::json)-removed)::boolean其他如 MySQL用JSON_EXTRACT(metadata, $.removed)datahub_database_reader.py支持 PostgreSQL/MySQL/MariaDB 的服务端游标流式读取stream_resultsTrueyield_perbatch_size并可按query_timeout设置statement_timeoutPG或max_execution_timeMySQL当include_all_versionsTrue时启用VersionOrderer同一createdon时间戳下版本 0 的 Aspect 被延后到最后输出保证最新版本后写、旧版本先写的稳定顺序datahub_database_reader.py解析时通过ASPECT_MAP将 JSON 元数据还原为对应的 Aspect 类对象preserve_system_metadataTrue时会复制 systemMetadata并剥离isNoOp标记解析失败计入num_database_parse_errors与database_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 客户端以DeserializingConsumerAvroDeserializer消费主题反序列化依赖Schema Registry这就是前置条件要求直连 schema registry 的原因关键消费参数auto.offset.resetearliest从头开始、enable.auto.commitFalse手动管理进度配合有状态检查点consumer group 固定为datahub_source-{pipeline_name}datahub_kafka_reader.pyon_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_apiTrue时连接器改用 DataHub Graph API 拉取版本化 Aspectdatahub_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默认开启。状态由StatefulDataHubIngestionHandlerstate.py管理检查点结构为class DataHubIngestionState(CheckpointStateBase): database_createdon_ts: NonNegativeInt 0 # 数据库侧进度毫秒时间戳 kafka_offsets: Dict[int, NonNegativeInt] # Kafka 各分区已消费的 offset运行逻辑datahub_source.py从上次检查点恢复from_createdon数据库与from_offsetsKafka数据库每处理一条记录就更新database_createdon_ts每commit_state_interval默认 1000条提交一次Kafka 每消费一条 MCL 就按offset 1更新对应分区进度只有当没有解析错误或显式设置commit_with_parse_errorsTrue时才推进检查点避免带病提交导致数据永久丢失。这保证了迁移任务中途失败后可安全重跑数据库从上次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_aspects因exclude_aspects跳过的 MCL 数num_timeseries_deletions_dropped被丢弃的时序 DELETE 变更数num_timeseries_soft_deleted_aspects_dropped因实体软删除而被丢弃的时序 Aspect 数其中两个丢弃指标对应 datahub_source.py 的过滤逻辑ChangeTypeClass.DELETE的时序变更被跳过当include_soft_deleted_entitiesFalse时先从数据库查出软删除 URN 列表再从 Kafka 流中过滤掉这些实体的时序 Aspect。九、迁移实操建议与注意事项先建索引再迁移确认源库存在timeIndexcreatedon索引否则先执行CREATE INDEX timeIndex ON metadata_aspect_v2 (createdon);保持默认 urn_pattern保留对 Ingestion Source / Secret / Settings 等环境专属 URN 的排除避免复制加密凭据优先双连接配置同时配置database_connection与kafka_connection才能完整迁移版本化与时序两类 Aspect若只配一个另一类会被跳过日志中会出现 Skipping ingestion of versioned aspects... 等提示注意 MySQL 方言约束MySQL 场景必须写scheme: mysqlpymysql保留默认 exclude_aspectsdatahubIngestionRunSummary、datahubIngestionCheckpoint、testResults属于运行期状态迁移它们没有业务价值默认排除是合理选择利用有状态增量由于默认开启 stateful ingestion迁移失败后直接重跑同一 Recipe 即可续传无需清空目标端核对运行报告迁移完成后对照上文表格中的各项计数确认无大量解析错误、且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.pyKafka 消费实现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),仅供参考
返回列表