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

如何在单个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配置:

  1. 新增校验与执行任务

    • 添加validate_sql1_result任务:负责执行sql_query1后校验结果是否符合预期(比如结果行数达标、特定字段值满足规则等),校验通过则任务成功,不通过则标记失败。
    • 添加execute_sql2_task任务:专门用来执行sql_query2。
  2. 更新processing_modules列表
    将新增任务加入模块列表:

    processing_modules:
    - 'export_bqtobqtmp_task'
    - 'export_bqtmptogcs_task'
    - 'pre_shell_execution'
    - 'post_shell_execution'
    - 'validate_sql1_result'
    - 'execute_sql2_task'
    
  3. 重构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。

  4. 任务逻辑细节

    • validate_sql1_result可结合BigQueryOperator执行查询,再用PythonOperator判断结果;也可以直接用SQL语句返回布尔值,通过任务返回码判定是否通过。
    • execute_sql2_task用对应的数据执行算子(比如BigQueryOperator)来运行sql_query2即可。

内容的提问来源于stack exchange,提问作者Vishakha Waghulde

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 04:25:01