ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

DataHub Actions 快速上手:用实时元数据事件驱动你的数据治理流水线

DataHub Actions 快速上手:用实时元数据事件驱动你的数据治理流水线 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 CLIDataHub Actions 的 CLI 命令是基础datahubCLI 命令的扩展因此官方建议先安装datahubCLIpython3 -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 ActionsActions 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.mdentityUrn是发生变更的实体唯一标识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 V1EntityChangeEvent_v1实体上的元数据变更如打标签、加术语、设 Domain、Schema 字段增删等参考 docs/actions/events/entity-change-event.mdMetadata Change Log V1MetadataChangeLogEvent_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 keyOR——事件只需匹配任意一个列出的类型跨 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 反序列化调用之前就丢弃符合条件的 MetadataChangeLogMCL消息。对于只消费 MCL 流中一小部分数据的流水线这可以显著降低 CPU 开销。为什么只对 MCL 生效MCL 消息的entityType、aspectName、entityUrn、changeType是 Kafka 原始消息顶层的 Avro 字段无需反序列化即可访问而 EntityChangeEventECE消息位于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 onceKafka 事件源实现了ack方法只有当事件成功走完 Transformers 并进入 Action未抛错时框架才会调用ack同步提交 Kafka Consumer Offset。因此框架默认提供至少一次at-least once处理语义——极端情况下若提交 offset 失败事件在重启后可能被重放。失败处理由 Pipeline 的failure_mode决定CONTINUE默认处理失败的事件被写入failed_events.log死信文件事件源继续推进主题 offsetTHROW失败即触发 Pipeline 错误并终止流水线在提交 offset 之前终止消息不会被标记为已处理。事件源配置参数Kafka 事件源的核心配置如下完整参数表见 docs/actions/sources/kafka-event-source.md字段必填默认值说明connection.bootstrap✅N/AKafka bootstrap 地址如localhost:9092connection.schema_registry_url✅N/ASchema Registry 地址如http://localhost:8081connection.consumer_config❌{}透传给底层 Kafka Consumer 的任意键值配置如security.protocol、SSL 证书等topic_routes.mcl❌MetadataChangeLog_v1MetadataChangeLog 事件所在主题topic_routes.pe❌PlatformEvent_v1PlatformEvent 事件所在主题async_commit_enabled❌true是否使用异步后台周期offset 提交批量处理型 Action 可设为falseasync_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.ymlCLI 还支持以下常见场景详见 docs/actions/README.md同时运行多条流水线重复使用-c参数即可每条配置对应一条独立 Pipelinedatahub 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 与 TransformerHello 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),仅供参考
返回列表