如何读取Dataplex数据质量结果表的查询并执行存储至Looker可查表格?
解决方案:Dataplex数据质量查询自动化与Looker集成
一、整体流程
核心逻辑:从Dataplex数据质量结果表提取校验失败数据的查询语句 → 批量执行查询 → 将结果写入BigQuery(Looker原生支持的存储)→ 在Looker中可视化分析
二、具体实现步骤
1. 读取Dataplex中的DQ查询语句
Dataplex的数据质量检测结果默认存储在BigQuery中,直接通过BigQuery客户端读取包含查询语句的字段(比如示例中的row_filter_query):
from google.cloud import bigquery def fetch_dq_failure_queries(): bq_client = bigquery.Client() # 替换为你的Dataplex DQ结果表路径 dq_results_table = "your-gcp-project.your-dataset.dataplex_dq_validation_results" # 过滤出有效查询(排除空值) fetch_query = f""" SELECT rule_id, row_filter_query, target_table FROM `{dq_results_table}` WHERE row_filter_query IS NOT NULL AND row_filter_query != '' """ query_job = bq_client.query(fetch_query) return [dict(row) for row in query_job.result()]
2. 执行查询并写入Looker可访问的表
将每个DQ查询的结果写入统一的BigQuery结果表,同时附加规则ID、源表、执行时间等元数据,方便Looker筛选分析:
def execute_queries_and_store_results(dq_queries): bq_client = bigquery.Client() # 替换为你的目标数据集(需配置Looker权限) target_dataset = "your-gcp-project.looker_dq_failure_data" bq_client.create_dataset(target_dataset, exists_ok=True) for query_item in dq_queries: rule_id = query_item["rule_id"] source_table = query_item["target_table"] failure_query = query_item["row_filter_query"] # 增强查询,添加元数据字段 enhanced_query = f""" SELECT '{rule_id}' AS dq_rule_id, '{source_table}' AS source_table_name, CURRENT_TIMESTAMP() AS execution_timestamp, * FROM ({failure_query}) """ # 写入按天分区的结果表,提升查询性能 dest_table = f"{target_dataset}.dq_failed_rows$DATE(CURRENT_TIMESTAMP())" job_config = bigquery.QueryJobConfig( destination=dest_table, write_disposition=bigquery.WriteDisposition.WRITE_APPEND, time_partitioning=bigquery.TimePartitioning(type_=bigquery.TimePartitioningType.DAY) ) bq_client.query(enhanced_query, job_config=job_config).result()
3. 配置Looker访问
- 在Looker中添加BigQuery数据源连接,指向存储DQ失败数据的数据集
- 创建Looker视图/探索,基于
dq_failed_rows表,可按规则ID、源表、执行时间等维度筛选,快速定位不符合质量规则的字段和数据
三、Cloud Function复用说明:无需为每张表单独创建
只需要一个Cloud Function即可处理所有表的DQ查询,理由:
- 从Dataplex结果表读取的查询已包含目标表信息,代码可动态适配任意表
- 可通过Cloud Scheduler定时触发该Function(如每日凌晨执行),实现自动化
- 若需针对特定表做特殊逻辑,只需在代码中添加条件判断,无需新建Function
四、优化建议
- 错误处理:捕获查询执行异常,将失败记录写入单独的错误日志表
- 增量执行:记录已处理的规则ID,避免重复执行相同查询
- 资源管控:设置BigQuery查询优先级为
BATCH,避免占用生产资源
内容的提问来源于stack exchange,提问作者Bastián SN
相关产品推荐
相关产品推荐

