4.3 自定义连接器开发


4.3 自定义连接器开发

本节摘要:当源系统没有官方连接器——公司自研的数据平台、老旧的报表系统、某个业务中台——自定义连接器是唯一的路。本节拆解官方摄入框架的插件契约:一个来源类、一个配置类、一套 yield 语义,然后用一个可运行的精简示例走通开发全程,并给出增量采集、错误上报、测试打包的工程要点。读完你应当能在两三天内为任意"能查到元数据"的系统写出可维护的连接器。本节承接 4.2 的标准操作,是接入能力的进阶篇。

先判断要不要自研

动笔前先过两道检查。第一道:官方连接器列表里真的没有吗?很多源以"通用数据库连接器"或"自定义 SQL 源"的形式被覆盖,配个查询就能用,不必写代码。第二道:这个源的元数据值得长期维护吗?一次性迁移的存量系统,用脚本临时推一拨数据更划算——自研连接器的成本不在写,而在源系统升级后的持续适配。

两道检查都指向自研,再开工。自研的产出物是一个符合框架契约的 Python 包:框架负责配置解析、命令行入口、执行调度与报告汇总,你只实现"怎么连源"与"产出什么元数据"两件事。

框架契约:三个必答的问题

官方摄入框架把连接器抽象成对三个问题的回答。你是谁:来源类声明自己的类型名,配置文件里的 source.type 就靠它路由到你的代码。你需要什么:配置类用声明式结构定义配置项与校验规则,框架自动完成解析与报错。你产出什么:执行入口按 yield 语义逐条吐出元数据事件,框架负责批量打包、提交 GMS、统计成败。

这套契约的好处是 standardized 的工程红利:错误上报、重试、报告格式、命令行体验全部白得。自研代码的职责边界因此非常清晰——下面的骨架展示了这份最小实现:

from datahub.ingestion.api.source import Source, SourceReport from datahub.ingestion.api.common import PipelineContext from datahub.emitter.mce_builder import make_dataset_urn_with_platform from datahub.emitter.mcp import MetadataChangeProposalWrapper import datahub.emitter.mce_builder as builder from mycompany.lakehouse.meta_client import LakeHouseMetaClient class LakeHouseConfig(dict): """配置声明:endpoint 与 token 由环境变量注入""" endpoint: str token: str project_allow: list class LakeHouseSource(Source): """自研湖仓平台的元数据来源""" @classmethod def create(cls, config_dict, ctx: PipelineContext): return cls(config=LakeHouseConfig(config_dict), ctx=ctx) def get_workunits(self): client = LakeHouseMetaClient( endpoint=self.config.endpoint, token=self.config.token, ) for project in client.list_projects(): if project.name not in self.config.project_allow: continue for table in client.list_tables(project.id): dataset_urn = make_dataset_urn_with_platform( platform="lakehouse", name=f"{project.name}.{table.name}", env="PROD", ) props = builder.make_dataset_properties_aspect( description=table.comment or "", customProperties={ "storage_size": str(table.size_bytes), "partition_key": table.partition_key or "", }, ) yield MetadataChangeProposalWrapper( entityUrn=dataset_urn, aspect=props, ).as_workunit() def get_report(self) -> SourceReport: return self.report def close(self): pass

这段骨架只有几十行,但每个位置都有讲究。create 方法是框架的工厂入口,配置字典在这里变成类型化配置。get_workunits 是心脏:逐条 yield 出工作单元,注意每条元数据一个事件而不是攒一个巨大列表——框架的批量与背压机制依赖这个节奏。URN 构造函数指明了 3.1 节强调的三要素:平台名 lakehouse 受控注册、名称用项目加点分的原生坐标、环境段显式 PROD。customProperties 是 3.2 守则一的应用:体量与分区键这类简单属性不配拥有正式 Aspect,键值对足够。

增量采集:从全量到增量的升级路径

第一版连接器全量采集即可——先让数据进来,再谈效率。全量版稳定运行后,按两个信息源升级增量:源系统有变更时间戳或事件日志的,按水位线过滤(上次运行时间之后有变更的表才拉);没有的,用元数据的修改时间做近似过滤。增量带来的一个隐藏义务:删除检测。源系统里下线的表,全量模式下会被覆盖语义自然清理(配合平台侧的软删除标记),增量模式下永远不会出现在拉取结果里,需要定期跑一轮全量对账补删,否则地图上会积累僵尸实体。

另一个工程要点是分页与限速。湖仓类系统的元数据接口大多有分页与频控,连接器里老老实实按页翻、按限速走,别仗着运维身份硬拉——把源平台的元数据服务打挂,你会同时得罪两个团队。

错误上报与测试

单条失败不中断全局。某个表的详情接口报错时,yield 循环不该被打断:捕获异常、记入报告对象、继续下一条。报告里的失败条目会出现在摄入报告末尾,格式与官方连接器一致——4.4 的排错流程因此对自研连接器同样适用。

测试用录制回放。对着测试环境的源系统录一份接口响应存成样例文件,单元测试里让伪客户端返回样例数据,断言 yield 出的工作单元数量与 URN 格式正确。源系统升级时,更新样例文件重跑测试,五分钟内可知适配成本。这个投入在第一次源系统大版本升级时就会回本。

打包用标准轮子。连接器连同配置注册写成标准 Python 包,内部源里发一个私有版本,摄入机的安装命令与官方插件完全一致。自研连接器最忌讳的做法是"直接在摄入机上改框架代码"——升级即丢失,交接即黑盒。

⚠️ 自研连接器写出来的 URN 坐标,决定它在地图上落在哪个街区。平台名、环境段、命名分隔符在写第一行代码前就要与团队对齐——URN 规则一旦有数据入库就难以更改,返工等于重接。

本节要点回顾

  • 两道检查:先查官方连接器与通用源是否已覆盖,再评估源是否值得长期适配,然后才自研。
  • 框架契约:来源类答"你是谁"、配置类答"你需要什么"、yield 工作单元答"你产出什么",工程设施白得。
  • yield 纪律:一条元数据一个事件,让框架的批量与背压机制工作。
  • 增量三件套:水位线过滤、删除检测靠定期全量对账、分页限速尊重源平台。
  • 工程底线:单条失败记报告不中断、样例回放做测试、打包成标准轮子,绝不在摄入机上改框架。

自研与标准路径都已铺好。但接入总会出岔子——下一节是一份按症状组织的排错实录,收录了接入现场最常撞上的故障与处置。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U