如何基于python_callable结果为TriggerDagRunOperator的trigger_dag_id赋值?
问题
编写了如下代码,希望根据workflow_name触发不同的DAG,现咨询:TriggerDagRunOperator的trigger_dag_id参数应填入什么值?是留空字符串还是其他值?
代码示例
def trigger_check(context, dag_run_obj): workflow_name_var = context['ti'].xcom_pull(key='workflow_name') dag_run_obj.payload = { "workflow_name": workflow_name_var } if workflow_name_var == 'normal': dag_run_obj.trigger_dag_id = JOB_NAME elif workflow_name_var == 'reload': dag_run_obj.trigger_dag_id = JOB_NAME_RELOAD elif workflow_name_var == 'XXX': dag_run_obj.trigger_dag_id = XXXX return dag_run_obj trigger_check_task = TriggerDagRunOperator( task_id='trigger_check_task', python_callable= trigger_check, trigger_dag_id="", dag=dag )
回答
你可以把trigger_dag_id设为空字符串"",或者填一个任意占位值(比如"placeholder_dag")都可以——因为你在python_callable函数里已经通过dag_run_obj.trigger_dag_id动态指定了要触发的目标DAG ID,Operator里的这个初始参数值会被函数中的赋值覆盖。
需要注意的是:
- 要保证
workflow_name_var的所有可能取值都被if/elif分支覆盖,避免出现未设置trigger_dag_id的情况。如果有遗漏分支,系统会使用你在Operator里填的初始值,若初始值为空字符串可能触发报错。 - 保险起见,可以在分支最后加一个
else块,设置一个默认的 fallback DAG ID,或者直接抛出异常,避免因未匹配到分支导致的意外问题。
内容的提问来源于stack exchange,提问作者Eric
相关产品推荐
相关产品推荐

