:从零编写自定义 Connector 的完整指南)
为 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 能够从你的目标数据系统中抽取元数据。文档给出了两种完全不同的开发路径把自定义源贡献回 DataHub 项目需要完整走完后续的 19 步含依赖声明、插件注册、测试、文档、平台 Logo、前端 UI 接入等保证质量与可维护性。仅为自己的使用场景编写、暂不贡献回上游可以跳过第 48 步依赖与文档化相关步骤直接参考 如何在不 fork DataHub 的情况下使用自定义摄取源 的说明将自定义源作为独立 Python 包安装使用。无论走哪条路径前置条件都是先按照 metadata-ingestion 开发指南 完成本地开发环境搭建。该指南要求在宿主环境安装 Python 3.9 与 Java 17Gradle 对 Java 版本有严格要求然后在仓库根目录执行cd metadata-ingestion ../gradlew :metadata-ingestion:installDev source venv/bin/activate datahub version # 应输出 DataHub CLI version: unavailable (installed in develop mode)提示DataHub 官方推荐使用DataHub SkillsAI 辅助框架加速 Connector 开发可通过 DataHub Skills 指南 了解如何在几分钟内从简单配置生成生产级 Connector。本文专注手工实现的完整原理与细节。1. 建立配置模型Configuration ModelDataHub 的摄取框架使用 pydantic它继承了StatefulIngestionConfigBase最终仍以ConfigModel为基类用 pydanticField声明了path、file_extension、read_mode、aspect、count_all_before_starting等字段。配置模型的价值不止于运行时解析还承担着自动生成文档的职责。DataHub 遵循 pydantic 惯例通过字段的description属性书写富文档。例如from pydantic import Field from datahub.api.configuration.common import ConfigModel class LookerAPIConfig(ConfigModel): client_id: str Field(descriptionLooker API client id.) client_secret: str Field(descriptionLooker API client secret.) base_url: str Field( descriptionUrl 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( defaultNone, descriptionPopulates 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_fieldvalidator 只能抛出ValueError、TypeError或AssertionError内部专用配置标记hidden_from_docs。从 decorators.py 的实现可以看到config_class装饰器会给源类注入get_config_class()方法并在类未自定义create()时自动生成一个基于该配置类的create()类方法——这正是配置模型与源类之间“约定优于配置”的绑定机制。2. 设置 Reporter运行报告器Reporter 接口允许源在运行过程中上报统计信息、警告、失败及其他运行细节。默认使用SourceReport类部分源会继承并扩展它以加入领域专属字段。仓库中最典型的扩展示例是 file.py 中的FileSourceReport它继承了StaleEntityRemovalSourceReport新增了total_num_files、num_files_completed、percentage_completion、estimated_time_to_completion_in_minutes、total_bytes_on_disk等字段并提供了add_deserialize_time、add_parse_time、add_count_time、append_total_bytes_on_disk等统计方法与compute_stats()完成进度计算按已读字节数占总字节数的比例估算剩余时间。在现代 DataHub 中报告体系还支持结构化日志StructuredLogs见 source.py按StructuredLogLevelINFO/WARN/ERROR分层收集日志条目report_log方法对同标题同消息的条目做聚合去重并可用DATAHUB_REPORT_*_SAMPLE_SIZE环境变量控制抽样规模。这些报告信息最终会呈现在datahub ingest命令的输出以及 DataHub 前端的摄取运行详情中。3. 实现 Source 主类get_workunits_internal整个 Source 的核心是get_workunits_internal方法它产生一个元数据事件流——通常是 MCPMetadataChangeProposal对象——并将其包装进MetadataWorkUnit。官方推荐的参考实现同样是 file.py 的GenericFileSource。GenericFileSource展示了 Source 类应有的完整骨架__init__保存PipelineContext、配置并实例化self.reportcreate(cls, config_dict, ctx)类方法通常由config_class自动生成用FileSourceConfig.model_validate(config_dict)解析配置get_workunits_internal(self) - Iterable[MetadataWorkUnit]遍历文件对每个解析出的对象按类型分派MetadataChangeProposalWrapper走MetadataWorkUnit(id, mcpobj)原生MetadataChangeProposal走mcp_rawMetadataChangeEvent走mceobjget_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 methodsmce_builder.py用于常见操作——例如构造 URNmake_dataset_urn等、构造 Ownership/AuditStamp 等 Metadata 对象编写 Source 时应优先复用这些辅助方法避免手写底层模型。4. 声明依赖仅贡献上游时需要第 48 步仅在打算把源贡献回 DataHub 项目时需要执行。在 setup.py 的plugins变量中声明源所需的 pip 依赖。该变量是一个Dict[str, Set[str]]键为插件名值为依赖字符串集合例如plugins: Dict[str, Set[str]] { # Source plugins aerospike: {aerospike15.0.0,20.0.0}, athena: sql_common | { PyAthena[SQLAlchemy]2.6.0,3.0.0, sqlalchemy-bigquery1.5.0,2.0.0, tenacity!8.4.0,9.0.0, }, ... }sql_common、kafka_common、aws_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, ], }注册后会有两个立竿见影的效果运行datahub check plugins时新源会出现在已发现的插件列表中在 recipe 中可以直接使用注册的短别名作为source.type例如type: my-source而不必写完整模块路径。从源码看entry_points还承载了datahub.ingestion.transformer.plugins、datahub.token_provider.plugins等其他扩展点说明“插件化注册”是贯穿整个acryl-datahub包的设计主线console_scripts则把datahubCLI 入口绑定到 entrypoints.py 的main函数。6. 编写测试测试放在metadata-ingestion的tests目录框架使用 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 -PtestFiletests/unit/test_bigquery_source.py对于涉及“黄金文件golden files”快照断言的集成测试变更行为后可用--update-golden-files重新生成基线例如pytest tests/integration/dbt/test_dbt.py --update-golden-files此外仓库还提供datahub check与datahub ingest等 CLI 子命令用于手工验证源的行为。7. 编写文档让 Connector 自动生成与站点化DataHub 的 Connector 文档分为“自动生成”与“手写定制”两部分两者结合形成最终的官方文档页面。7.1 用装饰器让源类自描述在源类上使用以下装饰器定义见 decorators.py文档生成器即可自动收集元信息platform_name(File)声明该源产出元数据的平台名优先使用人类可读的平台名如 BigQuery 而不是 bigquery。装饰器会同时自动派生平台 id小写并替换空格为-并支持id与doc_order参数。config_class(FileSourceConfig)声明源使用的配置类。support_status(SupportStatus.GA)声明连接器支持状态取值为SupportStatus枚举——ALPHA早期集成社区维护可无通知变更、BETA由 Ingestion 团队维护但生产采用有限、GA团队维护且被广泛生产采用、期望稳定、UNKNOWN作者未声明。这些状态会以徽章形式呈现在文档中。capability(SourceCapability.XXX, 描述, supported...)声明连接器支持或明确不支持的关键能力。SourceCapability枚举见 source.py涵盖PLATFORM_INSTANCE、DOMAINS、DATA_PROFILING、USAGE_STATS、DESCRIPTIONS、LINEAGE_COARSE/LINEAGE_FINE、OWNERSHIP、DELETION_DETECTION、TAGS、SCHEMA_METADATA、CONTAINERS、TEST_CONNECTION、GLOSSARY_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, supportedFalse, ) 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 判定谓词”的映射表函数遍历该表并返回第一个匹配的平台名无匹配则返回externaldef 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),仅供参考