【发布时间】: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