DataHub Cloud v0.3.3 版本技术解析:按需断言、Slack 事件协作与断言级订阅
2026/9/17 10:32:52 网站建设 项目流程

DataHub Cloud v0.3.3 版本技术解析:按需断言、Slack 事件协作与断言级订阅

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

2024 年 6 月 25 日发布的 DataHub Cloud v0.3.3 是数据质量(Observe)能力的一次集中增强:它首次支持从 UI 中"按需"触发断言运行、为 Slack 中的事件(Incident)协作提供了更丰富的信息载体,并引入断言级订阅以精确控制质量通知的粒度。本文以官方发布说明为主体,结合仓库中的 Assertions API 教程、Incidents 指南、Subscriptions 教程 及 SDK 示例源码,逐项拆解这些新能力的概念、配置方式与底层调用方式,帮助读者在升级后第一时间用上这些功能。

版本概览与升级前提

发布信息

项目内容
版本号v0.3.3
发布可用日期2024-06-25
推荐 CLI/SDKv0.13.3
面向平台DataHub Cloud(SaaS)

强烈建议:升级 CLI/SDK 至 v0.13.3

发布说明中对 CLI/SDK 给出了明确的升级要求:凡是使用 DataHub CLI/SDK 的场景——包括终端命令行、GitHub Actions、Airflow 调度、Python SDK、Java SDK 等——都强烈建议升级到 v0.13.3。官方给出的理由是 CLI 中持续在推送修复,升级到统一版本也有助于平台侧提供更好的支持。

对应到本仓库,CLI 相关入口与用法可参考 docs/cli.md;SDK 侧的断言与订阅调用示例则存放在 metadata-ingestion/examples/library/,下文会逐一引用。

新能力一:按需运行断言(On-Demand Assertions)

v0.3.3 的核心变化之一,是允许用户从 DataHub Cloud UI 直接触发断言的按需运行,而不必等待既定调度周期。该功能默认不开放,需要联系 DataHub 团队在账号侧开启(发布说明原文:Reach out to DataHub team to enable this feature)。

断言运行的两种评估模式

理解"按需运行"前,先明确 DataHub Cloud 中断言是如何被评估的。根据 docs/managed-datahub/observe/assertions.md 的说明,断言有两种评估方式:

  • Active query(主动查询):DataHub Cloud 按调度周期直接对数据源执行 SQL 查询。适用于 Snowflake、Redshift、BigQuery、Databricks。
  • Ingestion-driven(摄取驱动):DataHub Cloud 基于摄取阶段已上报的元数据(Operation操作事件、DatasetProfile数据集画像、已摄取 Schema)进行评估,适用于任何平台,但评估节奏受摄取节奏约束。

按需运行适用于 DataHub Cloud 原生执行的断言;发布说明中同时提醒,若断言是由外部工具(3rd party runner)执行的,则相关 API 会返回错误。

GraphQL 层:三个按需运行 Mutation

Assertions API 教程 中提供了三个按需运行入口,共同约定如下:

  • saveResult/saveResults:是否将运行结果写入 DataHub 后端(UI 可见)。默认true;设为false则结果不落库。
  • async:是否异步运行。默认false(同步执行);异步模式(true)下 API 立即返回,结果需通过assertion(urn: String!)查询的runEvents字段获取。
  • 同步超时限制:同步运行单条断言目前最长30 秒
1. 运行单条断言:runAssertion
mutation runAssertion { runAssertion(urn: "urn:li:assertion:your-assertion-id", saveResult: true) { type nativeResults { key value } } }

其中type返回断言运行结果,取值为SUCCESSFAILUREERROR。成功的响应示例:

{ "data": { "runAssertion": { "type": "SUCCESS", "nativeResults": [ { "key": "Value", "value": "1382" } ] } } }
2. 运行一组断言:runAssertions
mutation runAssertions { runAssertions( urns: [ "urn:li:assertion:your-assertion-id-1" "urn:li:assertion:your-assertion-id-2" ] saveResults: true ) { passingCount failingCount errorCount results { urn result { type nativeResults { key value } } } } }

