为 DataHub 添加新的元数据摄取源(Metadata Ingestion Source):从零编写自定义 Connector 的完整指南
2026/9/18 6:13:18 网站建设 项目流程

为 DataHub 添加新的元数据摄取源(Metadata Ingestion Source):从零编写自定义 Connector 的完整指南

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

本文是一份面向开发者的实战指南,基于 DataHub 开源仓库metadata-ingestion模块中关于添加元数据摄取源的官方文档编写,并辅以仓库源码级证据进行纵深讲解。文章覆盖从配置模型、Reporter、Source 主类实现,到依赖声明、插件注册、测试、文档生成、SQLAlchemy URI 映射与前端 UI 接入的完整链路,既适合想要把自定义源贡献回 DataHub 社区的开发者,也适合只为自己团队内部使用而编写私有 Connector 的工程师。读完本文,你将能够独立编写、测试、打包并在 recipe 中运行一个全新的 DataHub 摄取源。

写在前面:两种添加路径,先选对方向

在 DataHub 中,为metadata-ingestion框架添加一个元数据摄取源(Source),本质上就是实现一套标准的 Python 接口,让 DataHub 能够从你的目标数据系统中抽取元数据。文档给出了两种完全不同的开发路径:

  1. 把自定义源贡献回 DataHub 项目:需要完整走完后续的 1~9 步(含依赖声明、插件注册、测试、文档、平台 Logo、前端 UI 接入等),保证质量与可维护性。
  2. 仅为自己的使用场景编写、暂不贡献回上游:可以跳过第 4~8 步(依赖与文档化相关步骤),直接参考 如何在不 fork DataHub 的情况下使用自定义摄取源 的说明,将自定义源作为独立 Python 包安装使用。

无论走哪条路径,前置条件都是先按照 metadata-ingestion 开发指南 完成本地开发环境搭建。该指南要求在宿主环境安装 Python 3.9+ 与 Java 17(Gradle 对 Java 版本有严格要求),然后在仓库根目录执行:

cd metadata-ingestion ../gradlew :metadata-ingestion:installDev source venv/bin/activate datahub version # 应输出 "DataHub CLI version: unavailable (installed in develop mode)"

提示:DataHub 官方推荐使用DataHub Skills(AI 辅助框架)加速 Connector 开发,可通过 DataHub Skills 指南 了解如何在几分钟内从简单配置生成生产级 Connector。本文专注手工实现的完整原理与细节。

1. 建立配置模型(Configuration Model)

DataHub 的摄取框架使用 pydantic,它继承了StatefulIngestionConfigBase(最终仍以ConfigModel为基类),用 pydanticField声明了pathfile_extensionread_modeaspectcount_all_before_starting等字段。

配置模型的价值不止于运行时解析,还承担着自动生成文档的职责。DataHub 遵循 pydantic 惯例,通过字段的description属性书写富文档。例如:

from pydantic import Field from datahub.api.configuration.common import ConfigModel class LookerAPIConfig(ConfigModel): client_id: str = Field(description="Looker API client id.") client_secret: str = Field(description="Looker API client secret.") base_url: str = Field( description="Url to your Looker instance: `https://company.looker.com:19999` or `https://looker.company.com`, or similar. Used for making API calls to Looker and constructing clickable dashboard and chart urls." ) transport_options: Optional[TransportOptionsConfig] = Field( default=None, description="Populates the [TransportOptions](https://github.com/looker-open-source/sdk-codegen/blob/94d6047a0d52912ac082eb91616c1e7c379ab262/python/looker_sdk/rtl/transport.py#L70) struct for looker client", )

这些description会被文档生成器(docGen)提取,渲染为该 Connector 的配置文档。需要说明的是,字段级文档目前还不支持内联 Markdown 或代码片段(仅支持链接语法)。

配置文件编写的仓库级准则

除官方文档外,developing.md 的 “Guidelines for Ingestion Configs” 一节给出了更具操作性的约束,编写配置类时建议一并遵守:

  • 命名与源系统术语一致:例如 Snowflake 源不应出现host_port,而应使用account_id;优先使用client_id/tenant_id而非含义模糊的id,使用access_secret而非secret
  • 过滤类配置统一用AllowDenyPatterns:需要过滤列表时使用该模式,且模式总是作用于实体的全限定名,命名规范为*_pattern(如table_pattern)。
  • 避免*_only式配置:用profile_table_level/profile_column_level代替profile_table_level_only
  • 所有配置必须带description;设置合理的默认值;不要在内置行为上叠加默认值(例如不应在schema_pattern默认 deny 中硬编码information_schema,而应由源实现自动过滤)。
  • 编码细节:密码、token 等敏感字段用SecretStr;字段重命名用pydantic_renamed_field辅助函数,字段废弃用pydantic_removed_field;validator 只能抛出ValueErrorTypeErrorAssertionError;内部专用配置标记hidden_from_docs

