如何通过GCP Cloud Composer执行含UniqueKey断言的Dataform数据质量任务?
问题描述
我需要借助Cloud Composer的DAG,使用Dataform的断言对BigQuery中的UniqueKey执行数据质量校验任务。目前编写的DAG可以执行Dataform工作流,但无法完成数据质量校验,相关代码如下:
from datetime import datetime from google.cloud.dataform_v1beta1 import WorkflowInvocation from airflow import models from airflow.models.baseoperator import chain from airflow.providers.google.cloud.operators.dataform import ( DataformCancelWorkflowInvocationOperator, DataformCreateCompilationResultOperator, DataformCreateWorkflowInvocationOperator, DataformGetCompilationResultOperator, DataformGetWorkflowInvocationOperator, ) DAG_ID = "dataform" PROJECT_ID = "PROJECT_ID" REPOSITORY_ID = "REPOSITORY_ID" REGION = "REGION" GIT_COMMITISH = "GIT_COMMITISH" with models.DAG( DAG_ID, schedule_interval='@once', # 根据需求修改 start_date=datetime(2022, 1, 1), catchup=False, # 根据需求修改 tags=['dataform'], ) as dag: create_compilation_result = DataformCreateCompilationResultOperator( task_id="create_compilation_result", project_id=PROJECT_ID, region=REGION, repository_id=REPOSITORY_ID, compilation_result={ "git_commitish": GIT_COMMITISH, }, ) create_workflow_invocation = DataformCreateWorkflowInvocationOperator( task_id='create_workflow_invocation', project_id=PROJECT_ID, region=REGION, repository_id=REPOSITORY_ID, workflow_invocation={ "compilation_result": "{{ task_instance.xcom_pull('create_compilation_result')['name'] }}" }, ) create_compilation_result >> create_workflow_invocation
问题排查与解决步骤
1. 确认Dataform仓库已定义UniqueKey断言
首先要保证Dataform项目中存在针对UniqueKey的校验断言,示例代码如下:
-- Dataform断言文件示例:检查目标表的unique_key列是否唯一 config { type: "assertion", tags: ["data_quality"] -- 给断言打标签,方便后续DAG指定执行 } SELECT unique_key, COUNT(*) AS record_count FROM `${project_id}.${dataset_id}.target_table` GROUP BY unique_key HAVING COUNT(*) > 1
如果Dataform仓库中没有这类断言,DAG执行工作流时自然不会触发数据质量校验。
2. 修改DAG,指定工作流执行断言任务
当前DAG仅触发默认工作流,可能未包含断言任务。需要在create_workflow_invocation的配置中添加invocation_config,明确指定要执行的断言任务或标签:
create_workflow_invocation = DataformCreateWorkflowInvocationOperator( task_id='create_workflow_invocation', project_id=PROJECT_ID, region=REGION, repository_id=REPOSITORY_ID, workflow_invocation={ "compilation_result": "{{ task_instance.xcom_pull('create_compilation_result')['name'] }}", "invocation_config": { # 方式1:通过标签筛选要执行的断言任务 "included_tags": ["data_quality"], # 方式2:直接指定断言任务名称 # "included_tasks": ["assert_unique_key_validation"] } }, )
3. 验证编译结果是否包含断言任务
可以添加DataformGetCompilationResultOperator查看编译后的工作流内容,确认断言任务已被正确编译:
get_compilation_result = DataformGetCompilationResultOperator( task_id="get_compilation_result", project_id=PROJECT_ID, region=REGION, repository_id=REPOSITORY_ID, compilation_result_id="{{ task_instance.xcom_pull('create_compilation_result')['name'].split('/')[-1] }}" ) # 更新任务依赖链 create_compilation_result >> get_compilation_result >> create_workflow_invocation
查看该任务的日志,可确认编译结果中是否包含目标断言任务。
4. 配置断言失败的处理策略
确保断言失败时,工作流会终止并触发Airflow任务失败:
- 在Dataform的断言配置中添加
fail_on_assertion_failure: true(默认通常为true) - 可在Airflow operator中配置重试策略,确保数据质量问题被及时识别
内容的提问来源于stack exchange,提问作者dsrebechi
相关产品推荐
相关产品推荐