响应中为每个 URN 返回一个结果对象;外部执行的断言会被静默地从结果集中省略。

3. 运行某数据资产的全部断言:runAssertionsForAsset
mutation runAssertionsForAsset { runAssertionsForAsset( urn: "urn:li:dataset:(urn:li:dataPlatform:snowflake,purchase_events,PROD)" saveResults: true ) { passingCount failingCount errorCount results { urn result { type nativeResults { key value } } } } }

进阶:按标签运行 + 动态参数

  • 按标签筛选运行:先通过addTagmutation 给断言打标签(resourceUrn指向断言 URN),再在runAssertionsForAsset中传入tagUrns参数,即可只运行带特定标签的断言子集,适合按"重要程度"分组执行。
  • 动态参数:在断言的 SQL 片段中写入${parameterName}占位符,运行时可传parameters: [{ key: "parameterName", value: "parameterValue" }]动态注入。这对阈值随业务时段变化的场景(如按一天中不同时段调整限额)非常有用。

Python SDK 侧:仓库中的真实示例

仓库 metadata-ingestion/examples/library/run_assertion.py 展示了最简调用方式:

from datahub.ingestion.graph.client import DatahubClientConfig, DataHubGraph graph = DataHubGraph(config=DatahubClientConfig(server="http://localhost:8080")) assertion_urn = "urn:li:assertion:6e3f9e09-1483-40f9-b9cd-30e5f182694a" # Run the assertion assertion_result = graph.run_assertion(urn=assertion_urn, save_result=True) print(f"Assertion result (SUCCESS / FAILURE / ERROR): {assertion_result.get('type')}")

批量与按资产运行的完整写法分别见 run_assertions.py(graph.run_assertions(urns=..., save_result=True))与 run_assertions_for_asset.py(graph.run_assertions_for_asset(urn=dataset_urn),支持tag_urns按标签过滤)。

典型使用场景:把断言运行嵌入生产数据管道——例如在批处理作业完成后、下游消费前,同步执行数据契约中的关键断言,实现"管道级熔断"(Pipeline Circuit Breaking),避免坏数据继续向下游传播。

新能力二:更丰富的 Slack 事件(Incident)消息

v0.3.3 对 Slack 中的事件协作做了三方面增强(见 docs/managed-datahub/slack/saas-slack-app.md 的 "Manage Data Incidents" 一节):

  1. 直接从 Slack 解决(Resolve)与重新打开(Reopen)事件
  2. 同一消息内实时反映事件最新状态
  3. 由断言触发的事件带有更丰富的细节(关联断言信息、影响面等)。

事件(Incident)的基本模型

在 docs/incidents/incidents.md 中,事件被定义为"标记数据资产处于不健康状态"的独立生命周期概念,包含状态(active / resolved)、标题、描述等字段。两大典型用途:

  • 健康状态沟通:将已知有问题的资产标记为进行中事件,消费方在 UI 中可看到健康徽章并跟踪进展;
  • 管道熔断(进阶):以事件为基础,编排工具可阻塞依赖存在活跃事件输入的数据管道。

用 API 自动化事件的创建与解决

事件同样可以通过 GraphQL 自动化操作。创建事件使用raiseIncidentmutation,返回新事件的 URN:

mutation raiseIncident { raiseIncident(input: { type: OPERATIONAL title: "Dataset Failed Quality Checks" description: "Dataset failed 2/6 Quality Checks for suite run id xy123mksj812pk23." resourceUrn: "urn:li:dataset:(urn:li:dataPlatform:kafka,SampleKafkaDataset,PROD)" }) }
{ "data": { "raiseIncident": "urn:li:incident:bfecab62-dc10-49a6-a305-78ce0cc6e5b1" } }

查询活跃事件则通过实体上的incidents字段,可按state(如ACTIVE)过滤、分页:

query dataset { dataset(urn: "urn:li:dataset:(abc)") { incidents(state: ACTIVE, start: 0, count: 10) { total incidents { urn title description status { state } } } } }

