【发布时间】:2018-07-04 05:39:57
【问题描述】:
最初使用dag callback(on_failure_callback 和on_success_callback)时,我认为它会在dag 完成时触发success 或fail 状态(因为它在dag 中定义)。
但是它似乎是在每个task instance 而不是dag run 上实例化的,所以如果一个 DAG 有 N 个任务,它将触发这些回调 N 次。
我正在尝试捕获任务 ID,因此发送到 slack。阅读另一个related question 我想出了以下内容:
def success_msg(context):
slack.slack_message(context['task_instance']); #send task-id to slack
def failure_msg(context):
slack.slack_message(context['task_instance']); #send task-id to slack
default_args = {
[...]
'on_failure_callback': failure_msg,
'on_success_callback': success_msg,
[...]
}
但它失败了,我应该如何解析上下文变量以便获得任务ID?
【问题讨论】:
-
如此处所述stackoverflow.com/questions/51147762/… 回调仅适用于 1.9.0 中的任务级别,我认为在 1.10.0 中,DAG 也可以是回调级别。
-
是的,它在 DAG 级别调用回调函数,但我无法在该回调函数中访问上下文我也在做上面写的同样的事情\有什么帮助吗?
标签: airflow