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)一一对应。
| 字段名 | 类型 | 说明 | 必填 |
|---|---|---|---|
| entityUrn | String | 被变更实体的唯一标识,例如 Dataset 的 URN | False |
| entityType | String | 被变更实体的类型,支持 dataset、chart、dashboard、dataFlow(Pipeline)、dataJob(Task)、domain、tag、glossaryTerm、corpGroup、corpUser 等 | False |
| category | String | 变更类别,与所执行的操作类型相关,例如 TAG、GLOSSARY_TERM、DOMAIN、LIFECYCLE 等 | False |
| operation | String | 在给定类别下对实体执行的操作,例如 ADD、REMOVE、MODIFY,合法操作集合见下文事件目录 | False |
| modifier | String | 应用到实体的修饰符,取值依赖 category,例如应用到 Dataset 或 Schema Field 上的 Tag URN | True |
| parameters | Dict | 附加键值参数,用于提供具体上下文,具体内容取决于事件的 category + operation 组合 | True |
| auditStamp.actor | String | 触发该变更的操作者 URN | False |
| auditStamp.time | Number | 事件对应的时间戳(毫秒) | 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)时发出。modifier与parameters.tagUrn均为被操作 Tag 的 URN,category为TAG。
{ "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)被添加到实体或从实体移除时发出,category为GLOSSARY_TERM,modifier与parameters.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)被添加或移除到实体时发出,category为DOMAIN。
{ "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)时发出,category为OWNER。注意此处的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(结构化属性)被添加到实体、从实体移除或其取值被修改时发出,category为STRUCTURED_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
当实体的弃用状态被修改时发出,category为DEPRECATION,operation为MODIFY,modifier与parameters.status取值为DEPRECATED或ACTIVE(即弃用/恢复可用)。
{ "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)时发出,category为TECHNICAL_SCHEMA。注意:Schema Field 的 URN 使用带括号的复合格式urn:li:schemaField:(<datasetUrn>,<fieldName>),且事件仍以数据集实体作为entityUrn/entityType,字段本身通过modifier与parameters.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区分:CREATE、SOFT_DELETE、HARD_DELETE。此类事件通常不带modifier与parameters。
创建:
{ "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 数组,每个元素含type、typeUrn、ownerUrn:
{ "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(术语名),resourceType为glossaryTerm:
{ "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 数组,元素含propertyUrn与values:
{ "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
当一个新的审批工作流表单请求被提交时发出,operation为CREATE:
{ "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
当工作流请求中的某个审批步骤完成(被批准、拒绝或要求补充信息)时发出,operation为MODIFY,parameters.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
当审批工作流请求全部完成(所有步骤通过,或在任一步骤被拒绝)时发出,operation为COMPLETED,parameters.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_v1与EntityChangeEvent_v1,并限定后者的category为DOCUMENTATION、entityType为schemaField:
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"且operation为ADD或REMOVE; - 读取
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 等。
消费端注意事项
parameters是自由格式:PDL 中为AnyRecord,框架不校验其内部结构,消费端需依据category+operation组合自行解析(event_registry.py)。- 字符串序列化的 JSON 字段:
propertyValues、domains、owners、fields、structuredProperties等均以字符串形式内嵌 JSON,读取后必须先json.loads反序列化再访问。 - Schema Field 的 URN 是复合格式:
urn:li:schemaField:(<datasetUrn>,<fieldName>),解析时注意括号与逗号分隔结构。 - Action Request 事件的目标资源在
parameters内:entityUrn指向提案本身(urn:li:actionRequest:...),真正受影响的数据资源通过parameters.resourceUrn/parameters.resourceType表达。 - 版本字段:
version为事件类型版本(整型),一般取0;解析时建议做防御性处理。 - 审批流事件带事件信封: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),仅供参考