You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何读取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访问

  1. 在Looker中添加BigQuery数据源连接,指向存储DQ失败数据的数据集
  2. 创建Looker视图/探索,基于dq_failed_rows表,可按规则ID、源表、执行时间等维度筛选,快速定位不符合质量规则的字段和数据

三、Cloud Function复用说明:无需为每张表单独创建

只需要一个Cloud Function即可处理所有表的DQ查询,理由:

  • 从Dataplex结果表读取的查询已包含目标表信息,代码可动态适配任意表
  • 可通过Cloud Scheduler定时触发该Function(如每日凌晨执行),实现自动化
  • 若需针对特定表做特殊逻辑,只需在代码中添加条件判断,无需新建Function

四、优化建议

  • 错误处理:捕获查询执行异常,将失败记录写入单独的错误日志表
  • 增量执行:记录已处理的规则ID,避免重复执行相同查询
  • 资源管控:设置BigQuery查询优先级为BATCH,避免占用生产资源

内容的提问来源于stack exchange,提问作者Bastián SN

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 02:35:58