【问题标题】:Airflow : ExternalTaskSensor doesn't trigger the task气流:ExternalTask​​Sensor 不会触发任务
【发布时间】:2019-06-05 06:56:45
【问题描述】:

我已经看到了关于 SO 的 thisthis 问题,并进行了相应的更改。但是,我的依赖 DAG 仍然卡在戳状态。下面是我的主 DAG:

from airflow import DAG
from airflow.operators.jdbc_operator import JdbcOperator
from datetime import datetime
from airflow.operators.bash_operator import BashOperator

today = datetime.today()

default_args = {
    'depends_on_past': False,
    'retries': 0,
    'start_date': datetime(today.year, today.month, today.day),
    'schedule_interval': '@once'
}

dag = DAG('call-procedure-and-bash', default_args=default_args)

call_procedure = JdbcOperator(
    task_id='call_procedure',
    jdbc_conn_id='airflow_db2',
    sql='CALL AIRFLOW.TEST_INSERT (20)',
    dag=dag
)

call_procedure

下面是我的依赖 DAG:

from airflow import DAG
from airflow.operators.jdbc_operator import JdbcOperator
from datetime import datetime, timedelta
from airflow.sensors.external_task_sensor import ExternalTaskSensor

today = datetime.today()

default_args = {
    'depends_on_past': False,
    'retries': 0,
    'start_date': datetime(today.year, today.month, today.day),
    'schedule_interval': '@once'
}

dag = DAG('external-dag-upstream', default_args=default_args)

task_sensor = ExternalTaskSensor(
    task_id='link_upstream',
    external_dag_id='call-procedure-and-bash',
    external_task_id='call_procedure',
    execution_delta=timedelta(minutes=-2),
    dag=dag
)

count_rows = JdbcOperator(
    task_id='count_rows',
    jdbc_conn_id='airflow_db2',
    sql='SELECT COUNT(*) FROM AIRFLOW.TEST',
    dag=dag
)

count_rows.set_upstream(task_sensor)

以下是主 DAG 执行后依赖 DAG 的日志:

[2019-01-10 11:43:52,951] {{external_task_sensor.py:91}} INFO - Poking for call-procedure-and-bash.call_procedure on 2019-01-10T11:45:47.893735+00:00 ... 
[2019-01-10 11:44:52,955] {{external_task_sensor.py:91}} INFO - Poking for call-procedure-and-bash.call_procedure on 2019-01-10T11:45:47.893735+00:00 ... 
[2019-01-10 11:45:52,961] {{external_task_sensor.py:91}} INFO - Poking for call-procedure-and-bash.call_procedure on 2019-01-10T11:45:47.893735+00:00 ... 
[2019-01-10 11:46:52,949] {{external_task_sensor.py:91}} INFO - Poking for call-procedure-and-bash.call_procedure on 2019-01-10T11:45:47.893735+00:00 ... 
[2019-01-10 11:47:52,928] {{external_task_sensor.py:91}} INFO - Poking for call-procedure-and-bash.call_procedure on 2019-01-10T11:45:47.893735+00:00 ... 
[2019-01-10 11:48:52,928] {{external_task_sensor.py:91}} INFO - Poking for call-procedure-and-bash.call_procedure on 2019-01-10T11:45:47.893735+00:00 ... 
[2019-01-10 11:49:52,905] {{external_task_sensor.py:91}} INFO - Poking for call-procedure-and-bash.call_procedure on 2019-01-10T11:45:47.893735+00:00 ... 

以下是master DAG执行的日志:

[2019-01-10 11:45:20,215] {{jdbc_operator.py:56}} INFO - Executing: CALL AIRFLOW.TEST_INSERT (20)
[2019-01-10 11:45:21,477] {{logging_mixin.py:95}} INFO - [2019-01-10 11:45:21,476] {{dbapi_hook.py:166}} INFO - CALL AIRFLOW.TEST_INSERT (20)
[2019-01-10 11:45:24,139] {{logging_mixin.py:95}} INFO - [2019-01-10 11:45:24,137] {{jobs.py:2627}} INFO - Task exited with return code 0

我的假设是,如果 master 运行良好,Airflow 应该触发依赖 DAG?我尝试过使用execution_delta,但这似乎不起作用。

另外,schedule_intervalstart_date 对于两个 DAG 都是相同的,所以不要认为这会造成任何问题。

我错过了什么吗?

【问题讨论】:

  • 你能找出原因吗?我现在也有类似的问题

标签: python airflow directed-acyclic-graphs airflow-scheduler


【解决方案1】:

您可能应该使用正时间增量:https://airflow.readthedocs.io/en/stable/_modules/airflow/sensors/external_task_sensor.html,因为当减去执行增量时,它最终会寻找在其自身运行 2 分钟后运行的任务。

但是增量并不是一个真正的范围,TI 必须在日期时间列表中具有匹配的 Dag ID、任务 ID、成功结果以及执行日期。当您将execution_delta 作为增量提供时,它是一个包含当前执行日期并减去 timedelta 的日期时间列表。

这可能归结为您要么删除 timedelta 以便两个执行日期匹配并且传感器将等到另一个任务成功,或者您的开始日期和计划间隔设置为基本上今天和@once执行日期彼此之间没有可预测的锁步。您可以尝试设置 datetime(2019,1,10)0 1 * * * 让它们在每天凌晨 1 点运行(同样没有 execution_delta)。

【讨论】:

  • 我删除了execution_delta 并将schedule_interval 设置为0 1 * * *。尽管如此,当上游完成时,它并没有触发 DAG。但是,当我动态更改开始日期(传感器正在执行时)时,它会以某种方式完成下游DAG。看起来这可能与DAGs 的开始日期有关,但我还无法弄清楚。
【解决方案2】:

确保两个 DAG 同时启动,并且不要手动启动任何一个 DAG。

【讨论】:

    【解决方案3】:

    希望您没有手动触发 DAG。如果你想测试它,让 DAG 按照计划运行,然后监控 DAG 运行。

    【讨论】:

    • 是否有任何手动运行外部任务感应的解决方案?
    【解决方案4】:

    由于夏季/冬季时间更改,我遇到了这个问题:“1 天前”表示“正好 24 小时前”,因此如果时区之间有夏令时更改,则 DAG 会卡住。

    解决此问题的一种方法是手动将其设置为成功。

    在这种情况下,另一种方法是使用execution_date_fn 参数并手动正确计算时间差。

    【讨论】:

      猜你喜欢
      • 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
      相关资源
      最近更新 更多