YAOTU INSIGHTS

DataHub Snowflake DMF Assertions 实战指南:用 YAML 声明式断言编译为 Snowflake Data Metric Functions

DataHub Snowflake DMF Assertions 实战指南:用 YAML 声明式断言编译为 Snowflake Data Metric Functions
DataHub Snowflake DMF Assertions 实战指南用 YAML 声明式断言编译为 Snowflake Data Metric Functions【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本指南围绕 DataHub 的 Open Assertion Compiler 展开讲解如何用简单的 YAML 文件声明数据质量断言并将其编译为可在 Snowflake 上原生执行的 Data Metric FunctionsDMF再通过 Snowflake 摄取把执行结果回传到 DataHub在数据集上下文中以历史时间线的形式呈现。读完本文你将掌握从 YAML 定义、CLI 注册、DMF 编译、Snowflake 侧注册执行到摄取断言结果的完整闭环流程并理解外部用户自建DMF 的接入方式与约束边界。该功能当前处于BETA状态。一、整体思路从 YAML 到 Snowflake DMFDataHub 的 Open Assertion Compiler 允许你以简单的 YAML 格式声明数据质量断言将断言编译成 Snowflake Data Metric FunctionsDMF在 Snowflake 环境中注册编译产物在常规 DataHub 摄取流程中把 DMF 执行结果拉回 DataHub作为普通 Assertion Results 展示在关联表的历史时间线上。编译与执行是分离的compile命令只负责把 YAML 断言翻译成 SQL 工件不连接 Snowflake、不执行任何查询注册与调度由 Snowflake 侧完成结果回传则发生在 DataHub 的 Snowflake ingestion 过程中。从源码看Snowflake 编译器的实现位于 compiler.py其内部由三部分组成SnowflakeMetricSQLGenerator生成指标 SQL、SnowflakeMetricEvalOperatorSQLGenerator生成评估条件 SQL、SnowflakeDMFHandler拼装CREATE DATA METRIC FUNCTION与ALTER TABLE ... ADD DATA METRIC FUNCTION语句。二、前置条件开始之前请确认以下条件全部满足拥有Snowflake Enterprise 账户且已启用 DMF 功能具备在 Snowflake 环境中预置provisionDMF的权限见下文权限章节具备在 Snowflake 环境中查询 DMF 结果的权限见下文权限章节已有一个完成 Snowflake 元数据摄取、可用的 DataHub 实例。如果尚未配置 Snowflake 摄取参考 Snowflake Quickstart Guide 上手已安装 DataHub CLI 并运行过datahub init。三、权限准备在 Snowflake 侧正确配置权限是整套流程成功的前提。以下是三组不同角色所需的最小权限集合。3.1 注册 DMF 所需的权限执行 DMF 注册与取摄取的服务账户service account必须具备PrivilegeObjectNotesUSAGEDatabase, schemaDMF 将被创建在该数据库与 Schema 中由 compile 命令中的DMF_SCHEMA配置指定。CREATE FUNCTIONSchema允许在 compile 命令配置的 Schema 中创建新的 DMF。EXECUTE DATA METRIC FUNCTIONAccount控制哪些角色可以使用与服务器无关server-agnostic的计算资源来调用系统 DMF。USAGEDatabase, schema查询中引用的目标表所在的数据库与 Schema。OWNERSHIPTable允许将 DMF 与目标表关联。USAGEDMF允许调用 compile 命令配置的 Schema 中的 DMF。同时必须授予的角色RoleNotesSNOWFLAKE.DATA_METRIC_USER使用 System DMFs 所需3.2 运行 DMF调度执行所需的权限由于定时调度的 DMF 以表属主table owner的角色运行表属主必须具备PrivilegeObjectNotesUSAGEDatabase, schemaDMF 将被创建在该数据库与 Schema 中由 compile 命令中的DMF_SCHEMA配置指定。USAGEDMF允许调用 compile 命令配置的 Schema 中的 DMF。EXECUTE DATA METRIC FUNCTIONAccount控制哪些角色可以使用与服务器无关的计算资源来调用系统 DMF。同时必须授予的角色RoleNotesSNOWFLAKE.DATA_METRIC_USER使用 System DMFs 所需3.3 查询 DMF 结果所需的权限此外执行 DataHub Ingestion 并查询 DMF 结果的服务账户还必须被授予以下系统应用角色RoleNotesDATA_QUALITY_MONITORING_VIEWER查询 DMF 结果表所需有关 Snowflake DMF 及其预置、查询所需权限的更多细节参见 Snowflake 官方数据质量文档。3.4 授权示例 SQL以下 SQL 给出了一个可直接套用的授权脚本覆盖上述三类角色assertion-service-role、table-owner-role、datahub_role-- 为 assertion-service-role 配置创建 DMF 并与表关联的权限 grant usage on database dmf-database to role assertion-service-role grant usage on schema dmf-database.dmf-schema to role assertion-service-role grant create function on schema dmf-database.dmf-schema to role assertion-service-role -- 授予 assertion-service-role 所有权及其余权限 grant role table-owner-role to role assertion-service-role -- 为 table-owner-role 配置按计划运行 DMF 的权限 grant usage on database dmf-database to role table-owner-role grant usage on schema dmf-database.dmf-schema to role table-owner-role grant usage on all functions in dmf-database.dmf-schema to role table-owner-role grant usage on future functions in dmf-database.dmf-schema to role table-owner-role grant database role SNOWFLAKE.DATA_METRIC_USER to role table-owner-role grant execute data metric function on account to role table-owner-role -- 为 datahub-role 配置查询 DMF 结果的权限 grant application role SNOWFLAKE.DATA_QUALITY_MONITORING_VIEWER to role datahub_role注意其中使用了USAGE ON FUTURE FUNCTIONS确保后续新建的 DMF 也自动对表属主可见。四、支持的断言类型DataHub Snowflake DMF Assertion Compiler 当前支持以下断言类型Freshness校验表是否在指定时间窗口内被更新Volume校验表行数等体量指标是否满足阈值Column校验列上的指标如空值数或列值约束Custom SQL自定义 SQL 返回的指标是否满足条件。注意Schema Assertions 目前不受支持。从源码看编译器对断言类型的支持体现在 metric_sql_generator.py 的singledispatchmethod分发中RowCountChangeVolumeAssertion与SqlMetricChangeAssertion增量/变更型断言会被显式判定为Unsupported assertion type并抛错这与文档中变更类断言不被支持的限制一致。五、完整的五步实战流程整个流程分为五步定义 YAML 断言 → 注册到 DataHub → 编译为 Snowflake DMF → 在 Snowflake 注册执行 → 摄取结果回传。Step 1用 Assertion YAML 文件定义数据质量断言断言以 YAML 声明仓库提供了一个完整的可运行示例 assertions_configuration.yml包含五种典型断言version: 1 namespace: test-config-id-1 assertions: # Freshness Assertion - entity: urn:li:dataset:(urn:li:dataPlatform:snowflake,test_db.public.test_assertions_all_times,PROD) type: freshness lookback_interval: 1 hour last_modified_field: col_timestamp schedule: type: cron cron: 0 * * * * meta: entity_qualified_name: TEST_DB.PUBLIC.TEST_ASSERTIONS_ALL_TIMES entity_schema: - col: col_date native_type: DATE # Volume Assertion - type: volume entity: urn:li:dataset:(urn:li:dataPlatform:snowflake,test_db.public.test_assertions_all_times,PROD) metric: row_count condition: type: less_than_or_equal_to value: 1000 schedule: type: cron cron: 0 * * * * meta: entity_qualified_name: TEST_DB.PUBLIC.TEST_ASSERTIONS_ALL_TIMES entity_schema: - col: col_date native_type: DATE # Field Metric Assertion - type: field entity: urn:li:dataset:(urn:li:dataPlatform:snowflake,test_db.public.test_assertions_all_times,PROD) field: col_date metric: null_count condition: type: equal_to value: 0 schedule: type: cron cron: 0 * * * * meta: entity_qualified_name: TEST_DB.PUBLIC.TEST_ASSERTIONS_ALL_TIMES entity_schema: - col: col_date native_type: DATE # Field Value Assertion - type: field entity: urn:li:dataset:(urn:li:dataPlatform:snowflake,test_db.public.purchase_event,PROD) field: quantity condition: type: between min: 0 max: 10 schedule: type: on_table_change meta: entity_qualified_name: TEST_DB.PUBLIC.PURCHASE_EVENT entity_schema: - col: quantity native_type: FLOAT # Custom SQL Metric Assertion - type: sql entity: urn:li:dataset:(urn:li:dataPlatform:snowflake,test_db.public.purchase_event,PROD) statement: select mode(quantity) from test_db.public.purchase_event condition: type: equal_to value: 5 schedule: type: on_table_change meta: entity_qualified_name: TEST_DB.PUBLIC.PURCHASE_EVENT entity_schema: - col: quantity native_type: FLOAT关键字段说明version配置文件格式版本固定为1namespace内部别名id标识该断言配置文件的唯一命名空间。此结构由 assertion_config_spec.py 中的AssertionsConfigSpec模型校验entity断言作用目标数据集的 URN必须与 DataHub 中已摄取的 Snowflake 数据集一致type断言类型取值为freshness/volume/field/sqlschedule调度方式支持croncron 表达式 时区与on_table_change表变更触发meta.entity_qualified_nameSnowflake 侧的完整表名DB.SCHEMA.TABLE用于生成 DMF 关联语句meta.entity_schemaDMF 参数所需的一个列名与原生类型。Snowflake 不允许创建不带列参数的自定义数据指标函数因此编译器需要从表 schema 中取任意一列作为ARGT TABLE(...)的参数见 compiler.py 中get_dmf_args的实现。Step 2用 DataHub CLI 注册断言将断言注册到 DataHub使其在 DataHub UI 中可见datahub assertions upsert -f examples/library/assertions_configuration.ymlupsert命令读取 YAML 文件为每条断言生成一个urn:li:assertion:id并发射对应的AssertionInfoMCP 到 DataHub。其 CLI 实现在 assertions_cli.py。注仓库中该 assertions CLI 已标记为 deprecated未来版本可能移除并建议关注替代方案编译与注册流程在当前版本仍然可用。Step 3用assertions compile编译为 Snowflake DMF接下来使用compile命令生成可在 Snowflake 注册的 SQL 代码datahub assertions compile -f examples/library/assertions_configuration.yml -p snowflake -x DMF_SCHEMAdb.schema-where-DMF-should-live命令参数参数说明-f, --file断言 YAML 文件路径必填-p, --platform编译目标平台Snowflake 为snowflake必填-o, --output-to编译产物输出目录可选默认当前工作目录下的target/-x, --extras平台相关的额外键值对格式keyvalue可多次传入。Snowflake 平台必须提供DMF_SCHEMAdb.schemaDMF_SCHEMA是必需的扩展参数编译器会严格校验其存在否则报错Must specify value for DMF schema using -x DMF_SCHEMAdb.schema见 compiler.py 的create方法。命令运行后会生成两个文件默认位于target/目录可用-o指定其他目录dmf_definitions.sql将注册到 Snowflake 的 DMF 创建 SQLdmf_associations.sql将 DMF 与目标表关联、并配置调度计划的 SQL。此外还会生成compile_report.json编译报告记录每条断言的编译成功/失败状态见 assertions_cli.py。dmf_definitions.sql内容示例该文件保存从 YAML 断言定义编译生成的 DMF 创建语句-- Example dmf_definitions.sql -- Start of Assertion 5c32eef47bd763fece7d21c7cbf6c659 CREATE or REPLACE DATA METRIC FUNCTION test_db.datahub_dmfs.datahub__5c32eef47bd763fece7d21c7cbf6c659 (ARGT TABLE(col_date DATE)) RETURNS NUMBER COMMENT Created via DataHub for assertion urn:li:assertion:5c32eef47bd763fece7d21c7cbf6c659 of type volume AS $$ select case when metric 1000 then 1 else 0 end from (select count(*) as metric from TEST_DB.PUBLIC.TEST_ASSERTIONS_ALL_TIMES ) $$; -- End of Assertion 5c32eef47bd763fece7d21c7cbf6c659 ....要点解读DMF 名称统一使用datahub__assertion_id前缀源码见 compiler.py 的get_dmf_name这个前缀在摄取阶段用于区分 DataHub 管理的 DMF 与外部 DMFCOMMENT中记录了断言 URN 与断言类型便于溯源函数体是一个返回 0/1 的 CASE 表达式metric 1000返回 1通过否则返回 0失败。所有类型断言的最终评估都会被统一折叠为 0/1 语义组装上述语句的模板见 dmf_generator.py 的create_dmf。dmf_associations.sql内容示例该文件保存将 DMF 与目标表关联、并配置调度计划的语句-- Example dmf_associations.sql -- Start of Assertion 5c32eef47bd763fece7d21c7cbf6c659 ALTER TABLE TEST_DB.PUBLIC.TEST_ASSERTIONS_ALL_TIMES SET DATA_METRIC_SCHEDULE TRIGGER_ON_CHANGES; ALTER TABLE TEST_DB.PUBLIC.TEST_ASSERTIONS_ALL_TIMES ADD DATA METRIC FUNCTION test_db.datahub_dmfs.datahub__5c32eef47bd763fece7d21c7cbf6c659 ON (col_date); -- End of Assertion 5c32eef47bd763fece7d21c7cbf6c659 ....调度方式由断言 YAML 中的schedule决定编译器会将其翻译为 Snowflake 的调度语法见 compiler.py 的get_dmf_scheduleYAML schedule生成的 DATA_METRIC_SCHEDULEtype: on_table_changeTRIGGER_ON_CHANGEStype: cron如cron: 0 * * * *时区UTCUSING CRON 0 * * * * UTCtype: interval如 60 分钟60 MIN同时编译器会校验同一张表上的所有断言必须使用完全相同的调度否则抛错见 compiler.py这与下文注意事项中 Snowflake 的限制保持一致。Step 4在 Snowflake 环境中注册编译产物将 Step 3 生成的两个 SQL 文件在 Snowflake 中执行。可以直接在 Snowflake UI 中运行也可以使用 SnowSQL CLIsnowsql -f dmf_definitions.sql snowsql -f dmf_associations.sql:::note 在表上调度 Data Metric Function 会产生 SnowflakeServerless Credit 用量计费细节参考 Snowflake 官方计费说明。若某条断言不再使用请务必通过dmf_associations.sql中对应的方式 DROP 掉相关 Data Metric Function避免持续产生费用。 :::Step 5运行摄取把 DMF 结果回传到 DataHubDMF 注册完成后Snowflake 会在目标表更新时TRIGGER_ON_CHANGES或按固定计划自动执行它们。要把生成的数据质量断言结果回传到 DataHub需要以特殊配置标志运行 DataHub 摄取include_assertion_results: true# Your DataHub Snowflake Recipe source: type: snowflake config: # ... include_assertion_results: True # ...然后照常执行摄取datahub ingest -c snowflake.yml摄取过程中DataHub 会查询 Snowflake 中存储的最新 DMF 结果将其转换为 DataHub Assertion Results并在摄取期间通过 CLI 或 UI以普通断言的形式报告回来。从源码看摄取器在 snowflake_assertion.py 中查询 snowflake_query.py 定义的dmf_assertion_resultsSQL从SNOWFLAKE.LOCAL.DATA_QUALITY_MONITORING_RESULTS视图读取结果并按MEASUREMENT_TIME时间窗口对应 recipe 中的start_time/end_time过滤默认情况下include_externally_managed_dmfsFalse查询会通过METRIC_NAME ilike datahub\_\_% escape \过滤只取 DataHub 创建的 DMF结果行中VALUE被映射为断言运行状态VALUE1→ SUCCESS通过VALUE0→ FAILURE失败其他值 → ERROR见 snowflake_assertion.py每条结果生成一个AssertionRunEvent工作单元并携带MEASUREMENT_TIME作为时间戳DataHub UI 据此绘制历史时间线。单元测试 test_snowflake_assertion.py 验证了上述查询行为默认查询包含datahub前缀过滤开启外部 DMF 后则不附加任何ilike过滤。六、摄取外部用户自建DMF除了 DataHub 创建的 DMF你还可以摄取自己用 Snowflake 原生方式创建的 Data Metric Function 的结果。外部指那些未经过 DataHub assertion compiler、直接创建在 Snowflake 中的 DMF——它们独立于 DataHub 的管理范围之外。6.1 为什么要摄取外部 DMF存量 DMF在采用 DataHub 之前就已经存在于 Snowflake 的 DMF希望在 DataHub 中看到其结果而无需重建自定义逻辑需要 DataHub assertion compiler 不支持的 DMF 逻辑例如复杂的多表检查团队工作流不同团队直接在 Snowflake 中管理 DMF但你希望在 DataHub 中获得集中可见性渐进式采用在完全迁移到 DataHub 管理的断言之前先在 DataHub 中监控现有的数据质量检查。6.2 启用外部 DMF 摄取在 Snowflake recipe 中添加include_externally_managed_dmfs标志source: type: snowflake config: # ... connection config ... # 启用断言结果摄取必需 include_assertion_results: true # 启用外部 DMF 摄取新增 include_externally_managed_dmfs: true # 断言结果的时间窗口 start_time: -7 days两个标志必须同时开启外部 DMF 摄取才能生效。这一约束在配置模型层被强制校验include_externally_managed_dmfs: True而include_assertion_results未开启时配置校验会直接抛出ValueError见 snowflake_config.py。6.3 外部 DMF 的要求外部 DMF 必须以1表示 SUCCESS、以0表示 FAILURE。DataHub 将 SnowflakeDATA_QUALITY_MONITORING_RESULTS表中的VALUE列解释为VALUE 1→ 断言PASSEDVALUE 0→ 断言FAILED这是因为 DataHub 无法解释任意的返回值例如 100 个空值行——这到底是好还是坏。你必须把通过/失败逻辑构建到 DMF 本身中。:::warning 如果我的 DMF 返回其他值怎么办 如果 DMF 返回 0 或 1 以外的值DataHub 会把断言结果标记为ERRORVALUE 1→PASSEDVALUE 0→FAILEDVALUE ! 0 且 VALUE ! 1如 5、100、-1→ERRORERROR 状态表示该 DMF 未按 DataHub 摄取的要求正确配置。可通过以下方式识别这些情况检查摄取日志中的警告例如DMF my_dmf returned invalid value 100. Expected 1 (pass) or 0 (fail). Marking as ERROR.在 DataHub UI 中查找处于 ERROR 状态的断言 :::这一映射逻辑在 snowflake_assertion.py 中有直接实现日志中的警告文案与文档描述一致。示例正确编写外部 DMF错误示例—— 返回原始计数DataHub 无法解释CREATE DATA METRIC FUNCTION my_null_check(ARGT TABLE(col VARCHAR)) RETURNS NUMBER AS $$ SELECT COUNT(*) FROM ARGT WHERE col IS NULL $$; -- 返回0, 5, 100 等 —— DataHub 无法判定通过/失败正确示例—— 返回 1通过或 0失败CREATE DATA METRIC FUNCTION my_null_check(ARGT TABLE(col VARCHAR)) RETURNS NUMBER AS $$ SELECT CASE WHEN COUNT(*) 0 THEN 1 ELSE 0 END FROM ARGT WHERE col IS NULL $$; -- 返回无空值时返回 1通过存在空值时返回 0失败正确示例—— 带阈值CREATE DATA METRIC FUNCTION my_null_check_threshold(ARGT TABLE(col VARCHAR)) RETURNS NUMBER AS $$ SELECT CASE WHEN COUNT(*) 10 THEN 1 ELSE 0 END FROM ARGT WHERE col IS NULL $$; -- 返回空值 ≤10 时返回 1通过空值 10 时返回 0失败6.4 外部 DMF 与 DataHub 创建 DMF 的差异AspectDataHub-Created DMFsExternal DMFsNaming以datahub__为前缀任意名称Definition通过datahub assertions compile创建在 Snowflake 中手动创建Assertion Type基于 YAML 定义Freshness、Volume 等CUSTOMSourceNATIVE在 DataHub 中定义EXTERNALURN Generation从 DMF 名称提取datahub__guid由 Snowflake 的REFERENCE_ID生成从源码看外部 DMF 的 URN 通过SnowflakeExternalDmfKey基于REFERENCE_IDDMF-表-列关联的唯一标识与可选platform_instance生成确定性 GUID见 snowflake_assertion.py 与 测试用例从而保证同一REFERENCE_ID多次摄取生成相同的断言 URN、不同REFERENCE_ID互不冲突。而对于datahub__前缀的 DMF则直接从名称中提取 GUID。6.5 外部 DMF 在 DataHub UI 中的呈现外部 DMF 在 DataHub 中以如下形式出现Assertion TypeCUSTOMSourceEXTERNALPlatform InstanceSnowflake platform instance若已配置DescriptionExternal Snowflake DMF: {dmf_name}Custom Propertiessnowflake_dmf_nameDMF 函数名snowflake_reference_idSnowflake 对 DMF-表绑定的唯一标识snowflake_dmf_columnsDMF 作用的列列表逗号分隔这些字段在 snowflake_assertion.py 的_create_assertion_info_workunit中被写入AssertionInfo的customAssertion与customProperties单列 DMF 还会生成对应的 schema field URN从而把断言挂到具体列上。你可以在 DataHub UI 中关联数据集的Quality标签页查看外部 DMF 断言它们会与 DataHub 创建的断言一同展示通过/失败历史。七、注意事项Caveats以下限制源自 Snowflake 平台本身使用前需知晓目前 Snowflake 最多支持1000 个 DMF-表关联因此你无法为 Snowflake 定义超过 1000 条断言目前 Snowflake不允许在 DMF 定义中使用 JOIN 查询或非确定性函数因此 SQL 断言或过滤器部分不能使用这些特性目前一张表上调度的所有 DMF必须遵循完全相同的调度计划因此不能为同一张表上的不同断言设置不同的调度目前 DMF仅支持常规表不支持动态表dynamic tables或外部表external tables。其中同一表必须同调度的限制在编译器侧已被强制执行见 compiler.py仅支持常规表也与 Snowflake 摄取配置中table_types的取值BASE TABLE/EXTERNAL TABLE相互印证见 snowflake_config.py。八、FAQ原文档中该章节标注为 Coming soon!暂未提供具体问答内容。如果你在使用中遇到问题可以结合上文提到的源码路径编译器、摄取器、查询与测试定位行为细节或在 DataHub 官方渠道寻求支持。九、总结与后续探索本文完整走通了YAML 声明断言 →assertions upsert注册 →assertions compile编译 → Snowflake 注册与调度 →include_assertion_results摄取结果的 Snowflake DMF 断言全链路并说明了外部 DMF 的接入要求与 0/1 返回值约定。想深入了解的读者可以继续在仓库中探索编译产物生成与调度翻译compiler.py、dmf_generator.py、metric_sql_generator.py断言结果摄取与外部 DMF 处理snowflake_assertion.py、snowflake_query.py摄取配置项定义与校验snowflake_config.py单元测试test_snowflake_assertion.pyYAML 断言定义示例assertions_configuration.yml四种断言类型的详细概念Freshness、Volume、Column、Custom SQLSnowflake 摄取入门overview、setup、configuration。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考