与 v0.3.2 的能力衔接

值得说明的是,事件与断言的联动并非 v0.3.3 首次引入:v0.3.2 发布说明(docs/managed-datahub/release-notes/v_0_3_2.md)中已提到"Field assertions 在配置后会于出错时自动拉起事件""Group 所有者会在 Slack 事件通知中被标注"。v0.3.3 是在此基础上把 Slack 侧的交互与信息密度进一步补强——用户可以在 Slack 消息内直接完成事件的解决与重开,无需跳转回 UI。

新能力三:断言级订阅(Assertion-Level Subscriptions)

v0.3.3 将订阅粒度细化到了"单条断言":你可以订阅某条特定断言通过(passes)、失败(fails)或出错(errors out)时的通知。在此之前,订阅通常以数据集为粒度,粒度细化后可以显著减少通知噪音。

订阅模型:数据集级 vs 断言级

根据 Subscriptions 教程,订阅可以在两个层级创建:

  • 数据集级:数据集上的任何变更(弃用、Schema 变更、所有者变更等)以及该数据集上全部断言的变更都会触发通知;
  • 断言级:仅影响特定断言,适合只关心关键质量门禁的场景。

与断言相关的变更类型(Change Types)包括:

变更类型触发时机
ASSERTION_PASSED断言运行通过(失败或出错之前会被抑制,语义上更接近"恢复通知")
ASSERTION_FAILED断言运行失败
ASSERTION_ERROR断言运行报错

除此之外,订阅还覆盖 Schema 变更(OPERATION_COLUMN_ADDED/REMOVED/MODIFIED)、操作元数据(行插入/更新/删除)、事件状态(INCIDENT_RAISED/INCIDENT_RESOLVED)、弃用(DEPRECATED/UNDEPRECATED)、摄取状态(INGESTION_SUCCEEDED/INGESTION_FAILED)等各类事件类型。

Python SDK 示例:订阅与退订

仓库 metadata-ingestion/examples/library/subscription_create.py 给出了完整的订阅代码,包括数据集级订阅、按指定变更类型订阅单条断言、以及订阅到群组:

from datahub.sdk import DataHubClient client = DataHubClient.from_env() # 数据集级:订阅全部断言变更(entity_change_types 默认覆盖数据集全部可用变更类型) client.subscriptions.subscribe( urn="urn:li:dataset:(urn:li:dataPlatform:snowflake,purchases,PROD)", subscriber_urn="urn:li:corpuser:john.doe", ) # 断言级:只订阅单条断言的通过/失败 client.subscriptions.subscribe( urn="urn:li:assertion:your-assertion-id", subscriber_urn="urn:li:corpuser:john.doe", entity_change_types=["ASSERTION_PASSED", "ASSERTION_FAILED"], ) # 群组级:订阅失败与出错 client.subscriptions.subscribe( urn="urn:li:assertion:your-assertion-id", subscriber_urn="urn:li:corpGroup:data-team", entity_change_types=["ASSERTION_FAILED", "ASSERTION_ERROR"], )

退订使用subscriptions.unsubscribe,支持选择性移除特定变更类型或整体退订,示例见 subscription_remove.py。操作前提是调用方具备数据集上的Manage User Subscriptions权限;若订阅目标为群组,调用方还需是该群组成员。

通知的最终触达

订阅事件经由 DataHub 的通知体系触达用户:根据 docs/managed-datahub/observe/assertions.md 的 Alerts 章节,断言告警可通过Slack 私信(DM)或团队频道事件告警(incident alerts)AWS EventBridge等渠道投递;企业级用户还可通过 DataHub Actions 框架 订阅 Kafka 上的 Assertion Change 事件,构建自定义响应动作。

新能力四:结构化属性(Structured Properties)Schema 变更与删除

v0.3.3 为结构化属性补齐了 Schema 变更(schema change)与删除(delete)能力。结构化属性是附加在数据集、DataJob 等逻辑实体上的带类型、带约束的属性集合,是 DataHub 中实现元数据标准化(如数据保留期、PII 标记)的基础设施,完整用法见 Structured Properties 教程。

