【Atlas】Atlas 是否支持数据质量规则的定义与执行?
Apache Atlas 2.4.0 是否支持数据质量规则的定义与执行?——元数据驱动 DQ 治理的真相与实践
用户问题原文:Atlas 是否支持数据质量规则的定义与执行?
本文将直面这一高频误解,系统性地澄清 Apache Atlas 2.4.0 在数据质量(Data Quality, DQ)领域的实际能力边界,并基于 IoT 设备指标元数据治理 的真实场景,构建一套 以 Atlas 为元数据中枢、联动外部 DQ 引擎(如 Great Expectations、Deequ)的生产级解决方案。我们将通过源码剖析、架构图、配置示例与验证命令,揭示“元数据驱动 DQ”的正确打开方式。
一、问题引入:IoT 设备上报的温度字段为何出现 -999?
某工业物联网平台的数据团队收到告警:设备指标表 iot_device_metrics_hudi 中的 temperature 字段出现大量 -999 异常值,导致下游预测模型失效。
团队尝试在 Atlas 中为该字段添加“数据质量规则”,却发现:
- Atlas UI 无“数据质量”配置入口;
- REST API 无 DQ 规则相关端点;
- 官方文档未提及 DQ 执行能力。
根本原因在于:Apache Atlas 本身不提供数据质量规则的定义存储与执行引擎。它仅能作为 DQ 元数据的载体,记录“某字段应满足什么规则”以及“最近一次 DQ 检查结果如何”。
关键界定:
- “定义”:指 DQ 规则的声明(如 “temperature 应在 -50 到 100 之间”)。
- “执行”:指对实际数据运行规则并生成结果(Pass/Fail + 统计值)。
结论前置:Atlas 支持 DQ 元数据的存储(即“定义”的持久化),但不负责“执行”。
二、原理解析:Atlas 的 DQ 能力边界与设计哲学
2.1 官方立场与源码佐证
Apache Atlas 项目 从未将数据质量执行纳入核心功能。其设计哲学是 “元数据管理平台”而非“数据处理引擎”。
-
GitHub Issue 明确表态:
ATLAS-3128 中,社区成员询问 DQ 集成方案,Committer 回复:“Atlas can store metadata about data quality rules and results, but execution should be handled by dedicated DQ tools.”
-
源码结构验证:
在apache/atlas仓库中,无任何包路径包含quality、dq、validation等关键词。核心模块聚焦于:repository/:Entity 存储webapp/:REST APIaddons/:Hive/Kafka 等 Hook
通俗类比:
Atlas 就像 医院的电子病历系统(EMR)——它可以记录“患者需每日测血压(规则定义)”和“今日血压 120/80(检查结果)”,但 不负责拿血压计测量(执行)。
技术本质差异:EMR 是信息记录系统;血压计是专用医疗设备。Atlas 是元数据存储;DQ 引擎是专用计算框架。
2.2 Atlas 如何承载 DQ 元数据?
虽然不执行 DQ,但 Atlas 提供了两种机制存储 DQ 相关信息:
机制 1:自定义 Type System 扩展
通过定义新的 Entity Type 来描述 DQ 规则与结果。
-
DQ Rule Entity:
{ "name": "dq_rule", "superTypes": ["Referenceable"], "typeVersion": "1.0", "attributeDefs": [ {"name": "ruleType", "typeName": "string"}, // e.g., "range_check" {"name": "minValue", "typeName": "float"}, {"name": "maxValue", "typeName": "float"}, {"name": "targetField", "typeName": "string"} // qualifiedName of field ] } -
DQ Result Entity:
{ "name": "dq_result", "superTypes": ["Referenceable"], "typeVersion": "1.0", "attributeDefs": [ {"name": "ruleGuid", "typeName": "string"}, // 关联 dq_rule {"name": "status", "typeName": "string"}, // PASS/FAIL {"name": "timestamp", "typeName": "long"}, {"name": "metricValue", "typeName": "float"} // 如 null_count ] }
机制 2:利用 Classification 附加 DQ 属性
为数据资产打上 DQ 相关标签,并携带属性。
// Classification: DQ_VALIDATED
{
"name": "DQ_VALIDATED",
"attributeDefs": [
{"name": "lastCheckTime", "typeName": "date"},
{"name": "validityScore", "typeName": "float"}
]
}
核心限制:
这些 Entity/Classification 仅是静态元数据。Atlas 不会自动触发 DQ 作业,也不会验证dq_result是否与实际数据一致。
三、架构全景:元数据驱动的 DQ 治理流水线
3.1 整体架构图(Mermaid)
颜色说明:
- #333:数据源
- #00f:Atlas 核心
- #f96:DQ 执行引擎
- #0f0:用户交互
3.2 核心组件职责
| 组件 | 职责 | 技术选型 |
|---|---|---|
| Atlas | 存储 DQ 规则定义、检查结果、关联数据资产 | Apache Atlas 2.4.0 |
| DQ Scheduler | 定时触发 DQ 作业,从 Atlas 获取规则 | Airflow / Cron |
| DQ Engine | 执行实际数据校验 | Great Expectations / Deequ / Soda Core |
| Data Source | 提供待校验数据 | Hudi / Hive / ClickHouse |
四、实战配置:构建 IoT 场景下的 DQ 元数据闭环
4.1 步骤 1:扩展 Atlas Type System
创建 dq_rule 和 dq_result 类型:
# 创建 dq_rule 类型
curl -u admin:admin -X POST \
-H "Content-Type: application/json" \
-d '{
"entityDefs": [{
"name": "dq_rule",
"superTypes": ["Referenceable"],
"typeVersion": "1.0",
"attributeDefs": [
{"name": "ruleType", "typeName": "string"},
{"name": "minValue", "typeName": "float"},
{"name": "maxValue", "typeName": "float"},
{"name": "targetField", "typeName": "string"}
]
}]
}' \
http://atlas:21000/api/atlas/v2/types/typedefs
# 创建 dq_result 类型
curl -u admin:admin -X POST \
-H "Content-Type: application/json" \
-d '{
"entityDefs": [{
"name": "dq_result",
"superTypes": ["Referenceable"],
"typeVersion": "1.0",
"attributeDefs": [
{"name": "ruleGuid", "typeName": "string"},
{"name": "status", "typeName": "string"},
{"name": "timestamp", "typeName": "long"},
{"name": "metricValue", "typeName": "float"}
]
}]
}' \
http://atlas:21000/api/atlas/v2/types/typedefs
4.2 步骤 2:为 IoT 字段定义 DQ 规则
假设 iot_device_metrics_hudi 表的 temperature 字段 qualifiedName 为:hudi.iot_db.iot_device_metrics_hudi.temperature@prod
# 获取字段 GUID
FIELD_GUID=$(curl -s -u admin:admin \
"http://atlas:21000/api/atlas/v2/entity/uniqueAttribute/type/hudi_column?attr:qualifiedName=hudi.iot_db.iot_device_metrics_hudi.temperature@prod" \
| jq -r '.entity.guid')
# 创建 DQ Rule Entity
RULE_GUID=$(curl -s -u admin:admin -X POST \
-H "Content-Type: application/json" \
-d '{
"entity": {
"typeName": "dq_rule",
"attributes": {
"ruleType": "range_check",
"minValue": -50.0,
"maxValue": 100.0,
"targetField": "hudi.iot_db.iot_device_metrics_hudi.temperature@prod"
}
}
}' \
http://atlas:21000/api/atlas/v2/entity | jq -r '.mutatedEntities.CREATE[0].guid')
4.3 步骤 3:开发 DQ 执行脚本(Python + Great Expectations)
# dq_executor.py
import great_expectations as gx
from great_expectations.core import ExpectationSuite
import requests
import json
ATLAS_URL = "http://atlas:21000"
ATLAS_AUTH = ("admin", "admin")
HUDI_TABLE_PATH = "s3a://iot-bucket/iot_device_metrics_hudi"
def get_dq_rules_from_atlas():
# 查询所有 dq_rule
resp = requests.get(
f"{ATLAS_URL}/api/atlas/v2/search/basic?typeName=dq_rule",
auth=ATLAS_AUTH
)
return resp.json()["entities"]
def run_dq_check(rule):
context = gx.get_context()
datasource = context.sources.add_pandas("iot_datasource")
asset = datasource.add_dataframe_asset(name="iot_data")
# 从 Hudi 读取最新分区(简化)
df = spark.read.format("hudi").load(HUDI_TABLE_PATH).toPandas()
validator = context.get_validator(batch_request=asset.build_batch_request(dataframe=df))
# 应用规则
if rule["attributes"]["ruleType"] == "range_check":
min_val = rule["attributes"]["minValue"]
max_val = rule["attributes"]["maxValue"]
field = rule["attributes"]["targetField"].split(".")[-2] # extract field name
validator.expect_column_values_to_be_between(
column=field, min_value=min_val, max_value=max_val
)
result = validator.validate()
return result.success, result.results[0].result["unexpected_percent"]
def report_result_to_atlas(rule_guid, status, metric_value):
payload = {
"entity": {
"typeName": "dq_result",
"attributes": {
"ruleGuid": rule_guid,
"status": "PASS" if status else "FAIL",
"timestamp": int(time.time() * 1000),
"metricValue": metric_value
}
}
}
requests.post(f"{ATLAS_URL}/api/atlas/v2/entity", auth=ATLAS_AUTH, json=payload)
# 主流程
rules = get_dq_rules_from_atlas()
for rule in rules:
success, unexpected_pct = run_dq_check(rule)
report_result_to_atlas(rule["guid"], success, unexpected_pct)
⚠️ 危险操作警告:
DQ 脚本必须处理 网络超时、认证失败、Schema 变更 等异常,避免因单次失败导致整个 DQ 流水线中断。建议增加重试与死信队列机制。
4.4 步骤 4:调度 DQ 作业
在 Airflow 中创建 DAG:
# airflow_dag.py
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import timedelta
dag = DAG(
'iot_dq_validation',
schedule_interval=timedelta(hours=1),
start_date=datetime(2026, 4, 25),
catchup=False
)
run_dq = BashOperator(
task_id='run_dq',
bash_command='python /opt/dq/dq_executor.py',
dag=dag
)
4.5 验证端到端流程
验证点 1:确认 DQ Rule 已注册
curl -u admin:admin \
"http://atlas:21000/api/atlas/v2/entity/guid/$RULE_GUID"
# 预期输出包含 ruleType="range_check", minValue=-50.0
验证点 2:检查 DQ Result 是否上报
# 查询最近的 dq_result
curl -u admin:admin \
"http://atlas:21000/api/atlas/v2/search/basic?typeName=dq_result&sortBy=timestamp&sortOrder=DESC"
# 预期输出包含 status="FAIL", metricValue=5.2(表示 5.2% 异常)
验证点 3:在 Atlas UI 查看 DQ 状态
- 访问
http://atlas:21000 - 搜索
iot_device_metrics_hudi - 在字段详情页查看关联的
dq_result
五、高级模式:血缘驱动的 DQ 传播
5.1 场景:下游宽表继承上游 DQ 规则
当 iot_device_metrics_hudi 被加工为 iot_daily_summary,可自动复制 DQ 规则:
// 在 process Entity 中建立 relationship
{
"relationshipAttributes": {
"inputs": [{"guid": "<hudi_table_guid>"}],
"outputs": [{"guid": "<summary_table_guid>"}]
}
}
DQ Scheduler 可遍历血缘图,为下游表自动创建规则副本。
5.2 实现要点
- 使用 Atlas Relationship API 查询血缘:
GET /api/atlas/v2/lineage/entity/guid/{guid} - 规则复制时更新
targetField的 qualifiedName。
六、FAQ:高频关联问题解答
Q1:Atlas 能替代 Great Expectations 吗?
不能。Atlas 无数据读取、计算、断言能力。二者是互补关系:Atlas 管“规则是什么”,Great Expectations 管“规则是否满足”。
Q2:是否有开源项目实现 Atlas + DQ 集成?
- Marquez:侧重作业血缘,DQ 支持有限。
- OpenMetadata:内置 DQ 测试(基于 YAML),但执行仍需外部引擎。
- 自研是主流:金融/电商头部公司均采用“Atlas + 自研 DQ 调度器”模式。
Q3:DQ Result 如何影响数据消费?
- 策略联动:若
dq_result.status == FAIL,Ranger 可拒绝非授权访问。 - 数据目录标注:在数据地图中显示“DQ 风险等级”。
Q4:性能瓶颈在哪里?
- Atlas 写入:高频 DQ 结果上报可能导致 HBase 写入压力。
- 优化方案:批量上报(
/api/atlas/v2/entity/bulk)、结果采样(仅存 FAIL 结果)。
Q5:云上如何实现?
- AWS:Glue Data Catalog(Atlas 替代) + Glue Data Quality(执行引擎)。
- Azure:Purview(含 DQ 元数据) + Synapse Data Quality。
七、总结与最佳实践
- 认清边界:Atlas 是 DQ 元数据的“记事本”,不是“计算器”。
- IoT 场景最佳实践:
- 规则版本化:将
dq_rule定义存入 Git,与表 Schema 变更联动。 - 结果聚合:按天/周汇总
dq_result,生成 DQ 趋势报告。 - 告警分级:
FAIL触发 PagerDuty,WARNING发送 Slack。
- 规则版本化:将
避坑指南:
- 避免在 Atlas 中存储原始 DQ 日志(如每行错误详情),应仅存摘要指标。
- DQ 规则变更时,务必同步更新下游血缘表的规则副本。
作者署名:九师兄
注意:本文由 AI 辅助生成,技术细节请以官方文档为准。生产环境使用前务必充分测试。
更多推荐



所有评论(0)