【问题标题】:Airflow External sensor gets stuck at poking气流外部传感器卡在戳
【发布时间】:2018-10-13 19:40:16
【问题描述】:

我希望一个 dag 在另一个 dag 完成后开始。一种解决方案是使用外部传感器功能,您可以在下面找到我的解决方案。我遇到的问题是依赖的 dag 卡在戳,我检查了这个answer 并确保两个 dag 都按相同的时间表运行,我的简化代码如下: 任何帮助,将不胜感激。 领导者:

from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta


default_args = {
   'owner': 'airflow',
   'depends_on_past': False,
   'start_date': datetime(2015, 6, 1),
   'retries': 1,
   'retry_delay': timedelta(minutes=5),



 }

 schedule = '* * * * *'

 dag = DAG('leader_dag', default_args=default_args,catchup=False, 
 schedule_interval=schedule)

t1 = BashOperator(
   task_id='print_date',
   bash_command='date',
   dag=dag)

依赖的 dag:

from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta
from airflow.operators.sensors import ExternalTaskSensor


default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2018, 10, 8),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),


}
schedule='* * * * *'
dag = DAG('dependent_dag', default_args=default_args, catchup=False, 
schedule_interval=schedule)

 wait_for_task = ExternalTaskSensor(task_id = 'wait_for_task', 
 external_dag_id = 'leader_dag', external_task_id='t1', dag=dag)

 t1 = BashOperator(
     task_id='print_date',
     bash_command='date',
      dag=dag)

 t1.set_upstream(wait_for_task)

leader_dag 的日志:

依赖 dag 的日志:

【问题讨论】:

  • 感谢您在问题中提供的所有详细信息,尽管将日志作为文本提供会使搜索更加友好。我认为传感器存在根本性的误用,在我的答案末尾显示了修复,但充其量,我不会在 Airflow 中安排比*/10 * * * * 更频繁的任何事情。

标签: airflow airflow-scheduler


【解决方案1】:

首先leader_dag 中的task_id 被命名为print_date,但您使用任务wait_for_task 设置您的dependent_dag,该任务正在等待leader_dag 的名为t1 的任务。没有名为 t1 的任务。您在 py 文件中分配给它的内容不相关,也没有在 Airflow db 中使用,也没有被传感器横向使用。它应该正在等待任务名称print_date

第二个你的日志没有排队你在哪个 leader_dag 运行中显示了dependent_dag 正在等待什么。

最后,我不建议您使用 Airflow 来安排每分钟的任务。当然不是两个依赖的任务在一起。 考虑在 Spark 等不同的系统中编写流式作业,或为此滚动您自己的 Celery 或 Dask 环境。

您还可以通过在leader_dag 的末尾添加TriggerDagRunOperator 来触发dependent_dag 来避免ExternalTaskSensor,并通过将schedule_interval 设置为None 来删除计划。

我在您的日志中看到的是 2018-10-13T19:08:11 的领导者日志。这充其量是 execution_date 2018-10-13 19:07:00 的 dagrun,因为从 19:07 开始的分钟时间段在 19:08 结束,这是可以安排的最早时间。我看到调度和执行之间有大约 11 秒的延迟如果是这种情况。但是,Airflow 中可能存在数分钟的调度延迟。

我还看到了来自dependent_dag 的日志,该日志从 19:14:04 运行到 19:14:34,正在寻找对应的 19:13:00 dagrun 的完成情况。没有迹象表明您的调度程序没有足够的延迟,可以在 19:14:34 之前启动 leader_dag 的 19:13:00 dagrun。如果你展示它戳了 5 分钟左右,你会更好地说服我。当然,它永远不会感知到 leader_dag.t1,因为这不是您所显示的任务的名称。

所以,Airflow 有调度延迟,如果系统中有几个 1000 个 dag,它可能会高于 1 分钟,这样使用 catchup=False 您将在 IE 19 之后运行一些运行: 08、19:09 和一些跳过一分钟(或 6 分钟)的运行(例如 19:10 之后是 19:16)可能会发生,并且由于延迟在逐个日期的基础上有点随机,您可能会得到未对齐的运行传感器一直在等待,即使您有正确的任务 ID 等待:

 wait_for_task = ExternalTaskSensor(
     task_id='wait_for_task', 
     external_dag_id='leader_dag',
-    external_task_id='t1',
+    external_task_id='print_date',
     dag=dag)

【讨论】:

  • 感谢@dlamblin,我进行了与 task_id 相关的更改,并且成功了。我确实为测试目的安排了每分钟,但我没有想到调度程序的延迟。我不知道 TriggerDagOperator,我会检查一下。再次感谢您的详尽回答。
  • @sia 很高兴为您提供帮助;我认为它在技术上称为:TriggerDagRunOperator
  • @dlamblin 即使进行了您建议的更改,我也遇到了类似的问题。你能看到我在做什么错here吗?
  • 嗨@dlamblin,你的回答很棒!但是如果我将None 设置为调度间隔呢?因为我对自己的 DAG 的这种方法也有同样的问题,你能给我一个建议吗?
  • @ImamDigmi ExternalTask​​Sensor 可以采用确切的执行日期或返回您正在感应的任务的执行日期的可调用函数。无论执行是计划的执行还是手动触发的执行,这都应该有效。您可以使用可调用对象在数据库中查找目标 DAG DagRun 的最新执行日期,如果这是您想要的 None 间隔依赖 DAG。如果不以其中一种方式定义日期,它将默认为当前 DAG 的执行日期,在这种情况下,这对您来说几乎肯定是错误的。
【解决方案2】:

使用ExternalTaskSensor 时,您必须为两个 DAG 指定相同的开始日期。如果这不适用于您的用例,那么您需要在您的ExternalTaskSensor 中使用execution_deltaexecution_date_fn

【讨论】:

  • 我更改了开始日期,但问题仍然存在,如果两者的时间表相同,我会在文档中阅读,我不需要使用execution_delta
【解决方案3】:

简单的解决方案是执行以下操作:

1- 确保您在 ExternalTask​​Sensor 中设置的所有变量都设置正确,与您希望 master_dag 关注的 dag 的 task_id 和 dag_id 完全相同。

2- 使 master_dag 和 slave_dag(你要等待它的 dag)具有相同的 start_date,否则它将无法工作。如果您的奴隶从下午 22 点开始,而主人从 22:30 开始,那么您应该使用执行增量指定 30 分钟的差异。

如果您的错误无法通过以下方式解决,那么您的问题要么是基本问题,要么是您编写的 dag 方法太错误了..

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多