该教程覆盖的能力清单与本版本发布说明的表述相互印证:

  • 创建结构化属性(CLIdatahub properties upsert -f {properties_yaml},或 GraphQLcreateStructuredProperty);
  • 列出、读取、删除结构化属性;
  • 为数据集添加结构化属性;
  • 更新结构化属性并支持破坏性 Schema 变更(breaking schema changes)——对应 v0.3.3 的 "Structured Property Schema Change";
  • 按结构化属性做搜索与聚合。

属性定义示例(YAML):

- id: io.acryl.privacy.retentionTime qualified_name: io.acryl.privacy.retentionTime type: number cardinality: MULTIPLE display_name: Retention Time entity_types: - dataset - dataFlow description: "Retention Time is used to figure out how long to retain records in a dataset" allowed_values: - value: 30 description: 30 days, usually reserved for datasets that are ephemeral and contain pii - value: 365 description: Use this for non-sensitive data that can be retained for longer

从源码结构看,OpenAPI 层面向结构化属性暴露了/openapi/v3/entity/structuredProperty端点,相关行为在 EntityControllerTest 中可找到覆盖(如结构化属性 URN、实体上structuredProperties字段的读写断言),可用于验证删除与变更后的读写行为。

使用提示stringrich_textdateurn类型的属性值会被索引为 Elasticsearch / OpenSearch 的 keyword 字段,单个值默认上限为 32,766 UTF-8 字节(可配置),超出时默认被StructuredPropertiesValidator拒绝写入;建议短小精炼的值用于可搜索属性,大段自由文本应存入实体文档。

2.0 UI 修复与其它杂项

v0.3.3 同时包含一批 DataHub 2.0 UI 的体验修复:

  • 新的加载指示器(loading indicators);
  • 修复文本溢出(text overflows)问题;
  • 实体健康徽章样式统一(consistent entity health badges)。

这些修复属于新 UI 的打磨性改动,与上述数据质量功能共同构成该版本的用户可见变化。

上游 OSS DataHub 变更同步

v0.3.3 还同步拉取了自 v0.3.2 以来 OSS DataHub 上游仓库的变更(对应上游 commit 区间为6ed21bd92e9a58)。这部分内容属于 DataHub Cloud 对开源版本的持续跟随,涉及的底层能力(如实体注册、GraphQL resolver、OpenAPI 端点等)都沉淀在本仓库的各 Java 模块中(datahub-graphql-core、metadata-service/openapi-servlet、metadata-io 等),感兴趣可进一步深入源码阅读。

升级与验证清单

综合发布说明与上述各能力文档,升级到 v0.3.3 后建议按以下清单验证:

  1. 升级工具链:将 CLI/SDK 升级至 v0.13.3,覆盖终端、GitHub Actions、Airflow、Python/Java SDK 等所有使用入口。
  2. 启用并验证按需断言:联系 DataHub 团队开启该功能后,用runAssertion/runAssertions/runAssertionsForAsset跑通单条、批量与按资产三个入口,注意同步运行 30 秒超时与saveResult默认落库行为。
  3. 验证 Slack 事件闭环:确认事件消息可在 Slack 内解决与重开、状态实时刷新,并由断言失败触发的告警能携带丰富上下文。
  4. 收敛通知粒度:将高频数据集级订阅收敛为关键断言的断言级订阅,并核对ASSERTION_PASSED(恢复通知语义)、ASSERTION_FAILEDASSERTION_ERROR三类事件是否按预期投递。
  5. 演练结构化属性生命周期:验证属性 Schema 变更(含破坏性变更)与删除操作在 CLI / GraphQL / OpenAPI 三条通道上的行为一致。

至此,v0.3.3 的四大技术主题——按需断言、Slack 事件协作、断言级订阅、结构化属性 Schema 变更与删除——及其 API 与 SDK 用法已全部梳理完毕,可直接对照仓库中的教程与示例落地实践。

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

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询