DataHub Python SDK 批量创建带异常检测(Anomaly Detection)的数据质量断言实战指南
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
本篇技术指南讲解如何基于 DataHub Cloud Python SDK(acryl-datahub-cloud),在 DataHub 中以编程方式批量创建启用 Anomaly Detection(原称 Smart Assertions)的数据质量断言,覆盖表级与列级两类断言、表的发现与筛选、断言 URN 的存储与更新,以及面向生产环境的批处理与最佳实践。读完本文,你将掌握在数百张表、数千个列上规模化落地 Freshness / Volume / Column Metric / Custom SQL 四类智能断言的一整套可运行方案,并了解其底层实现原理。
背景:什么是带 Anomaly Detection 的断言
在进入批量创建之前,先理解 Anomaly Detection 的本质。根据 异常检测官方文档,Anomaly Detection 是可选的断言能力,用 AI 驱动的动态阈值替代固定阈值:你不再需要手工指定"行数必须介于 1000 到 2000",而是由模型学习底层指标的历史规律(趋势、季节性、典型波动),当最新取值落在"正常范围"之外时触发告警。
Anomaly Detection 当前支持四类断言,成熟度各不相同:
| 断言类型 | Anomaly Detection 阶段 | 关键约束 |
|---|---|---|
| Freshness 断言 | GA | 基于仓库查询或 ingestion 信号;配合 DataHuboperationaspect 可用于任何上报 Operations 的平台(含 Clickhouse、Oracle、Dremio 等) |
| Volume 断言 | GA | 基于仓库查询或 Dataset Profile;配合 Dataset Profile 可用于任何上报行数 profile 的平台(含 Iceberg、Postgres、MySQL 等) |
| Column Metric 断言 | Public Beta | 仅支持null_count、unique_count、empty_count、zero_count、negative_count指标 |
| Custom SQL 断言 | Public Beta | 需要活跃的仓库连接,仅限 Snowflake / Redshift / BigQuery / Databricks |
需要特别注意的是:Column Value 断言(如"值匹配正则""值属于集合")和 Schema 断言不支持 Anomaly Detection,前者是确定性校验、没有统计意义上的"正常",后者是离散事件而非异常。另外,新创建的 Anomaly Detection 断言会先收集14 天历史数据才正式开始告警——期间照常调度评估并记录结果,但不报失败;若想跳过等待期,可在创建时启用 Backfill Assertion History,从仓库回填指标历史,让预测从第一天起就可用。
为什么需要批量创建断言
通过 Python SDK 批量创建断言的核心价值在于:
- 规模化数据质量:跨数百甚至数千张表应用一致的断言策略;
- 自动化断言管理:基于元数据模式(标签、域名、平台、命名规律)以编程方式创建与更新断言;
- 落地治理策略:确保所有关键表都有恰当的数据质量检查;
- 节省时间:避免在 UI 中逐条手工创建。
前置条件
批量创建前需要满足以下条件:
安装 DataHub Cloud Python SDK:
pip install acryl-datahub-cloud需要说明的是,
sync_smart_*_assertion这一组 API 属于 Cloud SDK 的扩展能力。在当前开源仓库的 assertion_client.py 中,AssertionClient.__getattr__的实现明确说明:未安装acryl-datahub-cloud时,client.assertions上任何 Cloud-only 方法都会抛出SdkUsageError并提示先用pip install acryl-datahub-cloud安装。也就是说,仅安装开源acryl-datahub时只能使用sync_custom_assertion与report_assertion_result(可参考 sync_custom_assertion.py 示例),批量智能断言必须安装 Cloud SDK。配置有效的 DataHub Cloud 凭证:server URL 与具备相应权限的 access token;
发起调用的 actor 必须对目标数据集拥有Edit Assertions与Edit Monitors权限;
目标数据集必须已存在于 DataHub 实例中。如果尝试为不存在的实体创建断言,GMS 会持续向日志上报错误。
本指南目标
本指南将演示如何用 DataHub Cloud Python SDK 以编程方式创建大量启用 Anomaly Detection 的断言。
整体流程概览
批量断言创建的完整流程分七步:
- 发现表:通过搜索或直接指定表 URN 找到目标数据集;
- 创建表级断言:为每张表添加 Freshness 与 Volume 断言;
- 获取列信息:读取每张表的 schema 细节;
- 创建列级断言:为相关列添加 Column Metric 断言;
- 创建订阅:为数据集或断言创建订阅以便接收变更通知;
- 存储断言 URN:保存断言标识符供后续更新;
- 更新已有断言:基于存储的 URN 对断言参数做增量调整。
环境初始化:连接 DataHub
from datahub.sdk import DataHubClient client = DataHubClient(server="<your_server>", token="<your_token>")参数说明:
- server:DataHub GMS 服务器地址
- 本地:
http://localhost:8080 - 托管(hosted):
https://<your_datahub_url>/gms
- 本地:
- token:需先从 DataHub 实例中生成个人访问令牌(Personal Access Token)。
也可以先设置DATAHUB_GMS_URL、DATAHUB_GMS_TOKEN环境变量,或运行datahub init生成~/.datahubenv文件,再通过from_env()初始化:
from datahub.sdk import DataHubClient client = DataHubClient.from_env()并行处理的重要注意事项
- 对同一数据集的批量断言创建务必在单线程中执行,避免竞态条件(race conditions);
- 对同一数据集的订阅 API 调用也必须在单线程中执行;
- 如果直接订阅断言,请确保脚本按数据集维度单线程运行。
Step 1:发现目标表
方案 A:显式指定表 URN
from datahub.metadata.urns import DatasetUrn # Define specific tables you want to add assertions to table_urns = [ "urn:li:dataset:(urn:li:dataPlatform:snowflake,database.schema.users,PROD)", "urn:li:dataset:(urn:li:dataPlatform:snowflake,database.schema.orders,PROD)", "urn:li:dataset:(urn:li:dataPlatform:snowflake,database.schema.products,PROD)", ] # Convert to DatasetUrn objects datasets = [DatasetUrn.from_string(urn) for urn in table_urns]方案 B:按命名模式搜索表
更全面的搜索能力与筛选选项参见 Search API 文档:
from datahub.sdk.search_filters import FilterDsl from datahub.metadata.urns import DatasetUrn # Search for tables matching criteria def find_tables_by_pattern(client, platform="snowflake", name_pattern="production_*"): """Find tables matching a specific pattern.""" # Create filters for datasets on a specific platform with name pattern filters = FilterDsl.and_( FilterDsl.entity_type("dataset"), FilterDsl.platform(platform), FilterDsl.custom_filter("name", "EQUAL", [name_pattern]) ) # Use the search client to find matching datasets urns = list(client.search.get_urns(filter=filters)) return [DatasetUrn.from_string(str(urn)) for urn in urns] # Use the search function datasets = find_tables_by_pattern(client, platform="snowflake", name_pattern="production_*")从源码看,search_filters.py 中的FilterDsl提供了and_/or_/not_、entity_type、entity_subtype、platform、domain、container、env、owner、glossary_term、tag、has_custom_property、soft_deleted、custom_filter等静态工厂方法,可自由组合出精确的搜索谓词。其中FilterDsl.and_在编译阶段通过笛卡尔积合并各子句(见_And.compile的实现 search_filters.py),例如(A or B) and (C or D)会被展开为(A and C) or (A and D) or (B and C) or (B and D),因此理论上可表达任意复杂度的布尔查询。
方案 C:按标签或域名获取表
def find_tables_by_tag(client, tag_name="critical"): """Find tables with a specific tag.""" # Create filters for datasets with a specific tag filters = FilterDsl.and_( FilterDsl.entity_type("dataset"), FilterDsl.custom_filter("tags", "EQUAL", [f"urn:li:tag:{tag_name}"]) ) # Use the search client to find matching datasets urns = list(client.search.get_urns(filter=filters)) return [DatasetUrn.from_string(str(urn)) for urn in urns] # Find all tables tagged as "critical" critical_datasets = find_tables_by_tag(client, "critical")补充说明:源码中的
FilterDsl.tag()工厂方法(search_filters.py)内部使用_TagFilter,其校验器要求传入的必须是urn:li:tag:前缀的合法 tag URN(见 search_filters.py),因此上例中手写f"urn:li:tag:{tag_name}"的拼接方式与FilterDsl.tag("urn:li:tag:critical")等价。此外,custom_filter的condition参数支持EQUAL、CONTAIN、START_WITH、END_WITH、GREATER_THAN、LESS_THAN等条件,可用于数值或时间戳字段的范围过滤。
Step 2:创建表级断言
创建断言前先准备一个 URN 注册表,用于存放后续创建出的断言标识符:
# Storage for assertion URNs (for later updates) assertion_registry = { "freshness": {}, "volume": {}, "smart_sql": {}, "column_metrics": {} }Freshness 断言 + Anomaly Detection
def create_freshness_assertions(datasets, client, registry): """Create Freshness assertions with Anomaly Detection for multiple datasets.""" for dataset_urn in datasets: try: freshness_assertion = client.assertions.sync_smart_freshness_assertion( dataset_urn=dataset_urn, display_name=f"Freshness Anomaly Monitor", # Detection mechanism - information_schema is recommended detection_mechanism="information_schema", # AI sensitivity setting sensitivity="medium", # options: "low", "medium", "high" # Tags for grouping (supports urns or plain tag names!) tags=["automated", "freshness", "data_quality"], # Enable the assertion enabled=True ) # Store the assertion URN for future reference registry["freshness"][str(dataset_urn)] = str(freshness_assertion.urn) print(f"✅ Created freshness assertion for {dataset_urn.name}: {freshness_assertion.urn}") except Exception as e: print(f"❌ Failed to create freshness assertion for {dataset_urn.name}: {e}") # Create freshness assertions for all datasets create_freshness_assertions(datasets, client, assertion_registry)Volume 断言 + Anomaly Detection
def create_volume_assertions(datasets, client, registry): """Create Volume assertions with Anomaly Detection for multiple datasets.""" for dataset_urn in datasets: try: volume_assertion = client.assertions.sync_smart_volume_assertion( dataset_urn=dataset_urn, display_name=f"Volume Anomaly Monitor", # Detection mechanism options detection_mechanism="information_schema", # AI sensitivity setting sensitivity="medium", # Tags for grouping tags=["automated", "volume", "data_quality"], # Schedule (optional - defaults to hourly) schedule="0 */6 * * *", # Every 6 hours # Enable the assertion enabled=True ) # Store the assertion URN registry["volume"][str(dataset_urn)] = str(volume_assertion.urn) print(f"✅ Created volume assertion for {dataset_urn.name}: {volume_assertion.urn}") except Exception as e: print(f"❌ Failed to create volume assertion for {dataset_urn.name}: {e}") # Create volume assertions for all datasets create_volume_assertions(datasets, client, assertion_registry)Custom SQL 断言 + Anomaly Detection(Public Beta)
def create_smart_sql_assertions(datasets, client, registry): """Create Custom SQL assertions with Anomaly Detection for multiple datasets.""" # Define SQL queries to run on each table sql_queries = { "row_count": "SELECT COUNT(*) FROM {table_name}", "null_check": "SELECT COUNT(*) FROM {table_name} WHERE id IS NULL", "active_records": "SELECT COUNT(*) FROM {table_name} WHERE status = 'active'", } for dataset_urn in datasets: registry["smart_sql"][str(dataset_urn)] = {} for query_name, query_template in sql_queries.items(): try: table_name = dataset_urn.name statement = query_template.format(table_name=table_name) sql_assertion = client.assertions.sync_smart_sql_assertion( dataset_urn=dataset_urn, display_name=f"SQL Anomaly Monitor - {query_name}", statement=statement, # AI-powered sensitivity setting sensitivity="medium", # options: "low", "medium", "high" # Tags for grouping tags=["automated", "anomaly_detection", query_name], # Schedule schedule="0 */6 * * *", # Every 6 hours # Enable the assertion enabled=True ) registry["smart_sql"][str(dataset_urn)][query_name] = str(sql_assertion.urn) print(f"✅ Created Custom SQL anomaly monitor '{query_name}' for {dataset_urn.name}: {sql_assertion.urn}") except Exception as e: print(f"❌ Failed to create Custom SQL anomaly monitor '{query_name}' for {dataset_urn.name}: {e}") # Create Custom SQL anomaly monitors for all datasets create_smart_sql_assertions(datasets, client, assertion_registry)关键参数说明
sensitivity(灵敏度):"low"/"medium"/"high"三档。灵敏度越高,模型对数据的拟合越紧、越容易触发告警;越低则容忍更大的数据波动。这对应 异常检测文档 中"Tuning"一节所述的灵敏度调优手段。detection_mechanism:检测机制。information_schema表示通过仓库的 information_schema 查询信号;列级断言中还会见到all_rows_query_datahub_dataset_profile(基于 DataHub Dataset Profile 信号)。schedule:cron 表达式调度,可选项(Volume 与 Custom SQL 默认每小时一次)。tags:支持传入 URN 或普通标签名——普通标签名会自动转换为urn:li:tag:<name>形式,这是本指南末尾会重点强调的易用性特性。
Step 3:获取列信息
要为列创建断言,必须先读取数据集的 schema:
def get_dataset_columns(client, dataset_urn): """Get column information for a dataset.""" try: # Get dataset using the entities client dataset = client.entities.get(dataset_urn) if dataset and hasattr(dataset, 'schema') and dataset.schema: return [ { "name": field.field_path, "type": field.native_data_type, "nullable": field.nullable if hasattr(field, 'nullable') else True } for field in dataset.schema.fields ] return [] except Exception as e: print(f"❌ Failed to get columns for {dataset_urn}: {e}") return [] # Get columns for each dataset dataset_columns = {} for dataset_urn in datasets: columns = get_dataset_columns(client, dataset_urn) dataset_columns[str(dataset_urn)] = columns print(f"📊 Found {len(columns)} columns in {dataset_urn.name}")这里通过client.entities.get(dataset_urn)拉取实体,从dataset.schema.fields中读取每个字段的field_path(列名)、native_data_type(原生数据类型,如VARCHAR、INTEGER)与nullable标记。列的类型信息是下一步"按规则筛选列"的关键输入。
Step 4:创建列级断言
Column Metric 断言 + Anomaly Detection(Public Beta)
def create_column_assertions(datasets, columns_dict, client, registry): """Create Column Metric assertions with Anomaly Detection for multiple datasets and columns.""" # Define rules for which columns should get which assertions assertion_rules = { # Null count checks for critical columns "null_checks": { "column_patterns": ["id", "*_id", "user_id", "email"], "metric_type": "null_count", }, # Unique count checks for ID columns "unique_checks": { "column_patterns": ["*_id", "email", "username"], "metric_type": "unique_count", }, # Empty count checks for string columns "empty_checks": { "column_patterns": ["name", "description", "title"], "metric_type": "empty_count", }, } for dataset_urn in datasets: dataset_key = str(dataset_urn) columns = columns_dict.get(dataset_key, []) if not columns: print(f"⚠️ No columns found for {dataset_urn.name}") continue registry["column_metrics"][dataset_key] = {} for column in columns: column_name = column["name"] column_type = column["type"].upper() # Apply assertion rules based on column name and type for rule_name, rule_config in assertion_rules.items(): if should_apply_rule(column_name, column_type, rule_config): try: assertion = client.assertions.sync_smart_column_metric_assertion( dataset_urn=dataset_urn, column_name=column_name, metric_type=rule_config["metric_type"], display_name=f"{rule_name.replace('_', ' ').title()} - {column_name}", # Detection mechanism for column metrics detection_mechanism="all_rows_query_datahub_dataset_profile", # Tags (plain names automatically converted to URNs) tags=["automated", "column_quality", rule_name], enabled=True ) # Store assertion URN if column_name not in registry["column_metrics"][dataset_key]: registry["column_metrics"][dataset_key][column_name] = {} registry["column_metrics"][dataset_key][column_name][rule_name] = str(assertion.urn) print(f"✅ Created {rule_name} assertion for {dataset_urn.name}.{column_name}") except Exception as e: print(f"❌ Failed to create {rule_name} assertion for {dataset_urn.name}.{column_name}: {e}") def should_apply_rule(column_name, column_type, rule_config): """Determine if a rule should be applied to a column.""" import fnmatch # Check column name patterns for pattern in rule_config["column_patterns"]: if fnmatch.fnmatch(column_name.lower(), pattern.lower()): return True # Add type-based rules if needed if rule_config.get("column_types"): return any(col_type in column_type for col_type in rule_config["column_types"]) return False # Create column assertions create_column_assertions(datasets, dataset_columns, client, assertion_registry)这段代码演示了一个典型的"规则引擎"式列筛选:assertion_rules定义了按列名通配符(fnmatch)匹配的规则集,should_apply_rule决定某列是否命中规则,命中后调用sync_smart_column_metric_assertion创建断言。如需扩展,可在rule_config中增加"column_types"键来按数据类型(如VARCHAR、NUMERIC)补充筛选。
Step 5:创建订阅
关于如何在数据集或断言上创建订阅,参见订阅 SDK 教程。
注意:批量创建订阅时,必须单线程执行以避免竞态条件。另外,强烈建议在数据集级别创建订阅而不是为单个断言逐一创建订阅,这样后续的持续管理会简单得多。
Step 6:存储断言 URN
断言创建完成后,建议将 URN 注册表持久化到文件,便于未来更新与审计。
保存到文件
import json from datetime import datetime def save_assertion_registry(registry, filename=None): """Save assertion URNs to a file for future reference.""" if filename is None: timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") filename = f"assertion_registry_{timestamp}.json" # Add metadata registry_with_metadata = { "created_at": datetime.now().isoformat(), "total_assertions": { "freshness": len(registry["freshness"]), "volume": len(registry["volume"]), "column_metrics": sum( len(cols) for cols in registry["column_metrics"].values() ) }, "assertions": registry } with open(filename, 'w') as f: json.dump(registry_with_metadata, f, indent=2) print(f"💾 Saved assertion registry to {filename}") return filename # Save the registry registry_file = save_assertion_registry(assertion_registry)从文件加载(用于后续更新)
def load_assertion_registry(filename): """Load assertion URNs from a previously saved file.""" with open(filename, 'r') as f: data = json.load(f) return data["assertions"] # Later, load for updates # assertion_registry = load_assertion_registry("assertion_registry_20240101_120000.json")Step 7:更新已有断言
由于sync_smart_*_assertion系列是"同步(sync)"语义,传入已存在的urn即可原地更新,无需先删除再创建:
def update_existing_assertions(registry, client): """Update existing assertions using stored URNs.""" # Update freshness assertions for dataset_urn, assertion_urn in registry["freshness"].items(): try: updated_assertion = client.assertions.sync_smart_freshness_assertion( dataset_urn=dataset_urn, urn=assertion_urn, # Provide existing URN for updates # Update any parameters as needed sensitivity="high", # Change sensitivity tags=["automated", "freshness", "data_quality", "updated"], enabled=True ) print(f"🔄 Updated freshness assertion {assertion_urn}") except Exception as e: print(f"❌ Failed to update freshness assertion {assertion_urn}: {e}") # Update assertions when needed # update_existing_assertions(assertion_registry, client)高级模式
模式一:基于元数据条件的条件式断言创建
可以结合实体的标签、属性等元数据,对不同类型的表施加差异化策略——例如对打了critical标签的表使用更高灵敏度:
def create_conditional_assertions(datasets, client): """Create assertions based on dataset metadata conditions.""" for dataset_urn in datasets: try: # Get dataset metadata dataset = client.entities.get(dataset_urn) # Check if dataset has specific tags if dataset.tags and any("critical" in str(tag.tag) for tag in dataset.tags): # Create more stringent assertions for critical datasets client.assertions.sync_smart_freshness_assertion( dataset_urn=dataset_urn, sensitivity="high", detection_mechanism="information_schema", tags=["critical", "automated", "freshness"] ) # Check dataset size and apply appropriate volume checks if dataset.dataset_properties: # Create different volume assertions based on table characteristics pass except Exception as e: print(f"❌ Error processing {dataset_urn}: {e}")模式二:带错误处理与速率限制的批处理
面对大规模数据集,建议分批提交并在批次间加入延时,同时收集成功/失败明细:
import time from typing import List, Dict, Any def batch_create_assertions( datasets: List[DatasetUrn], client: DataHubClient, batch_size: int = 10, delay_seconds: float = 1.0 ) -> Dict[str, Any]: """Create assertions in batches with error handling and rate limiting.""" results = { "successful": [], "failed": [], "total_processed": 0 } for i in range(0, len(datasets), batch_size): batch = datasets[i:i + batch_size] print(f"Processing batch {i//batch_size + 1}: {len(batch)} datasets") for dataset_urn in batch: try: # Create assertion assertion = client.assertions.sync_smart_freshness_assertion( dataset_urn=dataset_urn, tags=["batch_created", "automated"], enabled=True ) results["successful"].append({ "dataset_urn": str(dataset_urn), "assertion_urn": str(assertion.urn) }) except Exception as e: results["failed"].append({ "dataset_urn": str(dataset_urn), "error": str(e) }) results["total_processed"] += 1 # Rate limiting between batches if i + batch_size < len(datasets): time.sleep(delay_seconds) return results # Use batch processing batch_results = batch_create_assertions(datasets, client, batch_size=5) print(f"Batch results: {batch_results['total_processed']} processed, " f"{len(batch_results['successful'])} successful, " f"{len(batch_results['failed'])} failed")最佳实践
1. 标签策略
- 使用一致的标签名对断言分组,例如
["automated", "freshness", "critical"]; - 普通标签名会自动转换为 URN:
"my_tag"→"urn:li:tag:my_tag",无需手工拼 URN; - 为不同类型的断言与优先级建立标签层级,便于后续按标签批量检索、订阅或审计。
2. 错误处理
- 始终用 try-catch 包裹断言创建逻辑;
- 将失败信息记录下来供事后排查;
- 为瞬时故障实现重试逻辑。
3. URN 管理
- 将断言 URN 存放到持久化位置(文件、数据库等);
- 文件名使用带时间戳的有意义命名;
- 记录断言创建的时间与原因等元数据。
4. 性能考量
后台架构面向大规模操作设计,但写入是异步提交到 Kafka 队列的,大规模操作可能存在明显延迟。若遇到问题,可参考以下建议:
- 错峰执行:避免大批量操作造成 Kafka lag 尖峰;
- 重跑 sync 前等待:更新前先等 GMS 完成上一轮处理,通过检查最近一条数据是否已反映在 GMS 中来避免不一致与重复;
- 监控处理状态:通过 DataHub UI 或 API 确认操作全部完成;
- 分批处理数据集:避免一次性压垮 API;
- 必要时在批次间加入延时。
5. 测试策略
- 先用一小部分数据集做试点;
- 在大规模处理前先验证断言创建是否正常;
- 用已存在的断言测试更新场景。
完整示例脚本
以下脚本把上述步骤串成一个可直接参考的端到端流程(实际运行时请先补齐create_freshness_assertions、create_volume_assertions、get_dataset_columns、create_column_assertions、save_assertion_registry等上文定义的函数):
#!/usr/bin/env python3 """ Complete example script for bulk creating assertions with Anomaly Detection enabled. """ import json import time from datetime import datetime from typing import List, Dict, Any from datahub.sdk import DataHubClient from datahub.ingestion.graph.client import DataHubGraph from datahub.metadata.urns import DatasetUrn def main(): # Initialize the DataHub client client = DataHubClient( server="https://your-datahub-instance.com", token="your-access-token", ) # The client provides both search and entity access # Define target datasets table_urns = [ "urn:li:dataset:(urn:li:dataPlatform:snowflake,prod.analytics.users,PROD)", "urn:li:dataset:(urn:li:dataPlatform:snowflake,prod.analytics.orders,PROD)", "urn:li:dataset:(urn:li:dataPlatform:snowflake,prod.analytics.products,PROD)", ] datasets = [DatasetUrn.from_string(urn) for urn in table_urns] # Registry to store assertion URNs assertion_registry = { "freshness": {}, "volume": {}, "column_metrics": {} } print(f"🚀 Starting bulk assertion creation for {len(datasets)} datasets") # Step 1: Create table-level assertions print("\n📋 Creating freshness assertions...") create_freshness_assertions(datasets, client, assertion_registry) print("\n📊 Creating volume assertions...") create_volume_assertions(datasets, client, assertion_registry) # Step 2: Get column information and create column assertions print("\n🔍 Analyzing columns and creating column assertions...") dataset_columns = {} for dataset_urn in datasets: columns = get_dataset_columns(client, dataset_urn) dataset_columns[str(dataset_urn)] = columns create_column_assertions(datasets, dataset_columns, client, assertion_registry) # Step 3: Save results print("\n💾 Saving assertion registry...") registry_file = save_assertion_registry(assertion_registry) # Summary total_assertions = ( len(assertion_registry["freshness"]) + len(assertion_registry["volume"]) + sum(len(cols) for cols in assertion_registry["column_metrics"].values()) ) print(f"\n✅ Bulk assertion creation complete!") print(f" 📈 Total assertions created: {total_assertions}") print(f" 🕐 Freshness assertions: {len(assertion_registry['freshness'])}") print(f" 📊 Volume assertions: {len(assertion_registry['volume'])}") print(f" 🎯 Column assertions: {sum(len(cols) for cols in assertion_registry['column_metrics'].values())}") print(f" 💾 Registry saved to: {registry_file}") if __name__ == "__main__": main()延伸阅读与底层线索
- 理解 Anomaly Detection 的完整机制(14 天学习期、灵敏度调优、训练数据回看窗口、异常反馈、时间序列分桶等),见异常检测文档及其子页面:Freshness 断言、Volume 断言、Column Metric 断言、Custom SQL 断言、Backfill Assertion History;
- 搜索客户端的完整过滤能力(
F.platform、F.env、F.entity_type、F.domain、F.soft_deleted、F.has_custom_property以及and_/or_/not_逻辑组合),见Search API 文档与 search_filters.py; - 在数据集或断言上创建订阅的用法与可用变更类型,见订阅教程;
- 开源 SDK 自带的 CUSTOM 断言上报示例(用于外部监控工具自报数据质量),见 sync_custom_assertion.py。
综上,借助 DataHub Cloud Python SDK 的sync_smart_*_assertion系列接口,配合搜索筛选、规则化列匹配、URN 注册表持久化与批处理节奏控制,即可在规模化场景下稳定地落地带 Anomaly Detection 的智能数据质量监控。其中标签名的自动 URN 转换特性(普通标签名 →urn:li:tag:<name>)进一步简化了断言的分类组织与后续管理。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考