news 2026/9/15 21:17:50

DataHub Entity Change Event V1 完整参考:事件结构、全量事件目录与 Actions 消费实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
DataHub Entity Change Event V1 完整参考:事件结构、全量事件目录与 Actions 消费实战

DataHub Entity Change Event V1 完整参考:事件结构、全量事件目录与 Actions 消费实战

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

Entity Change Event(实体变更事件)是 DataHub 在数据集、图表、看板等实体上发生元数据变更(添加标签、移除归属、修改弃用状态等)时统一发出的EntityChangeEvent_v1事件流。本文以 DataHub 官方参考文档为骨架,结合仓库中的 PDL 模型、Java 生成器与 Python Actions 框架源码,完整覆盖事件通用字段、全部 15 类事件与 7 类 Action Request 提案事件的示例载荷,并给出基于datahub-actions消费与过滤该事件的实战配置。读完本文,你将能够理解并解析任何一条 Entity Change Event,并能编写自己的事件驱动自动化动作。

Event Type 与触发场景

Entity Change Event 的事件类型标识为:

EntityChangeEvent_v1

datahub-actions的事件注册表中,该类型被正式注册为ENTITY_CHANGE_EVENT_V1_TYPE = "EntityChangeEvent_v1"(见 event_registry.py),并与MetadataChangeLogEvent_v1(元数据变更日志)、RelationshipChangeEvent_v1(关系变更)共同组成 DataHub 面向 Actions 的三大事件流。

每当 DataHub 上的实体(dataset、dashboard、chart 等)发生特定变更时,系统便会发出该事件。触发来源包括:用户在 UI 上添加/移除标签、术语、归属、域,修改弃用状态,增删数据集 Schema 字段,以及实体的创建、软删除与硬删除等生命周期操作。

事件通用结构(Common Fields)

不同场景下生成的事件虽然载荷各异,但共享同一组公共字段。事件模型在仓库中由 PDL 定义,见 EntityChangeEvent.pdl,注释中明确说明这些字段与实体注册表(entity registry)一一对应。

字段名类型说明必填
entityUrnString被变更实体的唯一标识,例如 Dataset 的 URNFalse
entityTypeString被变更实体的类型,支持 dataset、chart、dashboard、dataFlow(Pipeline)、dataJob(Task)、domain、tag、glossaryTerm、corpGroup、corpUser 等False
categoryString变更类别,与所执行的操作类型相关,例如 TAG、GLOSSARY_TERM、DOMAIN、LIFECYCLE 等False
operationString在给定类别下对实体执行的操作,例如 ADD、REMOVE、MODIFY,合法操作集合见下文事件目录False
modifierString应用到实体的修饰符,取值依赖 category,例如应用到 Dataset 或 Schema Field 上的 Tag URNTrue
parametersDict附加键值参数,用于提供具体上下文,具体内容取决于事件的 category + operation 组合True
auditStamp.actorString触发该变更的操作者 URNFalse
auditStamp.timeNumber事件对应的时间戳(毫秒)False

PDL 模型中还定义了version字段(整型,表示事件类型版本),不过普通场景下该字段在序列化输出时多数情况为默认值,Action Request 类事件的示例中会显式出现"version": 0

关于parameters需要特别说明:它是 PDL 中的AnyRecord(任意 JSON),因此事件模型本身不做结构校验。datahub-actions在反序列化时将其从对象中弹出并直接注入底层字典(见 event_registry.py 的注释),这意味着消费端必须自行理解parameters内部的序列化 JSON 格式(即 PDL 序列化后的 JSON)。这也是下方事件目录中同一 category 下 parameters 内容各不相同的原因。

事件目录:全量事件与示例载荷

以下逐一给出每种场景下 Entity Change Event 的触发说明与完整示例。所有示例的时间戳与 URN 均为演示值,实际操作中auditStamp.actor会指向真实触发人,entityUrn会指向真实实体。

Tag 相关:Add / Remove Tag Event

当 Tag 被添加到实体上(Add Tag Event)或从实体上移除(Remove Tag Event)时发出。modifierparameters.tagUrn均为被操作 Tag 的 URN,categoryTAG

