ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

DataHub Actions 自定义 Transformer 开发指南:从基类扩展到生产级事件转换

DataHub Actions 自定义 Transformer 开发指南:从基类扩展到生产级事件转换 DataHub Actions 自定义 Transformer 开发指南从基类扩展到生产级事件转换【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读本文基于 DataHub Actions 框架官方指南整理围绕如何编写一个自定义 Transformer事件转换器展开从继承Transformer基类、实现create与transform两个核心方法开始到将自定义 Transformer 安装到 Python 运行时、在 Actions 配置文件中引用并运行最后介绍将其贡献回 DataHub 核心库的规范流程。读完本文你将掌握 DataHub Actions 事件处理链路中转换/过滤环节的完整自定义方法并理解其底层注册与调用机制能够独立编写、调试和发布属于自己的 Transformer。全文代码与配置均来自当前仓库datahub-actions模块的真实实现可直接对照验证。一、Transformer 在 DataHub Actions 中的定位在动手写代码之前先明确 Transformer 在整个 DataHub Actions 框架中所处的位置。Actions 框架是一个事件驱动的自动化框架它订阅数据平台产生的各类元数据变更事件并据此执行用户自定义的动作Action。一个完整的 Action Pipeline 由四段组成Event Source从 Kafka、DataHub Cloud 等来源消费事件将原始事件封装为EventEnvelopeFilter基于事件类型与事件体对事件做前置筛选filters段Transformer对通过筛选的事件做转换或二次过滤可多个串联成链Action消费最终事件并执行具体动作。从 Pipeline 源码 可以看出Pipeline 对每个事件的处理顺序是先执行 filters再执行 transforms若转换结果非空则交给 action最后向事件源确认ack# 先应用过滤器 if not self._execute_filters(enveloped_event): return retval # 然后转换事件 transformed_event self._execute_transformers(enveloped_event) # 若事件非空则调用 action if transformed_event is not None: retval self._execute_action(transformed_event)而_execute_transformers的实现见 pipeline.py会按配置顺序依次调用每个 Transformer并把前一个 Transformer 的输出作为后一个的输入一旦某个 Transformer 返回None则整个转换链短路事件被丢弃curr_event enveloped_event for transformer in self.transforms: transformed_event self._execute_transformer(curr_event, transformer) if transformed_event is None: # 如果转换器过滤了事件则短路 return None curr_event transformed_event return curr_event因此Transformer 本质上是事件在被 Action 消费之前的最后一道加工与把关关卡既可以改写事件内容也可以丢弃过滤不符合条件的事件。开发者可以按需实现自己的语义与配置。二、Step 1定义自己的 Transformer2.1 Transformer 抽象基类所有 Transformer 都必须继承datahub_actions.transform.transformer.Transformer。该基类定义在 transformer.py是一个抽象类metaclassABCMeta开发者需要覆写两个抽象方法方法类型职责返回值create(cls, config, ctx)类方法classmethod以 Actions 配置文件中提取的自由格式配置字典为输入实例化 Transformer返回Transformer实例transform(self, event)实例方法每当收到一个事件时被调用承载转换核心逻辑返回转换后的EventEnvelope或返回None表示过滤掉该事件基类的完整签名如下from abc import ABCMeta, abstractmethod from typing import Optional from datahub_actions.event.event_envelope import EventEnvelope from datahub_actions.pipeline.pipeline_context import PipelineContext class Transformer(metaclassABCMeta): classmethod abstractmethod def create(cls, config: dict, ctx: PipelineContext) - Transformer: Factory method to create an instance of a Transformer pass abstractmethod def transform(self, event: EventEnvelope) - Optional[EventEnvelope]: Transform a single Event. This method returns an instance of EventEnvelope, or None if the event has been filtered. 其中config是 YAML 配置文件中transform.config段对应的字典若无配置则为空字典{}ctxPipelineContext携带 Pipeline 名称、DataHub Graph 客户端等运行时上下文可用于在 Transformer 内查询或写入 DataHub 元数据入参event是EventEnvelope。它定义在 event_envelope.py由event_type事件类型决定事件体结构、event事件本体如 MetadataChangeLogEvent 等和meta任意元数据字典三部分组成并提供as_json()/from_json()序列化接口。2.2 编写第一个 Transformer原文档给出一个极简示例CustomTransformer它创建时打印配置、收到事件时打印事件并原样返回no-op。完整代码如下# custom_transformer.py from datahub_actions.transform.transformer import Transformer from datahub_actions.event.event_envelope import EventEnvelope from datahub_actions.pipeline.pipeline_context import PipelineContext from typing import Optional class CustomTransformer(Transformer): classmethod def create(cls, config_dict: dict, ctx: PipelineContext) - Transformer: # Simply print the config_dict. print(config_dict) return cls(config_dict, ctx) def __init__(self, ctx: PipelineContext): self.ctx ctx def transform(self, event: EventEnvelope) - Optional[EventEnvelope]: # Simply print the received event. print(event) # And return the original event (no-op) return event注意原文档示例中__init__仅接收ctxcreate里以cls(config_dict, ctx)调用看似参数不一致这属于示例中的笔误——实际开发时请保持create与__init__的参数签名一致例如def __init__(self, config_dict: dict, ctx: PipelineContext)。2.3 编写有实际意义的转换逻辑真实的 Transformer 远不止打印事件常见的用途包括过滤根据事件类型或事件体字段丢弃无关事件返回None改写给事件附加额外的元数据修改meta或事件体字段聚合/分流把一类事件转换成另一类事件后继续向下传递。仓库内置的FilterTransformer是学习过滤型 Transformer的最佳参考实现见 filter_transformer.pyclass FilterTransformer(Transformer): def __init__(self, config: FilterTransformerConfig): self.config: FilterTransformerConfig config classmethod def create(cls, config_dict: dict, ctx: PipelineContext) - Transformer: config FilterTransformerConfig.model_validate(config_dict) return cls(config) def transform(self, env_event: EventEnvelope) - Optional[EventEnvelope]: # Match Event Type. if not match_util.matches(self.config.event_type, env_event.event_type): return None # Match Event Body. if self.config.event is not None: body_as_json_dict json.loads(env_event.event.as_json()) for key, val in self.config.event.items(): if not match_util.matches(val, body_as_json_dict.get(key)): return None return env_event这段代码展示了两个关键工程实践配置解析使用 Pydantic 模型FilterTransformerConfig字段event_type与可选event在校验配置的同时给出结构化错误提示比手写字典访问更健壮过滤语义transform在事件类型不匹配或事件体字段不匹配时返回None从而让 Pipeline 丢弃该事件。建议在自定义 Transformer 中同样使用 Pydantic 配置模型项目统一的ConfigModel基类来自datahub库并优先通过event.event_type判断事件类型、通过event.event.as_json()读取事件体。三、Step 2让框架发现你的 Transformer定义好类之后还需要让它对 DataHub Actions 框架可见。框架通过 Python 模块路径 类名来定位 Transformer因此只要模块可被 Python 运行时导入即可。3.1 最简单的方式与配置文件同目录直接把custom_transformer.py放在配置文件如custom_transformer_action.yaml所在的目录下。此时模块名与文件名一致即custom_transformer之后在配置中引用custom_transformer:CustomTransformer即可。3.2 进阶打包为 pip 包安装当 Transformer 需要被多个项目复用、或希望用全局唯一名称引用时可将其打包。在与 Transformer 相同的目录下创建setup.pyfrom setuptools import find_packages, setup setup( namecustom_transformer_example, version1.0, packagesfind_packages(), # if you dont already have DataHub Actions installed, add it under install_requires # install_requires[acryl-datahub-actions] )然后在包目录内执行安装开发模式安装便于迭代pip install -e .原文档还提到可用python setup.py即python setup.py install作为备选但现代 Python 环境更推荐pip install -e .。安装完成后类即可通过全限定名被引用custom_transformer_example.custom_transformer:CustomTransformer即包名.模块名:类名的格式。3.3 底层注册机制Type 字符串如何变成实例配置中的type字段之所以能直接写全限定模块路径:类名其机制在 transformer_registry.pytransformer_registry PluginRegistry[Transformer]() transformer_registry.register_from_entrypoint(datahub_actions.transformer.plugins) transformer_registry.register(__filter, FilterTransformer)transformer_registry是datahub库提供的PluginRegistry它会扫描所有已安装包中声明了datahub_actions.transformer.pluginsentry point 的插件并注册其全局别名同时内置注册了__filter这一系统级 Transformer即FilterTransformer对于没有注册别名的类型字符串PluginRegistry.get(type)会按 Python 的导入语义解析模块路径:类名。而在 Pipeline 构建阶段pipeline_util.py 中的create_transformer完成最终实例化def create_transformer(transform_config: TransformConfig, ctx: PipelineContext) - Transformer: transformer_type transform_config.type transformer_class transformer_registry.get(transformer_type) transformer_config transform_config.config if transform_config.config is not None else {} transformer_instance transformer_class.create(transformer_config, ctx) return transformer_instance即从配置中取出type字符串 → 通过 registry 解析出 Transformer 类 → 调用其create工厂方法传入 config 与 ctx→ 得到实例。TransformConfig模型定义于 pipeline_config.py仅含type: str与可选的config: Dict两个字段这也解释了为什么 YAML 中每个 transform 条目只需要这两个键。四、Step 3配置并运行包含自定义 Transformer 的 Action4.1 编写 Actions 配置文件创建custom_transformer_action.yaml在transform段引用刚开发好的 Transformer。type必须使用全限定 Python 模块与类名# custom_transformer_action.yaml name: custom_transformer_test source: type: kafka config: connection: bootstrap: ${KAFKA_BOOTSTRAP_SERVER:-localhost:9092} schema_registry_url: ${SCHEMA_REGISTRY_URL:-http://localhost:8081} transform: - type: custom_transformer_example.custom_transformer:CustomTransformer config: # Some sample configuration which should be printed on create. config1: value1 action: # Simply reuse the default hello_world action type: hello_world配置文件要点说明source事件源。示例使用 Kafka 事件源bootstrap与schema_registry_url支持${ENV_VAR:-default}形式的环境变量占位符localhost:9092与http://localhost:8081是缺省值transform列表类型支持配置多个 Transformer按顺序组成转换链Pipeline 会依次执行每个条目由type全限定名或已注册别名与config自由形式字典原样传给create组成。上面config1: value1即会在create中被打印出来action动作类型。示例直接复用内置的hello_world动作方便快速验证namePipeline 名称用于日志、统计与失败事件落盘目录命名。关于配置模型pipeline_config.py 中的PipelineConfig还支持enabled、filters、datahubDataHub 客户端连接配置以及options字段options可配置retry_count事件处理失败重试次数默认 0、failure_modeTHROW停止 Pipeline /CONTINUE跳过继续默认CONTINUE与failed_events_dir失败事件落盘目录默认/tmp/logs/datahub/actions这些能力同样适用于包含自定义 Transformer 的 Pipeline。4.2 启动 Action使用datahub actions命令以配置文件为参数启动datahub actions -c custom_transformer_action.yamlCLI 入口位于 entrypoints.py支持--debug输出 DEBUG 日志也可通过环境变量DATAHUB_DEBUG开启、--enable-monitoring启动 Prometheus 指标端点默认端口 8000等选项。若一切正常你的 Transformer 会开始接收并打印事件——create时打印传入的配置字典收到事件后打印EventEnvelope并原样放行给hello_world动作。4.3 观察运行效果在标准输出中你应当能看到类似序列Pipeline 启动时create(config_dict, ctx)被调用打印{config1: value1}每当 Kafka 事件源消费到一条元数据变更事件transform(event)被调用并打印事件对象事件随后被交给hello_world动作处理。若 Transformer 返回NonePipeline 会按 pipeline.py 的逻辑短路转换链并跳过 action同时累加统计计数increment_transformer_filtered_count。每个 Transformer 的处理数、过滤数、异常数等指标均由PipelineStats统计可通过监控端点观察。五、Step 4可选将 Transformer 贡献回 DataHub 核心库如果你的 Transformer 具有通用价值可以通过提交 PR 将其纳入 DataHub 提供的核心 Transformer 库。全部核心 Transformer 位于当前仓库datahub-actions/src/datahub_actions/plugin/transform目录下参考filter子目录中的 filter_transformer.py与setup.py中声明的datahub_actions.transformer.pluginsentry point 相对应。贡献时需要遵循两条硬性前提Testing测试为你的 Transformer 编写单元测试。可以在datahub-actions/tests/unit目录参照现有测试风格构造EventEnvelope与配置字典验证create的实例化行为以及transform在匹配/不匹配/改写等场景下的返回值含返回None的过滤分支Deduplication去重确认现有 Transformer 中没有功能等价、或稍加扩展即可实现同等功能的实现避免重复造轮子。当新 Transformer 被合入核心库后还需要在setup.py的entry_points段为其登记一个全局唯一的短名称例如my_transformer datahub_actions.plugin.transform.my_module:MyTransformer这样使用者就无需写完整模块路径可直接以该短名称在配置中引用。六、常见问题与调试建议报错Failed to create transformer with type xxx检查type字符串是否为包名.模块名:类名全限定格式、模块是否位于 Python 可导入路径当前目录或已安装包、类名是否拼写正确。该错误来自 pipeline_util.py 中create_transformer对create返回None的校验transform抛异常导致 Pipeline 中断Pipeline 会先重试retry_count次重试耗尽后依据failure_mode决定继续还是终止并将失败事件写入failed_events_dir下的failed_events.log路径规则见 pipeline.py可据此排查问题事件事件没有被处理确认事件源已消费到事件Kafka bootstrap 与 schema registry 地址可达并检查 Pipeline 日志确认事件是否在更早的 filter 阶段或前序 Transformer 处被过滤filtered计数会增加配置不生效确认transform是 YAML 列表以-开头config键名与create中读取的字典键完全一致。结语从本文可以看出开发一个 DataHub Actions Transformer 的完整路径是继承 Transformer 基类 并实现create/transform→ 通过目录放置或 pip 打包让框架可见 → 在 Actions 配置文件的transform段以全限定名引用 → 用datahub actions -c启动验证 → 若具备通用性再以带测试、无重复的标准贡献回核心库。理解 TransformerRegistry 的类型解析与 Pipeline 的链式执行语义将帮助你在多 Transformer 组合、事件过滤、异常处理等生产场景中写出更可靠的转换逻辑。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表