从 decorators.py 的实现可以看到,@config_class装饰器会给源类注入get_config_class()方法,并在类未自定义create()时,自动生成一个基于该配置类的create()类方法——这正是配置模型与源类之间“约定优于配置”的绑定机制。

2. 设置 Reporter(运行报告器)

Reporter 接口允许源在运行过程中上报统计信息、警告、失败及其他运行细节。默认使用SourceReport类,部分源会继承并扩展它以加入领域专属字段。

仓库中最典型的扩展示例是 file.py 中的FileSourceReport:它继承了StaleEntityRemovalSourceReport,新增了total_num_filesnum_files_completedpercentage_completionestimated_time_to_completion_in_minutestotal_bytes_on_disk等字段,并提供了add_deserialize_timeadd_parse_timeadd_count_timeappend_total_bytes_on_disk等统计方法与compute_stats()完成进度计算(按已读字节数占总字节数的比例估算剩余时间)。

在现代 DataHub 中,报告体系还支持结构化日志StructuredLogs(见 source.py)按StructuredLogLevel(INFO/WARN/ERROR)分层收集日志条目,report_log方法对同标题同消息的条目做聚合去重,并可用DATAHUB_REPORT_*_SAMPLE_SIZE环境变量控制抽样规模。这些报告信息最终会呈现在datahub ingest命令的输出以及 DataHub 前端的摄取运行详情中。

3. 实现 Source 主类:get_workunits_internal

整个 Source 的核心是get_workunits_internal方法,它产生一个元数据事件流——通常是 MCP(MetadataChangeProposal)对象——并将其包装进MetadataWorkUnit。官方推荐的参考实现同样是 file.py 的GenericFileSource

GenericFileSource展示了 Source 类应有的完整骨架:

  • __init__保存PipelineContext、配置,并实例化self.report
  • create(cls, config_dict, ctx)类方法(通常由@config_class自动生成)用FileSourceConfig.model_validate(config_dict)解析配置;
  • get_workunits_internal(self) -> Iterable[MetadataWorkUnit]遍历文件,对每个解析出的对象按类型分派:MetadataChangeProposalWrapperMetadataWorkUnit(id, mcp=obj),原生MetadataChangeProposalmcp_rawMetadataChangeEventmce=obj
  • get_report()返回报告实例;
  • 可选的test_connection静态方法实现连接测试能力(见 file.py),对应@capability(SourceCapability.TEST_CONNECTION, ...)

