【问题标题】:Airflow task callbacks are missed sometimes有时会错过 Airflow 任务回调
【发布时间】:2023-02-19 16:53:44
【问题描述】:

气流:2.1.2 - 执行器:KubernetesExecutor - 蟒蛇:3.7

我使用 Airflow 2+ TaskFlow API 编写了任务,并在 KubernetesExecutor 模式下运行 Airflow 应用程序。任务有成功和失败的回调,但有时它们会被遗漏。

我尝试通过 DAG 上的 default_args 和直接在任务装饰器中指定回调,但看到相同的行为。

@task(
    on_success_callback=common.on_success_callback,
    on_failure_callback=common.on_failure_callback,
)
def delta_load_pstn(files):
    # doing something here

这是任务的结束日志

2022-04-26 11:21:38,494] Marking task as SUCCESS. dag_id=delta_load_pstn, task_id=dq_process, execution_date=20220426T112104, start_date=20220426T112131, end_date=20220426T112138
[2022-04-26 11:21:38,548] 1 downstream tasks scheduled from follow-on schedule check
[2022-04-26 11:21:42,069] State of this instance has been externally set to success. Terminating instance.
[2022-04-26 11:21:42,070] Sending Signals.SIGTERM to GPID 34
[2022-04-26 11:22:42,081] process psutil.Process(pid=34, name='airflow task runner: delta_load_pstn dq_process 2022-04-26T11:21:04.747263+00:00 500', status='sleeping', started='11:21:31') did not respond to SIGTERM. Trying SIGKILL
[2022-04-26 11:22:42,095] Process psutil.Process(pid=34, name='airflow task runner: delta_load_pstn dq_process 2022-04-26T11:21:04.747263+00:00 500', status='terminated', exitcode=<Negsignal.SIGKILL: -9>, started='11:21:31') (34) terminated with exit code Negsignal.SIGKILL
[2022-04-26 11:22:42,095] Job 500 was killed before it finished (likely due to running out of memory)

我可以在任务实例详细信息中看到回调已配置。

如果我实施在执行任务之前调用的 on_execute_callback,我会收到警报(在 Slack 中)。所以我猜这肯定是在处理回调之前杀死 pod 的原因。

【问题讨论】:

    标签: airflow airflow-2.x


    【解决方案1】:

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-03-14
      • 2020-06-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多