相关疑难解决方法(0)

在 TriggerDagRunOperator 中提供上下文

我有一个由另一个 dag 触发的 dag。我已经通过DagRunOrder().payload字典以与官方示例相同的方式将一些配置变量传递给了这个 dag 。

现在在这个 dag 中,我有另一个 dagTriggerDagRunOperator来启动第二个 dag,并希望通过这些相同的配置变量。

我已经成功地访问了有效载荷变量,PythonOperator如下所示:

def run_this_func(ds, **kwargs):
    print("Remotely received value of {} for message and {} for day".format(
        kwargs["dag_run"].conf["message"], kwargs["dag_run"].conf["day"])
    )

run_this = PythonOperator(
    task_id='run_this',
    provide_context=True,
    python_callable=run_this_func,
    dag=dag
)
Run Code Online (Sandbox Code Playgroud)

但是相同的模式在以下情况下不起作用TriggerDagRunOperator:

def trigger(context, dag_run_obj, **kwargs):
    dag_run_obj.payload = {
        "message": kwargs["dag_run"].conf["message"],
        "day": kwargs["dag_run"].conf["day"]
    }
    return dag_run_obj

trigger_step = TriggerDagRunOperator(
    task_id="trigger_modelling",
    trigger_dag_id="Dummy_Modelling",
    provide_context=True,
    python_callable=trigger,
    dag=dag
)
Run Code Online (Sandbox Code Playgroud)

它会产生关于使用的警告provide_context:

INFO - Subtask: /usr/local/lib/python2.7/dist-packages/airflow/models.py:1927: PendingDeprecationWarning: Invalid …
Run Code Online (Sandbox Code Playgroud)

python airflow

3
推荐指数
2
解决办法
1万
查看次数

标签 统计

airflow ×1

python ×1