【发布时间】:2021-06-04 12:53:12
【问题描述】:
我正在尝试弄清楚如何实现工作流,以便传感器任务等待外部 dag 完成,只等待一定天数。这是一项日常工作,所以我希望传感器工作等待 3 天,然后在第四天发送一封电子邮件,然后等待或执行其他任务。
有人可以帮助阐明如何实现这一目标吗?另外,我们如何将“天数计数器”从一天传递到另一天?非常感谢您的帮助。
【问题讨论】:
标签: python airflow-scheduler airflow
我正在尝试弄清楚如何实现工作流,以便传感器任务等待外部 dag 完成,只等待一定天数。这是一项日常工作,所以我希望传感器工作等待 3 天,然后在第四天发送一封电子邮件,然后等待或执行其他任务。
有人可以帮助阐明如何实现这一目标吗?另外,我们如何将“天数计数器”从一天传递到另一天?非常感谢您的帮助。
【问题讨论】:
标签: python airflow-scheduler airflow
您可以将ExternalTaskSensor 与以下配置一起使用:
timeout = 3 * 24 * 60 * 60 - 3 天(以秒为单位),之后传感器将失效poke_interval = 12 * 60 * 60 - 传感器检查间隔 12 小时,您可以将其调整为每小时检查一次。当您检查外部 dag 状态时,它将减少次数mode = "reschedule" - 这样传感器不会占用 3 天的工作槽,它将被调度,执行,如果不满足条件,它将被重新调度到下一个 poke_interval 秒执行。将此模式用于长时间运行的任务是一种很好的做法。此外,您可以将等待的 DAG 构建为 wait_task >> [success_task , fail_task] where
wait_task 是您的传感器success_task 有触发规则all_success 并在传感器成功时遵循fail_task 和 all_failed 触发规则并处理传感器最终返回 false 或超时时的场景【讨论】:
execution_date_fn,可用于感知远程 dag 的某些执行。例如,将一个函数附加到此参数,该函数返回远程 dag 的三个执行日期列表。
time out
float('inf')。但在这种情况下,请记住使用reschedule 模式。
float('inf'),这样它就永远不会超时,如果您不设置它,则将使用默认值 7 天。对于运行时间超过一分钟的内容,始终建议使用reschedule 模式。