【问题标题】:Airflow - Parse task-id from dag context callbackAirflow - 从 dag 上下文回调中解析任务 ID
【发布时间】:2018-07-04 05:39:57
【问题描述】:

最初使用dag callbackon_failure_callbackon_success_callback)时,我认为它会在dag 完成时触发successfail 状态(因为它在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


【解决方案1】:

您可以从上下文中使用任务对象访问任务。

context['task'] 应该是执行此操作的适当方式。要获取任务名称,请使用task_id

context['task'].task_id

要查找上下文中可用的更多对象,您可以在此处浏览列表:https://airflow.apache.org/docs/apache-airflow/stable/macros-ref.html

【讨论】:

  • 确实有效。感谢您的回答和文档参考。
  • 你如何调试它的项目?我很好奇从 dag 和任务实例中获取项目。我会打印它们,但是当我像 python dag.py 一样运行时,它不会记录...
  • 请为此打开一个新的堆栈溢出问题。还包括示例,因为我不确定您在问什么。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-06-19
  • 1970-01-01
  • 2022-06-14
  • 1970-01-01
  • 1970-01-01
  • 2023-02-03
相关资源
最近更新 更多