一个值得注意的细节:get_workunits_internal中通过fs_registry按路径 schema(如file://s3://等)解析文件系统实现,在AUTO读取模式下,超过 100MB(_minsize_for_streaming_mode_in_bytes = 100 * 1000 * 1000)的文件自动切换为流式解析(ijson),小文件则整批json.load——这说明 DataHub 的源框架对大数据量摄取做了充分的性能考量。

元数据事件模型从哪来?

MetadataChangeEventClass等元数据模型定义在由代码生成得到的metadata-ingestion/src/datahub/metadata/schema_classes.py(该目录文件为构建期生成、不入库)。此外,仓库还提供了一批convenience methods(mce_builder.py)用于常见操作——例如构造 URN(make_dataset_urn等)、构造 Ownership/AuditStamp 等 Metadata 对象,编写 Source 时应优先复用这些辅助方法,避免手写底层模型。

4. 声明依赖(仅贡献上游时需要)

第 4~8 步仅在打算把源贡献回 DataHub 项目时需要执行。

在 setup.py 的plugins变量中声明源所需的 pip 依赖。该变量是一个Dict[str, Set[str]],键为插件名,值为依赖字符串集合,例如:

plugins: Dict[str, Set[str]] = { # Source plugins "aerospike": {"aerospike>=15.0.0,<20.0.0"}, "athena": sql_common | { "PyAthena[SQLAlchemy]>=2.6.0,<3.0.0", "sqlalchemy-bigquery>=1.5.0,<2.0.0", "tenacity!=8.4.0,<9.0.0", }, ... }

sql_commonkafka_commonaws_common等是仓库预定义的公共依赖集合,可以通过|运算符组合。在依赖管理上有两条硬性要求(见 developing.md):

  • 尽量不锁定版本;确需限制时使用区间(如>=1.2.3,<2.0.0)或负向约束(如!=1.2.7),且每个上界/负向约束都必须附注释说明原因;对频繁破坏性变更的包(如 Great Expectations、Airflow)可加“防御性上界”并定期复核放宽。
  • 修改setup.py后需要重新生成锁文件:运行../gradlew :metadata-ingestion:updateLockFile(执行setup.py → pyproject.toml → uv.lock → constraints.txt的完整链路),并用../gradlew :metadata-ingestion:checkLockFile校验;CI 的check任务会自动执行该校验,PR 中带过期生成文件会直接失败。

5. 启用可发现性:注册 entry point

在 setup.py 的entry_points变量中,将新源注册到datahub.ingestion.source.plugins分组下:

entry_points = { "console_scripts": ["datahub = datahub.entrypoints:main"], "datahub.ingestion.source.plugins": [ "file = datahub.ingestion.source.file:GenericFileSource", "bigquery = datahub.ingestion.source.bigquery_v2.bigquery:BigqueryV2Source", "kafka = datahub.ingestion.source.kafka.kafka:KafkaSource", # 新增:别名 = 模块路径:类名 "my-source = datahub.ingestion.source.my_source:MySource", ], }

注册后会有两个立竿见影的效果:

  1. 运行datahub check plugins时,新源会出现在已发现的插件列表中;
  2. 在 recipe 中可以直接使用注册的短别名作为source.type,例如type: my-source,而不必写完整模块路径。

从源码看,entry_points还承载了datahub.ingestion.transformer.pluginsdatahub.token_provider.plugins等其他扩展点,说明“插件化注册”是贯穿整个acryl-datahub包的设计主线;console_scripts则把datahubCLI 入口绑定到 entrypoints.py 的main函数。

6. 编写测试

测试放在metadata-ingestiontests目录,框架使用 pytest。仓库将测试分为两类(见 developing.md):

  • 单元测试pytest -m 'not integration'
  • 基于 Docker 的集成测试pytest -m 'integration'

常用测试命令速查:

../gradlew :metadata-ingestion:installDevTest # 安装全部 dev/test 依赖 pytest -vv # 运行全部测试 pytest -m 'not integration' # 仅单元测试 # 通过 gradle 运行 ../gradlew :metadata-ingestion:testQuick ../gradlew :metadata-ingestion:testFull ../gradlew :metadata-ingestion:testSingle -PtestFile=tests/unit/test_bigquery_source.py

对于涉及“黄金文件(golden files)”快照断言的集成测试,变更行为后可用--update-golden-files重新生成基线,例如:

pytest tests/integration/dbt/test_dbt.py --update-golden-files

此外仓库还提供datahub checkdatahub ingest等 CLI 子命令用于手工验证源的行为。

7. 编写文档:让 Connector 自动生成与站点化

DataHub 的 Connector 文档分为“自动生成”与“手写定制”两部分,两者结合形成最终的官方文档页面。

7.1 用装饰器让源类自描述

在源类上使用以下装饰器(定义见 decorators.py),文档生成器即可自动收集元信息:

  • @platform_name("File"):声明该源产出元数据的平台名,优先使用人类可读的平台名(如 BigQuery 而不是 bigquery)。装饰器会同时自动派生平台 id(小写并替换空格为-),并支持iddoc_order参数。
  • @config_class(FileSourceConfig):声明源使用的配置类。
  • @support_status(SupportStatus.GA):声明连接器支持状态,取值为SupportStatus枚举——ALPHA(早期集成,社区维护,可无通知变更)、BETA(由 Ingestion 团队维护但生产采用有限)、GA(团队维护且被广泛生产采用、期望稳定)、UNKNOWN(作者未声明)。这些状态会以徽章形式呈现在文档中。
  • @capability(SourceCapability.XXX, "描述", supported=...):声明连接器支持(或明确不支持)的关键能力。SourceCapability枚举(见 source.py)涵盖PLATFORM_INSTANCEDOMAINSDATA_PROFILINGUSAGE_STATSDESCRIPTIONSLINEAGE_COARSE/LINEAGE_FINEOWNERSHIPDELETION_DETECTIONTAGSSCHEMA_METADATACONTAINERSTEST_CONNECTIONGLOSSARY_TERMS等能力项。
  • 在类的docstring中书写富文档,支持 Markdown(官方文档中 docstring 里的:::提示块同样会被渲染)。

官方给出的完整示例:

from datahub.ingestion.api.decorators import ( SourceCapability, SupportStatus, capability, config_class, platform_name, support_status, ) @platform_name("File") @support_status(SupportStatus.GA) @config_class(FileSourceConfig) @capability( SourceCapability.PLATFORM_INSTANCE, "File based ingestion does not support platform instances", supported=False, ) @capability(SourceCapability.DOMAINS, "Enabled by default") @capability(SourceCapability.DATA_PROFILING, "Optionally enabled via configuration") @capability(SourceCapability.DESCRIPTIONS, "Enabled by default") @capability(SourceCapability.LINEAGE_COARSE, "Enabled by default") class FileSource(Source): """ The File Source can be used to produce all kinds of metadata from a generic metadata events file. :::note Events in this file can be in MCE form or MCP form. ::: """ ... source code goes here

仓库中 file.py 的实际实现与此模式完全一致(GenericFileSource标注了@platform_name("Metadata File")@support_status(SupportStatus.GA)TEST_CONNECTION能力)。

7.2 手写定制文档的组织方式

  • 复制 source-docs-template.md 并编辑相关内容;
  • 文档命名为<plugin>.md,放置到metadata-ingestion/docs/sources/<platform>/<plugin>.md(例如 Kafka 平台的文档位于metadata-ingestion/docs/sources/kafka/kafka.md);
  • 为插件提供一份 quickstart recipe,放置到metadata-ingestion/docs/sources/<platform>/<plugin>_recipe.yml(例如metadata-ingestion/docs/sources/kafka/kafka_recipe.yml);
  • 跨插件的平台级文档写在metadata-ingestion/docs/sources/<platform>/README.md(例如 BigQuery 平台的跨插件文档位于metadata-ingestion/docs/sources/bigquery/README.md)。

7.3 在本地查看生成的文档

第一步:生成摄取文档。在仓库根目录执行:

./gradlew :metadata-ingestion:docGen

成功结束后会输出统计信息,类似:

Ingestion Documentation Generation Complete ############################################ { "source_platforms": { "discovered": 40, "generated": 40 }, "plugins": { "discovered": 47, "generated": 47, "failed": 0 } } ############################################

生成的文档文件位于仓库根目录下的./docs/generated/ingestion/sources,可以找到你的源的 Markdown 文件检查渲染效果是否符合预期。

第二步:构建完整文档站点。在仓库根目录执行:

./gradlew :docs-website:build

构建成功后,从docs-website模块启动本地预览:

cd docs-website npm run serve

然后访问http://localhost:3000(或 npm 实际监听的端口)。你的源会出现在左侧边栏Metadata Ingestion / Sources分类下。

8. 添加 SQLAlchemy URI 映射(如适用)

如果你的源基于 SQLAlchemy 实现(即存在 sqlalchemy 数据源),需要在get_platform_from_sqlalchemy_uri函数中加入该源的映射。

值得注意的是,在当前仓库中该函数位于 sqlalchemy_uri_mapper.py,其核心是一个PLATFORM_TO_SQLALCHEMY_URI_TESTER_MAP——平台名到“URI 判定谓词”的映射表,函数遍历该表并返回第一个匹配的平台名,无匹配则返回"external"

def get_platform_from_sqlalchemy_uri(sqlalchemy_uri: str) -> str: for platform, tester in PLATFORM_TO_SQLALCHEMY_URI_TESTER_MAP.items(): if tester(sqlalchemy_uri): return platform return "external"

该映射的实际调用方包括 Kafka Connect 的源/目标连接器(见 source_connectors.py 与 sink_connectors.py),它们用 JDBC URL 反查目标平台并据此构造 DataHub URN。因此,为 SQLAlchemy 系源添加映射,能确保其在 Kafka Connect 等跨平台场景中被正确识别;若判定为external则会触发 warning 日志。

9. 添加平台 Logo

将平台 Logo 图片放入 datahub-web-react/src/images 目录,并在启动引导配置 contenteditable="false">【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询