【问题标题】:Cloud Composer Airflow task fails but functions complete successfullyCloud Composer Airflow 任务失败但功能成功完成
【发布时间】:2019-07-21 18:02:53
【问题描述】:

我正在尝试使用 Cloud Composer 编写我的第一个 Airflow 作业。我的 DAG 有三个任务,第一个任务成功完成,但第二个任务似乎失败并发出任何失败错误消息。我在第二个任务中使用PythonOperator。被调用的函数执行一个长时间运行的查询并进行轮询,直到查询完成。查询完成后,我收到一条消息,说数据已输出到正确的表,但 Airflow 将任务视为失败并再次重试任务。

我的 DAG default_args 如下所示:

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': today.strftime("%Y-%m-%d"),
    'email': ['email@email.com'],
    'email_on_failure': True,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=1),
    'dagrun_timeout': timedelta(minutes=30)
}

编辑:

这是我的 Python 可调用对象和 PythonOperator。 run_query 可调用对象在 Stackdriver 日志中生成输出,并指示实际功能已完成,但任务失败。

def run_query(**kwargs):
    ti = kwargs['ti']
    creds = ti.xcom_pull(key='key value 1', task_ids=t1_id)
    service = adh.get_service(creds)

    return adh.start_saved_query(service,
                                 kwargs['customer_id'],
                                 kwargs['query_name'],
                                 kwargs['start_date'],
                                 kwargs['end_date'],
                                 kwargs['project'],
                                 kwargs['dataset'],
                                 kwargs['table'],
                                 parameters=kwargs['parameters'])
run_adh_query = PythonOperator(
    task_id="task2",
    provide_context=True,
    python_callable=run_query,
    dag=dag,
    trigger_rule='all_success',
    op_kwargs={
        'customer_id': 01234,
        'query_name': 'queryName',
        'start_date': start_date.strftime("%Y-%m-%d"),
        'end_date': end_date.strftime("%Y-%m-%d"),
        'project': adh_project,
        'dataset': adh_dataset,
        'table': adh_table,
        'parameters': {
        'CONV_START_DATE': {'value': conv_start_date.strftime("%Y-%m-%d")},
        'CONV_END_DATE': {'value': end_date.strftime("%Y-%m-%d")},
        'LOOKBACK_DAYS': {'value': str(lookback_days)}
        }
    }
)

如果有任何提示,我将不胜感激!

【问题讨论】:

  • 消息来自日志?或电子邮件?也许任务在发送的消息和代码结束之间失败。你应该粘贴你的 PythonOperator 代码。
  • @howie 我只在 Airflow UI 中看到作业失败。 Stackdriver 日志仅在因失败重试开始时显示“正在尝试 2 或 2”。我添加了我的 python 可调用和 python 操作符代码。谢谢!

标签: python google-cloud-platform airflow google-cloud-composer


【解决方案1】:

我没有在您的代码中看到任何错误处理。

当长时间运行的查询和轮询失败时,引发 AirflowException,这将导致任务立即进入失败状态。

from airflow import AirflowException

ValueError 可用于失败和重试

【讨论】:

  • 谢谢。我在气流调度程序中发现了一个错误,上面写着“错误 - running]> 检测为僵尸”。然后是一个后续错误,指出此任务已标记为重试。我猜它认为这个任务已经死了,因为它是一个长期运行的任务。无论如何要进一步指定任务仍然存在并且正在工作?
【解决方案2】:

我在 GCP Cloud Composer [composer-1.11.0-airflow-1.10.9] 中遇到了同样的问题。

对于长时间运行的任务,很有可能(尤其是在使用 KubernetesPodOperator 时)该任务可以被 Airflow Scheduler 标记为 Zombie。

解决方案:-

我将调度程序中的 scheduler_zombie_task_threshold 配置参数值从 300(默认 5 分钟)增加到 1800(30 分钟)。在此之后,我能够运行任务长达 45 分钟,并且不会以失败状态结束。

如何更改参数 -

  1. 转到 GCP 云控制台。
  2. 导航到 Composer。
  3. 打开作曲家。
  4. 转到 AIRFLOW CONFIGURATION OVERRIDES 选项卡
  5. 输入值: 调度器 scheduler_zombie_task_threshold 1800

【讨论】:

    猜你喜欢
    • 2022-11-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-25
    • 1970-01-01
    • 1970-01-01
    • 2023-02-14
    相关资源
    最近更新 更多