如何在单个DAG中添加两个SQL查询,实现sql_query1条件触发sql_query2
现有DAG配置片段
export_query: sql_query1 ; export_query_filter: "" processing_modules: - 'export_bqtobqtmp_task' - 'export_bqtmptogcs_task' - 'pre_shell_execution' - 'post_shell_execution' processing_dependencies: - 'start >> generate_hashid >> export_bqtobqtmp_task >> export_bqtmptogcs_task >> pre_shell_execution >> post_shell_execution >> bq_delete_target_export_table >> end'
需求
先执行sql_query1,若其结果符合预期,则继续执行sql_query2,需在单个DAG中实现该逻辑。
实现方案
要实现这种基于SQL查询结果的分支判断,可按以下步骤调整DAG配置:
新增校验与执行任务
- 添加
validate_sql1_result任务:负责执行sql_query1后校验结果是否符合预期(比如结果行数达标、特定字段值满足规则等),校验通过则任务成功,不通过则标记失败。 - 添加
execute_sql2_task任务:专门用来执行sql_query2。
- 添加
更新processing_modules列表
将新增任务加入模块列表:processing_modules: - 'export_bqtobqtmp_task' - 'export_bqtmptogcs_task' - 'pre_shell_execution' - 'post_shell_execution' - 'validate_sql1_result' - 'execute_sql2_task'重构processing_dependencies依赖链
调整依赖关系,实现分支逻辑:processing_dependencies: - 'start >> generate_hashid >> export_bqtobqtmp_task >> validate_sql1_result' - 'validate_sql1_result >> execute_sql2_task >> export_bqtmptogcs_task >> pre_shell_execution >> post_shell_execution >> bq_delete_target_export_table >> end' - 'validate_sql1_result >> end'说明:
export_bqtobqtmp_task执行完sql_query1后,先跑校验任务;校验成功则执行sql_query2,再走原有后续流程;校验失败则直接结束DAG。任务逻辑细节
validate_sql1_result可结合BigQueryOperator执行查询,再用PythonOperator判断结果;也可以直接用SQL语句返回布尔值,通过任务返回码判定是否通过。execute_sql2_task用对应的数据执行算子(比如BigQueryOperator)来运行sql_query2即可。
内容的提问来源于stack exchange,提问作者Vishakha Waghulde
相关产品推荐
相关产品推荐

