DataHub Actions 快速上手:用实时元数据事件驱动你的数据治理流水线
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
导读
DataHub Actions 是 DataHub 提供的实时元数据响应框架,它让你能够在 Tag 添加/移除、术语(Glossary Term)增删、Schema 字段变更等事件发生的第一时间做出反应,将 DataHub 无缝接入到基于事件驱动的架构中。本文以docs/actions/quickstart.md为核心,完整讲解从环境准备、插件安装、Hello World 起步,到事件过滤(含多条件 AND/OR 组合与 MCL 预反序列化性能优化)的全过程,并结合仓库源码剖析框架的 Pipeline 架构、Kafka 事件源与内置 Action 的真实实现,读完后你将具备编写、过滤并运行自定义 Actions 流水线的实战能力。
前置条件:先装好 datahub CLI
DataHub Actions 的 CLI 命令是基础datahubCLI 命令的扩展,因此官方建议先安装datahubCLI:
python3 -m pip install --upgrade pip wheel setuptools python3 -m pip install --upgrade acryl-datahub datahub --version版本要求:Actions Framework 要求
acryl-datahub的版本不低于v0.8.34,低于该版本将无法正常使用 Actions 相关命令。相关说明可参见 docs/actions/README.md 与 docs/actions/quickstart.md。
安装 DataHub Actions
Actions Framework 以独立的 PyPI 包acryl-datahub-actions发布,安装后即随附一套开箱即用的 Filters、Transformers、Actions 与 Event Sources 组件库:
python3 -m pip install --upgrade pip wheel setuptools python3 -m pip install --upgrade acryl-datahub-actions # 验证安装是否成功(打印 Actions 框架版本) datahub actions version安装完成后,datahub命令即多出一组actions子命令,可通过datahub actions --help查看可用的选项。
Hello World:第一条实时响应流水线
DataHub 内置了一个 "Hello World" Action,它会把收到的每一个事件以 JSON 格式打印到控制台。这是验证框架连通性最快的方式,也是理解 Action 配置结构的最佳起点。
编写 Action 配置文件
创建一个hello_world.yaml文件:
# hello_world.yaml name: "hello_world" source: type: "kafka" config: connection: bootstrap: ${KAFKA_BOOTSTRAP_SERVER:-localhost:9092} schema_registry_url: ${SCHEMA_REGISTRY_URL:-http://localhost:8081} action: type: "hello_world"配置文件的整体骨架为:name(流水线唯一名称,必须唯一且保持稳定)+source(事件源及其连接配置)+action(要对事件执行的动作),各组件相互独立、可插拔。完整结构(含filters、transform、options、datahub可选段)参见 docs/actions/README.md。
其中${KAFKA_BOOTSTRAP_SERVER:-localhost:9092}是 shell 风格的默认值写法:如果环境变量KAFKA_BOOTSTRAP_SERVER未设置,则回退到localhost:9092。仓库自带的示例配置文件位于 datahub-actions/examples/hello_world.yaml,可直接复制使用。
启动流水线
datahub actions -c hello_world.yaml如果启动成功,控制台会输出:
Action Pipeline with name 'hello_world' is now running.此时登录你已连接的 DataHub 实例,执行任意元数据变更操作,例如:
- 添加 / 移除一个 Tag
- 添加 / 移除一个 Glossary Term
- 添加 / 移除一个 Domain
只要操作成功,控制台就会打印出实时事件。例如,当有人给数据集打上pii标签时,你会看到:
Hello world! Received event: { "event_type": "EntityChangeEvent_v1", "event": { "entityType": "dataset", "entityUrn": "urn:li:dataset:(urn:li:dataPlatform:hdfs,SampleHdfsDataset,PROD)", "category": "TAG", "operation": "ADD", "modifier": "urn:li:tag:pii", "parameters": {}, "auditStamp": { "time": 1651082697703, "actor": "urn:li:corpuser:datahub", "impersonator": null }, "version": 0, "source": null }, "meta": { "kafka": { "topic": "PlatformEvent_v1", "offset": 1262, "partition": 0 } } }上例为“向 Dataset 添加 pii 标签”时发出的事件。
这就是一次完整的 Actions 闭环:元数据变更 → 写入 Kafka 事件流 → Actions 框架消费 → Action 执行。事件中各字段的详细含义可参考 docs/actions/events/entity-change-event.md:entityUrn是发生变更的实体唯一标识,category表示变更类别(如TAG、GLOSSARY_TERM、DOMAIN、LIFECYCLE、TECHNICAL_SCHEMA等),operation表示具体操作(ADD/REMOVE/MODIFY等),auditStamp记录触发变更的操作者与时间戳。
Hello World 的源码实现
从源码看,HelloWorldAction是一个极简的Action基类实现,它通过act()方法接收EventEnvelope并打印其 JSON 表示,还额外支持一个可选的to_upper配置项(默认False,设为true时以大写形式打印事件)。核心实现见 datahub-actions/src/datahub_actions/plugin/action/hello_world/hello_world.py,更完整的参数说明见 docs/actions/actions/hello_world.md。
理解框架核心概念:Pipeline 与事件流
在动手配置更复杂的流水线之前,有必要先建立框架的整体心智模型(详见 docs/actions/concepts.md):
在 Actions Framework 中,事件自左向右持续流动。一条Pipeline是持续运行的进程,它按顺序执行以下职能:
- 从配置的Event Source轮询事件;
- 对事件应用配置的Filters(不匹配的事件被丢弃);
- 对事件应用配置的Transformations(可选,用于事件富化或生成新类型事件);
- 对最终事件执行配置的Action。
此外 Pipeline 还负责初始化、错误处理、重试与日志等横切逻辑。每个 Action 配置文件对应一条唯一的 Pipeline,各 Pipeline 拥有自己独立的事件源、过滤器、转换器与 Action,从而易于为关键任务独立维护状态。
框架围绕“事件”展开:事件对象只需满足“可转成 JSON、可从 JSON 转回”两个约束(以便失败时写入 failed events 文件);每个事件对应一个Event Type(如EntityChangeEvent_v1),可类比为“主题名”。目前框架支持两类事件:
- Entity Change Event V1(
EntityChangeEvent_v1):实体上的元数据变更,如打标签、加术语、设 Domain、Schema 字段增删等,参考 docs/actions/events/entity-change-event.md; - Metadata Change Log V1(
MetadataChangeLogEvent_v1):底层 aspect 变更日志流,参考 docs/actions/events/metadata-change-log-event.md。
事件过滤:只消费你关心的事件
默认情况下,Pipeline 会接收事件源产生的所有事件。多数生产场景下我们只关心特定类型、特定条件的事件,此时就需要配置过滤,让不匹配的事件在进入 Action 之前被静默丢弃(Action 永远不会收到这些事件)。
按事件类型过滤
如果只想消费EntityChangeEvent_v1类型的事件,可以在配置中加入filter段:
# hello_world.yaml name: "hello_world" source: type: "kafka" config: connection: bootstrap: ${KAFKA_BOOTSTRAP_SERVER:-localhost:9092} schema_registry_url: ${SCHEMA_REGISTRY_URL:-http://localhost:8081} filter: event_type: "EntityChangeEvent_v1" action: type: "hello_world"仅过滤EntityChangeEvent_v1类型的事件。
进阶过滤:按事件字段值匹配
除按类型过滤外,还可以通过event块按事件字段的值进行匹配。每个提供的字段都会与真实事件的对应值比较,事件必须全部匹配才会被转发给 Action(同一event块内多个字段为AND语义):
# hello_world.yaml name: "hello_world" source: type: "kafka" config: connection: bootstrap: ${KAFKA_BOOTSTRAP_SERVER:-localhost:9092} schema_registry_url: ${SCHEMA_REGISTRY_URL:-http://localhost:8081} filter: event_type: "EntityChangeEvent_v1" event: category: "TAG" operation: "ADD" modifier: "urn:li:tag:pii" action: type: "hello_world"该过滤条件只匹配“向实体添加 PII 标签”的事件。
单字段 OR 语义:数组取值
如果想对某个字段实现 “OR” 语义,只需为该字段提供一个取值数组。例如同时匹配“添加或移除 PII 标签”:
# hello_world.yaml name: "hello_world" source: type: "kafka" config: connection: bootstrap: ${KAFKA_BOOTSTRAP_SERVER:-localhost:9092} schema_registry_url: ${SCHEMA_REGISTRY_URL:-http://localhost:8081} filter: event_type: "EntityChangeEvent_v1" event: category: "TAG" operation: ["ADD", "REMOVE"] modifier: "urn:li:tag:pii" action: type: "hello_world"该过滤条件只匹配“向实体添加或移除 PII 标签”的事件。
现代过滤语法:filters:列表与event_type过滤器
上述filter:段是早期的过滤写法,目前仍受支持但已废弃(deprecated)。官方推荐使用更富表达力的filters:列表语法,其中内置的event_type过滤器兼具类型路由与字段谓词匹配能力(详见 docs/actions/concepts.md 与 docs/actions/README.md)。
语义规则
event_type过滤器把“事件类型字符串 → 可选的 body 谓词”映射为配置,事件在“类型出现在映射的 key 中且body 满足谓词(若有)”时通过。其完整语义为:
- 跨 event_type key:OR——事件只需匹配任意一个列出的类型;
- 跨 body 谓词列表项:OR——事件 body 只需满足任意一个谓词字典;
- 单个谓词字典内部的多字段:AND——每个键值对都必须匹配。
若省略event谓词,则只按类型匹配即通过。
典型用法示例
同时放行“schemaField 的 documentation aspect 变更(MCL 事件)”与“文档类别的 schemaField 变更(ECE 事件)”:
filters: - type: event_type config: filter: MetadataChangeLogEvent_v1: event: - entityType: schemaField aspectName: documentation EntityChangeEvent_v1: event: - category: DOCUMENTATION entityType: schemaField只放行EntityChangeEvent_v1(配合 Kafka 源的预反序列化过滤,可让所有 MCL 消息在反序列化前被丢弃):
filters: - type: event_type config: filter: EntityChangeEvent_v1: {}多个过滤器(filters列表内的多个 item)之间按顺序以AND语义求值:事件必须满足每一个过滤器才能继续向下游流动。这一判定逻辑在源码中由EventTypeFilter.matches()实现——它先检查事件类型是否出现在配置的 key 中,再对 body 谓词做“任意谓词满足”判定,见 datahub-actions/src/datahub_actions/plugin/filter/event_type_filter.py。
兼容性说明:当
filter:与filters:同时出现时,filter:段会被转换为一个最先执行的 Transformer,而filters:段在 Transformer 之前求值。建议尽早迁移到filters:以获取改进的 Pipeline 级语义与可选的性能收益。
性能优化:MCL 预反序列化过滤
当 Kafka 事件源配置了enable_mcl_pre_deserialization_filter: true时,框架会利用 Pipeline 的event_type过滤条件,在昂贵的MetadataChangeLogClass.from_obj()avrogen 反序列化调用之前就丢弃符合条件的 MetadataChangeLog(MCL)消息。对于只消费 MCL 流中一小部分数据的流水线,这可以显著降低 CPU 开销。
为什么只对 MCL 生效?MCL 消息的entityType、aspectName、entityUrn、changeType是 Kafka 原始消息顶层的 Avro 字段,无需反序列化即可访问;而 EntityChangeEvent(ECE)消息位于PlatformEvent信封内部、载荷为 JSON 编码,category、operation等字段必须等到信封反序列化之后才能读取,因此ECE 的投递永远不会受该开关影响。配置方法:
source: type: kafka config: # ... connection config ... enable_mcl_pre_deserialization_filter: true该优化对“只关心实体变更事件(如 Slack / Teams 通知)”的流水线尤其有效——所有 MCL 消息都可以在不解码的情况下被直接丢弃。
Kafka 事件源:状态、消费者组与投递语义
默认的事件源是Kafka Event Source,它使用 Kafka Consumer 订阅 DataHub 流出的MetadataChangeLog_v1与PlatformEvent_v1主题。每条 Action 会根据配置文件中的唯一name自动归入一个独立的消费者组(Consumer Group),这带来两个关键特性(详见 docs/actions/sources/kafka-event-source.md):
- 水平扩展:在多节点/多进程上用同一份配置文件运行同一
name的 Action,各实例会成为同一消费者组的成员,按分区负载均衡地分担事件流流量; - 状态化(stateful):框架会持续记录各主题的处理 offset。停止并重启 Action 后,它会先从上次停止的位置“补课”处理积压消息。如果 Action 计算开销较大、不想回放历史,最直接的办法是修改配置中的
name(新消费者组默认从日志末尾即 latest 策略开始消费)。
处理保证:at-least once
Kafka 事件源实现了ack方法:只有当事件成功走完 Transformers 并进入 Action(未抛错)时,框架才会调用ack同步提交 Kafka Consumer Offset。因此框架默认提供至少一次(at-least once)处理语义——极端情况下若提交 offset 失败,事件在重启后可能被重放。
失败处理由 Pipeline 的failure_mode决定:
CONTINUE(默认):处理失败的事件被写入failed_events.log死信文件,事件源继续推进主题 offset;THROW:失败即触发 Pipeline 错误并终止流水线(在提交 offset 之前终止),消息不会被标记为已处理。
事件源配置参数
Kafka 事件源的核心配置如下(完整参数表见 docs/actions/sources/kafka-event-source.md):
| 字段 | 必填 | 默认值 | 说明 |
|---|---|---|---|
connection.bootstrap | ✅ | N/A | Kafka bootstrap 地址,如localhost:9092 |
connection.schema_registry_url | ✅ | N/A | Schema Registry 地址,如http://localhost:8081 |
connection.consumer_config | ❌ | {} | 透传给底层 Kafka Consumer 的任意键值配置(如security.protocol、SSL 证书等) |
topic_routes.mcl | ❌ | MetadataChangeLog_v1 | MetadataChangeLog 事件所在主题 |
topic_routes.pe | ❌ | PlatformEvent_v1 | PlatformEvent 事件所在主题 |
async_commit_enabled | ❌ | true | 是否使用异步(后台周期)offset 提交;批量处理型 Action 可设为false |
async_commit_interval | ❌ | 10000 | 后台 offset 提交间隔(毫秒),仅异步模式生效 |
commit_retry_count | ❌ | 5 | 同步提交失败重试次数,仅同步模式生效 |
commit_retry_backoff | ❌ | 10.0 | 同步提交重试间隔(秒),仅同步模式生效 |
异步提交(默认)消除了每个事件约 3ms 的 Kafka 同步往返,高吞吐流水线可因此获得显著吞吐提升;代价是消费者崩溃时最多会重放async_commit_interval毫秒内的事件——对幂等的内置 Action 而言无碍。对于需要更紧 offset 控制的场景可改用同步提交(单 Pod 吞吐约 326 evt/sec,异步约 8,200 evt/sec,数据来自 docs/actions/sources/kafka-event-source.md 提供的基准)。
Schema Registry 配置
Kafka 事件源必须依赖 Schema Registry 反序列化事件。若使用 DataHub 自带的默认 Schema Registry,可用集群内地址:
source: type: "kafka" config: connection: bootstrap: ${KAFKA_BOOTSTRAP_SERVER:-localhost:9092} schema_registry_url: "http://datahub-datahub-gms:8080/schema-registry/api/"若在集群外运行,需将该内部地址映射为外部可访问的 URL。外部 Schema Registry(如 Confluent Cloud)则需提供完整 URL 与认证信息:
source: type: "kafka" config: connection: bootstrap: ${KAFKA_BOOTSTRAP_SERVER:-localhost:9092} schema_registry_url: "https://your-schema-registry-url" schema_registry_config: basic.auth.user.info: "${REGISTRY_API_KEY_ID}:${REGISTRY_API_KEY_SECRET}"Pipeline 进阶选项与 DataHub API 配置
在name、source、filter(s)、transform、action之外,配置文件中还有两段可选配置(完整示例见 docs/actions/README.md):
# 6. 可选:Pipeline 级选项(错误处理等) options: retry_count: 0 # 同一事件处理抛异常时的重试次数,默认 0 failure_mode: "CONTINUE" # 事件处理失败时的行为:'CONTINUE' 继续推进,'THROW' 停止流水线;无论哪种,失败事件都会写入 failed_events.log failed_events_dir: "/tmp/datahub/actions" # 记录失败事件的文件目录,默认 "/tmp/logs/datahub/actions" # 7. 可选:DataHub API 配置(部分 Action 需要) datahub: server: "http://localhost:8080" # DataHub API 地址 # token: <your-access-token> # 若 Metadata Service 开启认证则必填CLI 实战:多流水线、调试与优雅停机
除了基础用法datahub actions -c <config.yml>,CLI 还支持以下常见场景(详见 docs/actions/README.md):
- 同时运行多条流水线:重复使用
-c参数即可,每条配置对应一条独立 Pipeline:
datahub actions -c <config-1.yaml> -c <config-2.yaml>- 调试模式:追加
--debug打印更详细日志。注意调试模式会暴露敏感信息,日志不可外传,也不要在已开启 UI ingestion 的实例上对不可信用户开放:
datahub actions -c <config.yaml> --debug- 优雅停机:直接按下 Control-C,流水线会干净地关闭并打印一段处理结果摘要:
Actions Pipeline with name '<action-pipeline-name>' has been stopped.开发期间还可以使用独立的datahub-actionsCLI 入口(该入口由框架单独暴露),构建并运行示例流水线:
# 构建 datahub-actions 模块 ./gradlew datahub-actions:build # 进入虚拟环境 cd datahub-actions && source venv/bin/activate # 启动 hello world action datahub-actions actions -c ../examples/hello_world.yaml # 同时运行多条流水线 datahub-actions actions -c ../examples/executor.yaml -c ../examples/hello_world.yaml更进一步:自定义 Action 与 Transformer
Hello World 只是起点。框架的 Filter、Transformer、Action 全部是可插拔组件,你可以随时开发自己的实现:
- 自定义 Action:继承
Action基类并实现create()(根据配置字典实例化)、act()(事件处理核心逻辑)、close()(关闭时的清理钩子)三个方法,随后把模块放入配置同目录(或打包成 Python 包),在配置中以全限定名引用,例如type: "custom_action_example.custom_action:CustomAction",再以datahub actions -c custom_action.yaml运行。完整步骤见 docs/actions/guides/developing-an-action.md; - 自定义 Transformer:类似地继承 Transformer 基类,用于在事件进入 Action 前做富化或多步过滤,详见 docs/actions/guides/developing-a-transformer.md;
- 内置 Action 库:除 Hello World 外,框架还预置了 Executor(执行 ingestion 流程)、Slack(发送 Slack 通知)、Teams 等 Action,插件源码位于 datahub-actions/src/datahub_actions/plugin/action,各插件说明见 docs/actions/actions。
总结
从本指南可以看到,DataHub Actions 的快速上手路径非常清晰:安装acryl-datahub-actions→ 编写一份包含name、source、action的 YAML 配置 →datahub actions -c启动。在此之上,filters:现代过滤语法提供了类型路由、字段谓词与 AND/OR 组合的完整表达力,Kafka 事件源借助消费者组实现了状态化处理与水平扩展,failure_mode与async_commit_enabled等选项则让你能按业务需要权衡投递语义与吞吐。建议以 Hello World 打通端到端链路后,逐步为生产环境配置精确的事件过滤与合适的失败策略,再依据本文的源码指引深入扩展自己的 Action 与 Transformer。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考