将数据质量结果写入 Data Catalog 表
您可以将 AWS Glue 数据质量评估运行配置为自动将结果写入 AWS Glue Data Catalog 中的 Apache Iceberg 表。启用结果输出后,您可以直接使用相关工具查询数据质量结果,使用可视化工具构建控制面板,并在整个账户中集中维护数据质量结果的历史记录。
您可以将以下类型的数据质量结果写入 Data Catalog 表:
-
规则结果:规则集中每条规则的通过或失败结果,包括评估的指标和失败原因
-
分析结果:分析器收集的统计数据,包括标量值(例如平均值和标准差)和分布数据(直方图和值分布)
-
行级结果:每条记录的评估结果,用于确定数据集中哪些特定行通过或未通过每条规则
-
观测结果:异常检测预测,包括预期值、预测界限以及实际值是否被标记为异常
先决条件
要将数据质量结果写入 Data Catalog 表,您用于评估运行的 IAM 角色必须具有以下权限:
-
在 AWS Glue Data Catalog 中创建和更新数据库及表的权限
-
对存储 Iceberg 表数据的 Amazon S3 位置的写入权限
评估运行使用您指定的 IAM 角色写入结果表。这与有权访问源数据表的角色相同。
配置结果输出
您可以使用 StartDataQualityRulesetEvaluationRun API 的 --additional-run-options 参数或 AWS Glue ETL 作业中的 additional_options 参数来配置数据质量结果输出。默认情况下,AWS Glue 数据质量自动监测功能数据质量自动监测功能不会将结果写入 Data Catalog 表。您必须明确启用要写入的每种结果类型。
每种结果类型都有其自己的配置块,并具有共享 CatalogTableConfig 结构。如果您不提供 CatalogTableConfig,AWS Glue 数据质量自动监测功能会自动派生默认值,包括表名称和 Amazon S3 路径。
CatalogTableConfig 结构包含以下字段:
-
DatabaseName(可选):目标表所在的目录数据库的名称。如果未指定,则创建默认数据库。
-
TableName(可选):目标表的名称。如果未指定,则使用默认表名称。
-
S3Location(可选):存储表数据的 Amazon S3 位置。格式:
s3://。如果未指定,则结果将存储在默认位置。amzn-s3-demo-bucket/prefix/ -
CatalogId(可选):要在其中创建表的 AWS Glue Data Catalog 的 ID。如果未指定,则默认使用 AWS 账户 ID。
示例:配置规则结果和分析结果
aws glue start-data-quality-ruleset-evaluation-run \ --data-source '{ "GlueTable": { "DatabaseName": "my_database", "TableName": "my_table" } }' \ --role "arn:aws:iam::123456789012:role/GlueServiceRole" \ --ruleset-names '["my_ruleset"]' \ --additional-run-options '{ "DataQualityRuleResults": { "WriteDataQualityRuleResultsEnabled": true, "CatalogTableConfig": { "DatabaseName": "quality_results", "TableName": "rule_results" } }, "ProfilingResults": { "WriteProfilingResultsEnabled": true, "CatalogTableConfig": { "DatabaseName": "quality_results", "TableName": "profiles" } } }'
示例:配置行级结果
对于行级结果,您还可以指定要包含的记录类型以及要写入的最大行数。
aws glue start-data-quality-ruleset-evaluation-run \ --data-source '{ "GlueTable": { "DatabaseName": "my_database", "TableName": "my_table" } }' \ --role "arn:aws:iam::123456789012:role/GlueServiceRole" \ --ruleset-names '["my_ruleset"]' \ --additional-run-options '{ "RowLevelResults": { "MaxRowsToWrite": 5000, "ResultType": "FAILED_ONLY", "CatalogTableConfig": { "DatabaseName": "quality_results", "TableName": "row_level_results" } } }'
ResultType 参数接受以下值:
-
FAILED_ONLY:仅写入至少未通过一条数据质量规则的行。 -
PASSED_ONLY:仅写入通过所有数据质量规则的行。 -
ALL:写入所有行及其评估结果。
示例:在 AWS Glue ETL 作业中进行配置
在 AWS Glue ETL 作业中,您可以使用带点符号键的 additional_options 参数配置结果输出:
result = EvaluateDataQuality.process_rows( frame=dynamic_frame, ruleset=ruleset, publishing_options={ "dataQualityEvaluationContext": "my_context", "enableDataQualityResultsPublishing": True }, additional_options={ "observations.scope": "ALL", "dataQualityResultsPublishing.strategy": "BEST_EFFORT", "dataQualityResultsPublishing.resultsFormat.profilingResults.writeProfilingResultsEnabled": "true", "dataQualityResultsPublishing.resultsFormat.profilingResults.catalogTableConfig.databaseName": "my_db", "dataQualityResultsPublishing.resultsFormat.profilingResults.catalogTableConfig.tableName": "profiling_results", "dataQualityResultsPublishing.resultsFormat.profilingResults.catalogTableConfig.s3Location": "s3://amzn-s3-demo-bucket/profiling/", "dataQualityResultsPublishing.resultsFormat.profilingResults.catalogTableConfig.catalogId": "123456789012" } )
示例:配置观测结果
您可以像配置其他结果类型一样配置观测结果。观测结果需要启用异常检测功能(ObservationScope: ALL):
aws glue start-data-quality-ruleset-evaluation-run \ --data-source '{ "GlueTable": { "DatabaseName": "my_database", "TableName": "my_table" } }' \ --role "arn:aws:iam::123456789012:role/GlueServiceRole" \ --ruleset-names '["my_ruleset"]' \ --additional-run-options '{ "ObservationScope": "ALL", "ObservationResults": { "WriteObservationResultsEnabled": true, "CatalogTableConfig": { "DatabaseName": "quality_results", "TableName": "observation_results" } } }'
表架构
AWS Glue 数据质量自动监测功能将每种结果类型写入单独的 Iceberg 表。规则结果、分析结果(包括单独的分布结果表)和观测结果表按 catalog_id、database_name、table_name 和 day(stored_on) 进行分区,以实现高效查询。您可以直接对 stored_on 进行筛选以执行基于时间的查询,Iceberg 会自动处理分区修剪。
规则结果表
规则结果表存储在数据质量运行期间评估的每条规则的通过或失败结果。
| 列 | 类型 | 说明 |
|---|---|---|
dq_result_id |
STRING | 数据质量结果的唯一标识符。 |
rule_name |
STRING | 规则的名称(例如,Rule_1)。 |
rule_description |
STRING | 规则的 DQDL 表达式。 |
rule_result |
STRING | 评估结果:PASS 或 FAIL。 |
evaluation_message |
STRING | 描述失败原因的消息(如果适用)。 |
evaluated_metrics |
MAP<STRING, DOUBLE> | 规则评估的指标。 |
catalog_id |
STRING | 源表的目录 ID。 |
database_name |
STRING | 源表的数据库名称。 |
table_name |
STRING | 源表的名称。 |
ruleset_evaluation_run_id |
STRING | 评估运行的 ID。 |
started_on |
TIMESTAMP | 评估开始的时间。 |
completed_on |
TIMESTAMP | 评估完成的时间。 |
evaluated_rule |
STRING | 经操作数解析后得到的已求值规则表达式。 |
ruleset_name |
STRING | 产生此结果的规则集的名称。 |
数据剖析结果表
下表描述数据剖析结果表中的各列。此表存储分析器和规则(例如 Mean、StandardDeviation 和 Completeness)收集的标量统计数据。AWSGlue 数据质量自动监测功能将分布统计数据存储在单独的分布结果表中。
| 列 | 类型 | 说明 |
|---|---|---|
profile_id |
STRING | 数据质量剖析的唯一标识符。 |
statistic_id |
STRING | 统计数据的唯一标识符。 |
statistic_name |
STRING | 统计数据的名称(例如,Mean、Completeness) |
evaluation_level |
STRING | 评估统计数据的级别:Dataset、Column 或 Multicolumn。 |
statistics_value |
DOUBLE | 统计数据的标量值。 |
statistic_properties |
MAP<STRING, STRING> | 统计数据的其他属性。 |
columns_referenced |
ARRAY<STRING> | 统计数据引用的列。 |
referenced_datasets |
ARRAY<STRING> | 统计数据的参考数据集。 |
column_name |
STRING | 目标列名。 |
dq_result_id |
STRING | 数据质量结果标识符。 |
started_on |
TIMESTAMP | 评估开始的时间。 |
completed_on |
TIMESTAMP | 评估完成的时间。 |
stored_on |
TIMESTAMP | 记录写入表中的时间。 |
catalog_id |
STRING | 源表的目录 ID。 |
database_name |
STRING | 源表的数据库名称。 |
table_name |
STRING | 源表的名称。 |
region |
STRING | AWS 区域。 |
account_id |
STRING | AWS 账户 ID。 |
ruleset_evaluation_run_id |
STRING | 评估运行的 ID。 |
分布结果表
下表介绍了分布结果表中的各列。分布结果与标量分析统计数据分开存储,每个分箱或类别对应一行。您可以在 ProfilingResults.DistributionResults 块中配置此表。
| 列 | 类型 | 说明 |
|---|---|---|
statistic_id |
STRING | 分布统计数据的唯一标识符。 |
column_name |
STRING | 源列(例如,“age”或“department”)。 |
data_type |
STRING | 列的数据类型(例如,“LongType”、“StringType”)。 |
num_bins |
INT | 用于分布的分箱数。 |
bin_index |
INT | 基于 0 的分箱位置。 |
bin_label |
STRING | 对于分类列:不同的值。对于数字列,则为 NULL。 |
bin_lower_bound |
STRING | 对于数字列:分箱的下边缘。对于分类列,则为 NULL。 |
bin_upper_bound |
STRING | 对于数字列:分箱的上边缘。对于分类列,则为 NULL。 |
bin_count |
BIGINT | 此分箱的频率计数。 |
null_count |
INT | 从分布中排除的 NULL 值的数量。在一次运行中,给定统计数据每一行上的值相同。不存在空值时为 NULL。 |
tail_count |
INT | 排名前 20 个之外的类别值的聚合频率。在一次运行中,给定统计数据每一行上的值相同。对于数值直方图,则为 NULL。 |
profile_id |
STRING | 剖析标识符。 |
dq_result_id |
STRING | 数据质量结果标识符。 |
ruleset_evaluation_run_id |
STRING | 评估运行标识符。 |
started_on |
TIMESTAMP | 评估开始的时间。 |
completed_on |
TIMESTAMP | 评估完成的时间。 |
stored_on |
TIMESTAMP | 记录写入表中的时间。 |
catalog_id |
STRING | 源表的目录 ID。 |
database_name |
STRING | 源表的数据库名称。 |
table_name |
STRING | 源表的名称。 |
region |
STRING | AWS 区域。 |
account_id |
STRING | AWS 账户 ID。 |
行级结果表
下表描述行级结果表中的各列。您可以使用此表来标识不符合数据质量规则的特定记录。
| 列 | 类型 | 说明 |
|---|---|---|
| 源列 | 可变。 | 原始源数据中的所有列。 |
data_quality_rules_pass |
ARRAY<STRING> | 此记录通过的规则。 |
data_quality_rules_fail |
ARRAY<STRING> | 此记录未通过的规则。 |
data_quality_rules_skip |
ARRAY<STRING> | 此记录中跳过的规则。 |
data_quality_evaluation_result |
STRING | 此记录的总体评估结果:Passed 或 Failed。 |
dq_result_id |
STRING | 数据质量结果的唯一标识符。 |
ruleset_evaluation_run_id |
STRING | 评估运行的 ID。 |
started_on |
TIMESTAMP | 评估开始的时间。 |
completed_on |
TIMESTAMP | 评估完成的时间。 |
stored_on |
TIMESTAMP | 记录写入表中的时间。 |
catalog_id |
STRING | 源表的目录 ID。 |
database_name |
STRING | 源表的数据库名称。 |
table_name |
STRING | 源表的名称。 |
region |
STRING | AWS 区域。 |
account_id |
STRING | AWS 账户 ID。 |
观测结果表
观测结果表存储每次评估运行中每个统计数据的异常检测预测结果。该表包含所有预测结果:异常值、正常值和跳过的预测。这使您可以呈现带有预测区间的连续趋势图。
| 列 | 类型 | 说明 |
|---|---|---|
statistic_id |
STRING | 正在监控的统计数据的标识符。 |
statistic_name |
STRING | 监控的统计数据的名称。 |
prediction_outcome |
STRING | 异常检测结果:ANOMALY、NOT_ANOMALY 或 SKIPPED。 |
expected_value |
DOUBLE | 预测的预期值。跳过预测时为 NULL。 |
lower_bound |
DOUBLE | 预测范围的下限。跳过预测时为 NULL。 |
upper_bound |
DOUBLE | 预测范围的上限。跳过预测时为 NULL。 |
observation_message |
STRING | 对异常的描述(如果检测到)。 |
training_input |
STRING | 此数据点是否包含在异常检测模型中:INCLUDED 或 EXCLUDED。 |
ruleset_evaluation_run_id |
STRING | 评估运行的 ID。 |
recorded_on |
TIMESTAMP | 记录观测结果的时间。 |
stored_on |
TIMESTAMP | 记录写入表中的时间。 |
actual_value |
DOUBLE | 统计数据的实际观测值。 |
training_status |
STRING | 异常检测模型训练的状态(例如,PENDING、COMPLETED)。 |
recommended_rules |
STRING | 根据异常检测预测推荐的规则。 |
modified_rules |
STRING | 使用基于预测的更新阈值修改的规则。 |
catalog_id |
STRING | 源表的目录 ID。 |
database_name |
STRING | 源表的数据库名称。 |
table_name |
STRING | 源表的名称。 |
注意
观测结果表使用仅追加写入模型。使用 BatchPutDataQualityStatisticAnnotation API 排除数据点时,会追加一个新行,其中 training_input 被设置为 EXCLUDED。要查询每个观测值的最新状态,请使用 stored_on 时间戳来标识每个统计数据与运行组合的最新行。
注意
此表还存储分布溢出观测值,当超过 2% 的值落在冻结分箱边界之外时生成。这些行包含 statistic_name = 'Distribution',且 prediction_outcome 为 NULL。observation_message 字段包含溢出描述。
查询结果
数据质量评估完成后,您可以直接查询结果表。以下示例演示常见的查询模式。
示例:查找特定运行的失败规则
SELECT rule_name, rule_description, evaluation_message, evaluated_metrics FROM quality_results.rule_results WHERE ruleset_evaluation_run_id = 'dqr-12345678' AND rule_result = 'FAIL' ORDER BY rule_name;
示例:查看一段时间内的分析统计数据
SELECT stored_on, statistics_value FROM quality_results.profiles WHERE database_name = 'my_database' AND table_name = 'my_table' AND statistic_name = 'Mean' AND columns_referenced = ARRAY['salary'] ORDER BY stored_on;
示例:识别未通过特定规则的行
SELECT * FROM quality_results.row_level_results WHERE data_quality_evaluation_result = 'Failed' AND contains(data_quality_rules_fail, 'IsComplete "email"');
示例:查看数字直方图
SELECT bin_index, bin_lower_bound, bin_upper_bound, bin_count FROM quality_results.distributions WHERE column_name = 'salary' AND ruleset_evaluation_run_id = 'dqrun-abc123' ORDER BY bin_index;
示例:查看类别值分布
SELECT bin_label, bin_count FROM quality_results.distributions WHERE column_name = 'department' AND ruleset_evaluation_run_id = 'dqrun-abc123' ORDER BY bin_count DESC;
示例:跟踪一段时间内的类别频率
SELECT started_on, bin_count FROM quality_results.distributions WHERE column_name = 'status' AND bin_label = 'active' ORDER BY started_on;
示例:使用预测区间查看异常检测趋势
SELECT o.recorded_on, p.statistics_value AS actual_value, o.expected_value, o.lower_bound, o.upper_bound, o.prediction_outcome FROM quality_results.profiles p JOIN quality_results.observation_results o ON p.statistic_id = o.statistic_id AND p.ruleset_evaluation_run_id = o.ruleset_evaluation_run_id WHERE p.database_name = 'my_database' AND p.table_name = 'my_table' AND p.statistic_name = 'RowCount' AND p.stored_on >= DATE '2025-03-01' ORDER BY p.stored_on;
示例:注释后查询最新观测状态
由于观测结果表使用仅追加模型,排除注释会添加新行。使用重复数据删除查询获取每次观测的最新状态:
SELECT statistic_id, statistic_name, prediction_outcome, expected_value, lower_bound, upper_bound, training_input, stored_on FROM ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY statistic_id, ruleset_evaluation_run_id ORDER BY stored_on DESC ) AS rn FROM quality_results.observation_results WHERE database_name = 'my_database' AND table_name = 'my_table' ) WHERE rn = 1 ORDER BY stored_on;
注意事项
将数据质量结果写入 Data Catalog 表时,请牢记以下注意事项:
-
AWS Glue 数据质量自动监测功能以 Apache Iceberg 格式存储结果,该格式支持高效的时间旅行查询和分区裁剪。
-
单个结果表可以存储来自多个源表的结果。使用
catalog_id、database_name和table_name分区列筛选特定源的结果。 -
评估运行完成后,AWS Glue 数据质量自动监测功能会异步写入观测结果。观测结果出现在表中之前,可能会有短暂的延迟。
-
对于分布结果表中的分布统计数据,每个分箱或类别都存储为单独的行。例如,包含 20 个分箱的直方图会在表中为该统计数据生成 20 行。