{ "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 } }

移除标签时仅operation变为REMOVE,其余字段保持一致:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "TAG", "operation": "REMOVE", "modifier": "urn:li:tag:PII", "parameters": { "tagUrn": "urn:li:tag:PII" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

该事件的 Java 侧生成逻辑位于 GlobalTagsChangeEventGenerator.java,属于EntityChangeEventGenerator<GlobalTags>的实现,负责在GlobalTags方面(aspect)变更时产出上述事件。

Glossary Term 相关:Add / Remove Glossary Term Event

当业务术语(Glossary Term)被添加到实体或从实体移除时发出,categoryGLOSSARY_TERMmodifierparameters.termUrn为术语 URN。

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "GLOSSARY_TERM", "operation": "ADD", "modifier": "urn:li:glossaryTerm:ExampleNode.ExampleTerm", "parameters": { "termUrn": "urn:li:glossaryTerm:ExampleNode.ExampleTerm" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

移除时operation变为REMOVE

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "GLOSSARY_TERM", "operation": "REMOVE", "modifier": "urn:li:glossaryTerm:ExampleNode.ExampleTerm", "parameters": { "termUrn": "urn:li:glossaryTerm:ExampleNode.ExampleTerm" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

对应生成器为 GlossaryTermsChangeEventGenerator.java。

Domain 相关:Add / Remove Domain Event

当域(Domain)被添加或移除到实体时发出,categoryDOMAIN

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "DOMAIN", "operation": "ADD", "modifier": "urn:li:domain:ExampleDomain", "parameters": { "domainUrn": "urn:li:domain:ExampleDomain" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

移除域时operation变为REMOVE

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "DOMAIN", "operation": "REMOVE", "modifier": "urn:li:domain:ExampleDomain", "parameters": { "domainUrn": "urn:li:domain:ExampleDomain" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

Owner 相关:Add / Remove Owner Event

当新的所有者被指派到实体(Add Owner Event)或既有所有者被移除(Remove Owner Event)时发出,categoryOWNER。注意此处的parameters额外包含ownerType字段(如BUSINESS_OWNER),用于标明所有者类型。

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "OWNER", "operation": "ADD", "modifier": "urn:li:corpuser:jdoe", "parameters": { "ownerUrn": "urn:li:corpuser:jdoe", "ownerType": "BUSINESS_OWNER" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

移除所有者时operation变为REMOVE

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "OWNER", "operation": "REMOVE", "modifier": "urn:li:corpuser:jdoe", "parameters": { "ownerUrn": "urn:li:corpuser:jdoe", "ownerType": "BUSINESS_OWNER" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

对应生成器为 OwnershipChangeEventGenerator.java。

Structured Property 相关:Add / Remove / Modify Event

Structured Property(结构化属性)被添加到实体、从实体移除或其取值被修改时发出,categorySTRUCTURED_PROPERTY。三种操作的modifier均为属性 URN,而parameters.propertyValues是一个序列化为字符串的 JSON 数组(如"[\"value1\"]"),消费端需要先反序列化再使用。示例载荷中还会显式携带"version": 0

添加:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "STRUCTURED_PROPERTY", "operation": "ADD", "modifier": "urn:li:structuredProperty:prop1", "parameters": { "propertyUrn": "urn:li:structuredProperty:prop1", "propertyValues": "[\"value1\"]" }, "version": 0, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

移除:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "STRUCTURED_PROPERTY", "operation": "REMOVE", "modifier": "urn:li:structuredProperty:prop1", "version": 0, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

修改取值:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "STRUCTURED_PROPERTY", "operation": "MODIFY", "modifier": "urn:li:structuredProperty:prop1", "parameters": { "propertyUrn": "urn:li:structuredProperty:prop1", "propertyValues": "[\"value1\",\"value2\"]" }, "version": 0, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

Deprecation 相关:Modify Deprecation Event

当实体的弃用状态被修改时发出,categoryDEPRECATIONoperationMODIFYmodifierparameters.status取值为DEPRECATEDACTIVE(即弃用/恢复可用)。

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "DEPRECATION", "operation": "MODIFY", "modifier": "DEPRECATED", "parameters": { "status": "DEPRECATED" }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

对应生成器为 DeprecationChangeEventGenerator.java,它监听Deprecation方面(aspect)的变更并生成事件。

Dataset Schema Field 相关:Add / Remove Event

当数据集 Schema 新增字段(Add Dataset Schema Field Event)或移除字段(Remove Dataset Schema Field Event)时发出,categoryTECHNICAL_SCHEMA。注意:Schema Field 的 URN 使用带括号的复合格式urn:li:schemaField:(<datasetUrn>,<fieldName>),且事件仍以数据集实体作为entityUrn/entityType,字段本身通过modifierparameters.fieldUrn标识。

添加字段:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "TECHNICAL_SCHEMA", "operation": "ADD", "modifier": "urn:li:schemaField:(urn:li:dataset:abc,newFieldName)", "parameters": { "fieldUrn": "urn:li:schemaField:(urn:li:dataset:abc,newFieldName)", "fieldPath": "newFieldName", "nullable": false }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

移除字段:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "TECHNICAL_SCHEMA", "operation": "REMOVE", "modifier": "urn:li:schemaField:(urn:li:dataset:abc,newFieldName)", "parameters": { "fieldUrn": "urn:li:schemaField:(urn:li:dataset:abc,newFieldName)", "fieldPath": "newFieldName", "nullable": false }, "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

生命周期相关:Create / Soft-Delete / Hard-Delete Event

当实体被创建(Entity Create Event)、软删除(Entity Soft-Delete Event)或硬删除(Entity Hard-Delete Event)时发出,category均为LIFECYCLE,通过operation区分:CREATESOFT_DELETEHARD_DELETE。此类事件通常不带modifierparameters

创建:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

软删除:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "LIFECYCLE", "operation": "SOFT_DELETE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

硬删除:

{ "entityUrn": "urn:li:dataset:abc", "entityType": "dataset", "category": "LIFECYCLE", "operation": "HARD_DELETE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1649953100653 } }

Action Request Events(提案事件)

Action Request 事件代表对实体的变更提案(Proposals)——这类变更可能需要经过审批(approval)才会真正生效。其共同特征为:

  • entityType固定为actionRequest
  • entityUrn为提案本身的 URN(如urn:li:actionRequest:abc-123);
  • 使用LIFECYCLEcategory,配合CREATE(新建提案)或MODIFY(审批步骤推进)、COMPLETED(提案完结)等操作;
  • 具体提案内容全部承载在parameters中,其中actionRequestType标识提案类型,resourceUrn/resourceType标识被提案的目标资源。

Domain Association Request Event

当为实体提出域关联(domain association)提案时发出。parameters中的domains是序列化 JSON 字符串数组:

{ "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:abc-123", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1234567890 }, "version": 0, "parameters": { "domains": "[\"urn:li:domain:marketing\"]", "actionRequestType": "DOMAIN_ASSOCIATION", "resourceUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "resourceType": "dataset" } }

Owner Association Request Event

当为实体提出所有者关联(owner association)提案时发出。parameters.owners为序列化 JSON 数组,每个元素含typetypeUrnownerUrn

{ "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:def-456", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1234567890 }, "version": 0, "parameters": { "owners": "[{\"type\":\"TECHNICAL_OWNER\",\"typeUrn\":\"urn:li:ownershipType:technical_owner\",\"ownerUrn\":\"urn:li:corpuser:jdoe\"}]", "actionRequestType": "OWNER_ASSOCIATION", "resourceUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "resourceType": "dataset" } }

Tag Association Request Event

当为实体提出标签关联(tag association)提案时发出:

{ "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:ghi-789", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1234567890 }, "version": 0, "parameters": { "actionRequestType": "TAG_ASSOCIATION", "resourceUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "tagUrn": "urn:li:tag:pii", "resourceType": "dataset" } }

Create Glossary Term Request Event

当提出新建业务术语(glossary term)提案时发出。注意此提案针对的是术语本身而非既有实体,因此parameters中包含parentNodeUrn(父节点)与glossaryEntityName(术语名),resourceTypeglossaryTerm

{ "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:jkl-101", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1234567890 }, "version": 0, "parameters": { "parentNodeUrn": "urn:li:glossaryNode:123", "glossaryEntityName": "ExampleTerm", "actionRequestType": "CREATE_GLOSSARY_TERM", "resourceType": "glossaryTerm" } }

Term Association Request Event

当为实体提出术语关联(glossary term association)提案时发出:

{ "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:mno-102", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1234567890 }, "version": 0, "parameters": { "glossaryTermUrn": "urn:li:glossaryTerm:123", "actionRequestType": "TERM_ASSOCIATION", "resourceUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "resourceType": "dataset" } }

Update Description Request Event

当提出更新实体描述(description)提案时发出,parameters.description为拟写入的描述文本:

{ "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:pqr-103", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1234567890 }, "version": 0, "parameters": { "description": "Example description for a dataset.", "actionRequestType": "UPDATE_DESCRIPTION", "resourceUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "resourceType": "dataset" } }

Structured Property Association Request Event

当为实体提出结构化属性关联提案时发出,parameters.structuredProperties为序列化 JSON 数组,元素含propertyUrnvalues

{ "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:stu-104", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "actor": "urn:li:corpuser:jdoe", "time": 1234567890 }, "version": 0, "parameters": { "structuredProperties": "[{\"propertyUrn\":\"urn:li:structuredProperty:123\",\"values\":[\"value1\",\"value2\"]}]", "actionRequestType": "STRUCTURED_PROPERTY_ASSOCIATION", "resourceUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "resourceType": "dataset" } }

Data Access Workflows 审批流事件

以下三个事件服务于 DataHub 的 Data Access Workflows(数据访问工作流),全部围绕审批表单请求(actionRequestType: WORKFLOW_FORM_REQUEST)产生。其外层结构为event_type+event包裹格式(即事件信封 Envelope),事件体位于event键内。

特别提醒parameters.fields包含用户在提交审批请求表单时填写的所有表单字段,是一个序列化后的 JSON 对象字符串,访问前必须先反序列化!

Approval Workflow Request Create Event

当一个新的审批工作流表单请求被提交时发出,operationCREATE

{ "event_type": "EntityChangeEvent_v1", "event": { "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:88170f40-b4f5-4b05-a997-f5f680248b59", "category": "LIFECYCLE", "operation": "CREATE", "auditStamp": { "time": 1753829419490, "actor": "urn:li:corpuser:admin" }, "version": 0, "parameters": { "actorUrn": "urn:li:corpuser:admin", "actionRequestType": "WORKFLOW_FORM_REQUEST", "workflowUrn": "urn:li:actionWorkflow:7a6e7a24-7525-4f1e-8d9b-08ba562399d0", "expiresAtMs": 1756680611743, "fields": "{\"reason\":[\"Building Tableau Dashboard Q2\"],\"privileges\":[\"SELECT\"],\"projects\":[\"Orders Dashboard Q2\"],\"tags\":[\"urn:li:tag:__default_gold\"]}", "workflowId": "7a6e7a24-7525-4f1e-8d9b-08ba562399d0", "entityType": "dataset", "entityUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "entityName": "table", "qualifiedName": "example.table", "entityPlatformName": "Snowflake" } } }
Approval Workflow Request Step Complete Event

当工作流请求中的某个审批步骤完成(被批准、拒绝或要求补充信息)时发出,operationMODIFYparameters.stepResult标明该步骤的结果(如ACCEPTED),parameters.stepId标明步骤标识:

{ "event_type": "EntityChangeEvent_v1", "event": { "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:ce424fff-7fab-4731-b106-fab9d69238da", "category": "LIFECYCLE", "operation": "MODIFY", "auditStamp": { "time": 1753829625133, "actor": "urn:li:corpuser:admin" }, "version": 0, "parameters": { "actorUrn": "urn:li:corpuser:admin", "actionRequestType": "WORKFLOW_FORM_REQUEST", "stepId": "approval-step-2", "workflowUrn": "urn:li:actionWorkflow:7a6e7a24-7525-4f1e-8d9b-08ba562399d0", "expiresAtMs": 1754002407730, "stepResult": "ACCEPTED", "fields": "{\"reason\":[\"Test\"],\"privileges\":[\"SELECT\"],\"tags\":[\"urn:li:tag:NeedsDocumentation\"]}", "workflowId": "7a6e7a24-7525-4f1e-8d9b-08ba562399d0", "entityType": "dataset", "entityUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "entityName": "table", "qualifiedName": "example.table", "entityPlatformName": "Snowflake" } } }
Approval Workflow Request Complete Event

当审批工作流请求全部完成(所有步骤通过,或在任一步骤被拒绝)时发出,operationCOMPLETEDparameters.result为最终结果(如ACCEPTED):

{ "event_type": "EntityChangeEvent_v1", "event": { "entityType": "actionRequest", "entityUrn": "urn:li:actionRequest:553595ec-2295-4970-89a4-bb7d02e691c0", "category": "LIFECYCLE", "operation": "COMPLETED", "auditStamp": { "time": 1753829954539, "actor": "urn:li:corpuser:admin" }, "version": 0, "parameters": { "result": "ACCEPTED", "actorUrn": "urn:li:corpuser:admin", "actionRequestType": "WORKFLOW_FORM_REQUEST", "workflowUrn": "urn:li:actionWorkflow:7a6e7a24-7525-4f1e-8d9b-08ba562399d0", "expiresAtMs": 1756335474418, "fields": "{\"reason\":[\"test\"],\"privileges\":[\"SELECT\"],\"projects\":[\"test\"],\"tags\":[\"urn:li:tag:NeedsDocumentation\"]}", "operation": "COMPLETE", "workflowId": "7a6e7a24-7525-4f1e-8d9b-08ba562399d0", "entityType": "dataset", "entityUrn": "urn:li:dataset:(urn:li:dataPlatform:snowflake,example.table,PROD)", "entityName": "table", "qualifiedName": "example.table", "entityPlatformName": "Snowflake" } } }

源码视角:事件如何被生成与消费

后端生成:Java 侧 Event Generator

事件并非凭空产生,而是由元数据时间线(timeline)体系中的各类EntityChangeEventGenerator在对应方面(aspect)发生变更时计算生成。仓库 metadata-io/src/main/java/com/linkedin/metadata/timeline/eventgenerator 目录下可以看到与上文事件目录一一对应的实现:

  • GlobalTagsChangeEventGenerator.java——TAG 事件的 Add/Remove;
  • GlossaryTermsChangeEventGenerator.java——GLOSSARY_TERM 事件;
  • OwnershipChangeEventGenerator.java——OWNER 事件;
  • DeprecationChangeEventGenerator.java——DEPRECATION 事件;
  • 此外还有 DocumentationChangeEventGenerator.java、ApplicationsChangeEventGenerator.java 等,覆盖更多方面类型。

这些生成器通过EntityChangeEventGeneratorRegistryFactory(metadata-service/factories)注册,事件模型本身由 EntityChangeEvent.pdl 定义并编译为EntityChangeEventClass

前端消费:datahub-actions 的事件解析

datahub-actions中,EntityChangeEvent类继承自编译生成的EntityChangeEventClass,其from_json反序列化时会将parameters从对象中弹出并直接注入内部字典(event_registry.py)。为此,框架额外暴露了safe_parameters属性,统一返回parameters或内部注入的原始 JSON,供 Action 安全读取。

事件过滤:EventTypeFilter 的匹配语义

对 Entity Change Event 做路由过滤时,使用 event_type_filter.py 中定义的EventTypeFilter。其匹配语义如下:

  • 事件类型层面:OR——事件只要命中filter配置中的任意一个类型键即进入下一步判断;
  • 单个类型下的 body 谓词列表:OR——只要满足任意一条谓词字典即可通过;
  • 单条谓词字典内部的键值对:AND——所有键值必须同时匹配。

例如,同时过滤MetadataChangeLogEvent_v1EntityChangeEvent_v1,并限定后者的categoryDOCUMENTATIONentityTypeschemaField

filter: event_type: MetadataChangeLogEvent_v1: event: - entityType: schemaField aspectName: documentation - entityType: dataset aspectName: documentation EntityChangeEvent_v1: event: - category: DOCUMENTATION entityType: schemaField

实战:基于 Entity Change Event 编写 Actions

最小 Pipeline 配置

一个消费 Entity Change Event 的 Action pipeline 由四段构成:source(事件源)、filter(事件过滤)、action(处理动作)、datahub(回写端点)。仓库 examples 下的配置均遵循此结构,例如 snowflake_tag_propagation.yaml 展示了一个完整的标签/术语传播场景:

name: "snowflake_tag_propagation" 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: "snowflake_tag_propagation" config: tag_propagation: tag_prefixes: - classification term_propagation: target_terms: - Classification term_groups: - "Personal Information" snowflake: account_id: ${SNOWFLAKE_ACCOUNT_ID} warehouse: COMPUTE_WH username: ${SNOWFLAKE_USER_NAME} password: ${SNOWFLAKE_PASSWORD} role: ACCOUNTADMIN datahub: server: "http://localhost:8080"

其中filter.event_type直接指定EntityChangeEvent_v1,表示该 pipeline 只处理实体变更事件流。

Action 内部如何识别事件

以标签传播 Action 为例,tag_propagation_action.py 中的should_propagate方法展示了典型的事件判别逻辑:

  • 先判断event.event_type == "EntityChangeEvent_v1",并对事件对象做类型断言;
  • 再检查semantic_event.category == "TAG"operationADDREMOVE
  • 读取semantic_event.modifier作为被操作的 Tag URN;
  • 若配置了tag_prefixes(如classification),仅当 modifier 命中前缀时才触发传播;
  • 返回TagPropagationDirective(含 propagate、tag、operation、entity 四个字段)供后续执行。

这正是文档中事件目录所定义字段(category/operation/modifier)在真实代码中的直接使用方式,说明掌握了事件目录就能编写任何基于该事件流的 Action。同类实现还包括术语传播 term_propagation_action.py、Slack 通知 slack.py 等。

消费端注意事项

  1. parameters是自由格式:PDL 中为AnyRecord,框架不校验其内部结构,消费端需依据category+operation组合自行解析(event_registry.py)。
  2. 字符串序列化的 JSON 字段propertyValuesdomainsownersfieldsstructuredProperties等均以字符串形式内嵌 JSON,读取后必须先json.loads反序列化再访问。
  3. Schema Field 的 URN 是复合格式urn:li:schemaField:(<datasetUrn>,<fieldName>),解析时注意括号与逗号分隔结构。
  4. Action Request 事件的目标资源在parametersentityUrn指向提案本身(urn:li:actionRequest:...),真正受影响的数据资源通过parameters.resourceUrn/parameters.resourceType表达。
  5. 版本字段version为事件类型版本(整型),一般取0;解析时建议做防御性处理。
  6. 审批流事件带事件信封:Data Access Workflows 相关事件的外层是event_type+event信封格式,事件体位于event键内,与普通事件的扁平结构不同。

总结

Entity Change Event(EntityChangeEvent_v1)是 DataHub 事件驱动体系中最贴近业务语义的一类事件:它用category(TAG、GLOSSARY_TERM、DOMAIN、OWNER、STRUCTURED_PROPERTY、DEPRECATION、TECHNICAL_SCHEMA、LIFECYCLE)与operation(ADD、REMOVE、MODIFY、CREATE、SOFT_DELETE、HARD_DELETE、COMPLETED)的组合,把标签、术语、归属、域、结构化属性、弃用、Schema 与实体生命周期等元数据变更统一抽象为可编程的 JSON 载荷。配合 datahub-actions 的注册表解析与EventTypeFilter过滤语义,你可以基于它构建标签传播、术语传播、Slack 通知、审批流联动等自动化场景;而每个category在 eventgenerator 目录下都有对应的 Java 生成器实现,可进一步追溯事件的产生源头。

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/15 21:17:47

RS485转以太网实战:协议模式、EMS防护与SCADA接入

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/15 21:15:14

基于Hadoop+Spark的胆结石大数据分析系统设计与实现

1. 项目背景与核心价值这个毕业设计选题完美结合了医疗健康与大数据技术两大热门领域。胆结石作为一种常见消化系统疾病&#xff0c;全球发病率约10%-15%&#xff0c;我国部分地区甚至高达20%。传统临床研究受限于样本量小、维度单一等问题&#xff0c;而基于HadoopSpark的分布…

作者头像 李华
网站建设 2026/9/15 21:13:32

极验四代滑块验证码轨迹构造:物理模型与行为特征模拟实战

说实话&#xff0c;极验四代滑块验证码这玩意儿&#xff0c;我在很长一段时间里看见就头疼。前两篇我们聊了怎么定位缺口、怎么拿参数&#xff0c;但那都只是前戏。真正决定你能不能稳定跑通的&#xff0c;就是标题里写的这三个字&#xff1a;轨迹构造。你就算把缺口识别得再准…

作者头像 李华