尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

DataHub Cloud 事件接入 AWS EventBridge 实战指南:Entity Events API 事件结构、Event Bus 配置与规则路由

DataHub Cloud 事件接入 AWS EventBridge 实战指南:Entity Events API 事件结构、Event Bus 配置与规则路由 DataHub Cloud 事件接入 AWS EventBridge 实战指南Entity Events API 事件结构、Event Bus 配置与规则路由【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本指南以 DataHub 仓库中 setting-up-events-api-on-aws-eventbridge.md 为骨架完整讲解如何将 DataHub Cloud 的元数据变更事件Entity Events API通过 AWS EventBridge 实时接入自有云环境。读者将掌握EventBridge 事件的标准封装结构与过滤 Pattern 写法、创建带资源策略的专用 Event Bus、配置路由规则以及如何与 DataHub Actions 联动实现实时消费与处理。背景Entity Events API 与事件驱动架构Entity Events API 允许你将 DataHub 元数据图上发生的变更如给数据集打上 PII 标签、术语关联、字段新增、实体生命周期变更等以事件形式实时集成到更广泛的事件驱动架构中。在 Entity Events API 文档 中官方将其典型应用场景归纳为五类工作流集成Workflow Integration将 DataHub 流程接入组织内部的工单系统例如当某个 Dataset 被提出标签或术语建议时自动创建 Jira 工单。通知Notifications当 DataHub 发生变更时生成组织级通知例如当任何数据资产被添加 PII 标签时给治理团队发送邮件。元数据增强Metadata Enrichment在上游变更发生时触发下游元数据变更例如向血缘下游实体传播术语或标签。同步Synchronization把 DataHub 中的变更同步到第三方系统例如把 DataHub 中的标签变更反映到 Snowflake。审计Auditing审计谁在什么时间对 DataHub 做了什么变更。DataHub Cloud 通过 AWS EventBridge 向外部投递这些事件本指南即围绕在 AWS 侧创建接收链路展开。事件结构AWS EventBridge 标准封装与 Entity Event 载荷与所有 AWS EventBridge 事件一样DataHub Cloud 发出的事件会被一组标准字段包裹最值得关注的三个字段是字段含义source事件来源的唯一标识。DataHub Cloud 默认使用acryl.events。account事件产生的 AWS 账户。该值对应 DataHub Cloud 的 AWS 账户 ID由你的 DataHub Cloud 客户成功CustomerSuccess代表提供。detailEntity Event 载荷实际出现的字段即下面将要展开的事件正文。完整示例事件以下是一个为数据集添加 PII 标签事件被 EventBridge 封装后的完整 JSON{ version: 0, id: 6a7e8feb-b491-4cf7-a9f1-bf3703467718, detail-type: entityChangeEvent, source: acryl.events, account: 111122223333, time: 2017-12-22T18:43:48Z, region: us-west-1, detail: { entityUrn: urn:li:dataset:abc, entityType: dataset, category: TAG, operation: ADD, modifier: urn:li:tag:pii, parameters: { tagUrn: urn:li:tag:pii } } }detail-type 与 PDL 定义的对应关系源码佐证注意外层封装中的detail-type为entityChangeEvent。这一命名并非随意设置它直接对应元数据模型中对事件类型的 PDL 定义。在 EntityChangeEvent.pdl 中可以看到Event { name: entityChangeEvent } record EntityChangeEvent { entityType: string entityUrn: Urn category: string operation: string modifier: optional string parameters: optional Parameters auditStamp: AuditStamp version: int }也就是说EventBridge 外层封装里的detail-type直接来源于 PDL 注解中的事件名称entityChangeEvent而detail字段就是EntityChangeEventrecord 序列化后的 JSON 内容。另外在 DataHub Actions 的事件注册表中该类事件被登记为EntityChangeEvent_v1见 event_registry.py其EntityChangeEvent.from_json实现会从 JSON 中取出parameters并以AnyRecord任意 JSON形式注入事件对象因此parameters的内容可以随 category/operation 组合灵活变化。Entity Event 载荷通用字段一览无论事件类型如何detail中每个实体事件都遵循同一套公共字段结构字段类型说明是否必填entityUrnString被变更实体的唯一标识例如某个 Dataset 的 urn否原文标注 False即必填entityTypeString被变更实体的类型支持dataset、chart、dashboard、dataFlowPipeline、dataJobTask、domain、tag、glossaryTerm、corpGroup、corpUser等必填categoryString变更的类别与执行的操作种类相关例如TAG、GLOSSARY_TERM、DOMAIN、LIFECYCLE等必填operationString在给定 category 下对实体执行的操作例如ADD、REMOVE、MODIFY必填modifierString应用到实体上的修饰符取值取决于 category例如应用到 Dataset 或 Schema Field 的 Tag 的 urn可选parametersDict提供具体上下文的附加键值参数精确内容取决于事件的 category operation 组合可选auditStamp.actorString触发变更的操作者 urn必填auditStamp.timeNumber事件对应的毫秒级时间戳必填例如某 Tag 被添加到某 Dataset的事件会把上述字段填充为{ entityUrn: urn:li:dataset:abc, entityType: dataset, category: TAG, operation: ADD, modifier: urn:li:tag:PII, parameters: { tagUrn: urn:li:tag:PII }, auditStamp: { actor: urn:li:corpuser:jdoe, time: 1649953100653 } }事件类型目录category / operation / 实体类型速查DataHub Cloud 支持的事件类型完整目录如下详见 Entity Events API 文档事件categoryoperation适用实体类型添加 TagTAGADDdataset, dashboard, chart, dataJob, container, dataFlow, schemaField移除 TagTAGREMOVEdataset, dashboard, chart, dataJob, container, dataFlow, schemaField添加术语GLOSSARY_TERMADDdataset, dashboard, chart, dataJob, container, dataFlow, schemaField移除术语GLOSSARY_TERMREMOVEdataset, dashboard, chart, dataJob, container, dataFlow, schemaField添加 DomainDOMAINADDdataset, dashboard, chart, dataJob, container, dataFlow移除 DomainDOMAINREMOVEdataset, dashboard, chart, dataJob, container, dataFlow添加 OwnerOWNERADDdataset, dashboard, chart, dataJob, dataFlow, container, glossaryTerm, domain, tag移除 OwnerOWNERREMOVEdataset, dashboard, chart, dataJob, container, dataFlow, glossaryTerm, domain, tag添加描述DOCUMENTATIONADDdataset, dashboard, chart, dataJob, dataFlow, container, glossaryTerm, domain, tag, schemaField移除描述DOCUMENTATIONREMOVEdataset, dashboard, chart, dataJob, container, dataFlow, glossaryTerm, domain, tag, schemaField修改废弃状态DEPRECATIONMODIFYdataset, dashboard, chart, dataJob, dataFlow, container新增 Schema 字段TECHNICAL_SCHEMAADDdataset移除 Schema 字段TECHNICAL_SCHEMAREMOVEdataset实体创建LIFECYCLECREATEdataset, dashboard, chart, dataJob, dataFlow, glossaryTerm, domain, tag, container实体软删除LIFECYCLESOFT_DELETEdataset, dashboard, chart, dataJob, dataFlow, glossaryTerm, domain, tag, container实体硬删除LIFECYCLEHARD_DELETEdataset, dashboard, chart, dataJob, dataFlow, glossaryTerm, domain, tag, container断言运行完成RUNCOMPLETEDassertion数据进程实例启动RUNSTARTEDdataProcessInstance数据进程实例完成RUNCOMPLETEDdataProcessInstance创建 Action Request元数据提案LIFECYCLECREATEDactionRequestAction Request 状态变更LIFECYCLEPENDING/COMPLETEDactionRequestIncident 变更INCIDENTACTIVE/RESOLVEDincident各事件类型的 parameters 明细不同 category/operation 组合下parameters的内容各不相同关键参数如下TAG 相关tagUrn必填被添加/移除的 Tag urnfieldPath可选仅当实体类型为schemaField时出现parentUrn可选仅当实体类型为schemaField时出现指向字段所属的父 Dataset。GLOSSARY_TERM 相关termUrn必填被添加/移除的术语 urnfieldPath、parentUrn同上仅 schemaField 场景。DOMAIN 相关domainUrn必填被添加/移除的 Domain urn。OWNER 相关ownerUrn必填被添加/移除的 Owner urnownerType必填取值如TECHNICAL_OWNER、BUSINESS_OWNER、DATA_STEWARD、NONE等。DOCUMENTATION 相关description必填被添加/移除的描述文本。DEPRECATION 相关status必填实体新的废弃状态取值为DEPRECATED或ACTIVE。TECHNICAL_SCHEMA 相关fieldUrn必填新增/移除字段的 urnfieldPath必填字段路径可参考 Dataset 字段路径说明nullable必填布尔值。RUNassertion相关runResult必填SUCCESS或FAILURErunId必填平台原生运行标识asserteeUrn必填断言所作用的实体 urn。RUNdataProcessInstance相关runResult仅 COMPLETED取值SUCCESS、FAILURE、SKIPPED、UP_FOR_RETRYattempt可选尝试次数dataFlowUrn、dataJobUrn可选仅当运行与 Data Flow / Data Job 关联时填充parentInstanceUrn可选父 DataProcessInstance 的 urn。Action Request 创建LIFECYCLE/CREATED公共参数为actionRequestType必填取值为TAG_ASSOCIATION、TERM_ASSOCIATION、CREATE_GLOSSARY_NODE、CREATE_GLOSSARY_TERM、UPDATE_DESCRIPTION、resourceType、resourceUrn、subResourceType、subResource提案类型专属参数包括tagUrn标签关联、termUrn术语关联、glossaryEntityName/parentNodeUrn/description创建术语/节点。Action Request 状态变更LIFECYCLE/PENDING|COMPLETEDactionRequestStatus必填actionRequestResult仅当状态为COMPLETED时填充取值为ACCEPTED或REJECTED。Incident 变更entities必填与 Incident 关联的实体列表。以断言运行成功为例其完整事件为{ entityUrn: urn:li:assertion:abc, entityType: assertion, category: RUN, operation: COMPLETED, parameters: { runResult: SUCCESS, runId: 123, asserteeUrn: urn:li:dataset:def }, auditStamp: { actor: urn:li:corpuser:jdoe, time: 1649953100653 } }事件过滤使用 Event Pattern 精确定位目标事件EventBridge 的 Rule 通过Event Pattern决定哪些事件被路由到目标。DataHub Cloud 示例给出了一个非常典型的场景——只路由 PII 标签上的 Add Tag 事件{ source: [acryl.events], detail: { category: [TAG], parameters: { tagUrn: [urn:li:tag:pii] } } }这个 Pattern 的含义是只有当source等于acryl.events、detail.category等于TAG、且detail.parameters.tagUrn等于urn:li:tag:pii时事件才会被匹配。通过组合category、operation、entityType、parameters等任意字段你可以把流量精确收敛到真正关心的变更子集避免下游系统被无关事件淹没。Step 1创建专用 Event Bus 并配置资源策略DataHub Cloud 官方建议为 DataHub Cloud 创建专用的 Event Bus而不是混用既有总线操作步骤如下登录 AWS 控制台进入你将部署 EventBridge 的账户。搜索并进入EventBridge页面。切换到Event Buses标签页。点击Create Event Bus创建事件总线。为新总线命名例如acryl-events。定义资源策略Resource Policy。创建 Event Bus 时必须创建一条允许 DataHub Cloud 的 AWS 账户向该总线发布消息的策略即通过账户 ID 授予该账户PutEvents权限。示例策略如下{ Version: 2012-10-17, Statement: [{ Sid: allow_account_to_put_events, Effect: Allow, Principal: { AWS: arn:aws:iam::795586375822:root }, Action: events:PutEvents, Resource: event-bus-arn }] }其中需要你自行填写的字段是event-bus-arn新 Event Bus 的 AWS ARN。说明Principal中的账户号是文档给出的 DataHub Cloud 发布方账户示例实际生产配置时应以你的 DataHub Cloud 客户成功代表提供的账户 ID 为准。Step 2创建路由规则Routing RuleEvent Bus 定义完成后还需要创建一条规则Rule将入站事件路由到目标Destination例如 SQS 队列、Lambda 函数、Log Group 等。操作步骤如下导航到Rules标签页。点击Create Rule创建规则。为规则命名——通常依据该规则打算路由到的目标来命名。在Event Bus字段中选择Step 1中创建的 Event Bus。选择Rule with Event Pattern带事件模式的规则选项。点击Next下一步。在Event Source事件源中选择Other其他。可选定义 Sample Event示例事件——可以直接使用上文事件结构一节中的示例事件。定义匹配规则Rule Pattern——决定哪些 DataHub 事件会被当前规则路由可参考上文事件结构一节的示例 Pattern。定义Target目标——即匹配规则的事件应被路由到哪里。完成这十步后一个从 DataHub Cloud → Event Bridge → 目标SQS/Lambda/LogGroup 等的完整事件链路就已就绪。Step 3配置 DataHub Cloud 发送事件AWS 侧的接收链路就绪后还需要让 DataHub Cloud 知晓并向该总线投递事件。请将以下信息提供给你的 DataHub Cloud 客户成功CustomerSuccess代表新 Event Bus 的ARNEvent Bus 所在的AWS 区域region。收到并完成配置后DataHub Cloud 即会开始向你的 EventBridge 总线发送事件。消费端联动DataHub Actions 与 DataHub Cloud Event SourceEventBridge 把事件投递到目标之后最常见的消费方式之一是使用DataHub Actions框架进行响应式处理打 Slack/Teams、传播标签术语、写入 Snowflake 等。DataHub Actions 提供官方的 DataHub Cloud Event Source 配置见 datahub_cloud_event_source.py其最小配置如下name: unique-action-name datahub: server: https://your-organization.acryl.io token: your-datahub-cloud-token source: type: datahub-cloud action: # action configs其消费主题默认是PlatformEvent_v1即 Entity Change Event 流对应事件注册表中的EntityChangeEvent_v1见 event_registry.py。若要一并消费原始 Metadata Change LogMCL事件可通过topics参数追加MetadataChangeLog_Versioned_v1与MetadataChangeLog_Timeseries_v1。关于消费语义该事件源实现了 ack 函数仅当事件成功通过 Transformers 并进入 Action无任何错误时才提交消费偏移默认情况下框架提供at-least once至少一次处理语义——如果提交偏移失败重启 Action 时相关事件可能被重放。若 Action 管线的failure_mode为CONTINUE默认处理失败的事件会写入failed_events.log死信日志而不阻塞进度若为THROW失败会终止管线且不提交偏移消息不会被标记为已处理。此外还有lookback_days初次启动回溯天数、reset_offsets重置已存偏移、infinite_retry连接失败无限重试指数退避 2s 至 60s等高级参数可供调优。注意官方明确提示使用相同name部署多个 DataHub Cloud Event Source 属于未定义行为所有事件应由单个运行中的 Action 处理。常见问题与注意事项事件至少投递一次at-least once与多数事件系统一致EventBridge 与 DataHub Actions 消费侧都只保证至少一次投递下游消费方应具备幂等处理能力。account 字段的取值detail外层封装中的account是 DataHub Cloud 事件源账户的 ID由客户成功代表提供不要臆测填写。专用 Event Bus 的必要性官方明确推荐为 DataHub Cloud 创建专用总线便于通过资源策略精确控制 PutEvents 权限、按总线隔离流量。规则匹配粒度Event Pattern 可下钻到parameters层级如tagUrn建议在规则层面尽早过滤减少下游无谓负载。consumer 唯一性若使用 DataHub Actions 消费切勿以同一name并行部署多个消费者否则处理行为未定义详见 DataHub Cloud Event Source 文档。通过本指南的三步配置即可在 AWS 侧建立起完整、可过滤、可扩展的 DataHub Cloud 事件接收链路并借助 DataHub Actions 将元数据变更实时转化为组织内部的工作流、通知与自动化动作。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表