【发布时间】:2019-05-10 22:07:16
【问题描述】:
我们有很多在 Airflow 上运行的 DAG。当某些事情失败时,我们希望得到通知,或者做出特定的动作:我已经尝试通过装饰器
def on_failure_callback(f):
@wraps(f)
def wrap(*args, **kwargs):
try:
return f(*args, **kwargs)
except Exception as e:
return f"An exception {e} on ocurred on{f}"
return wrap
这可行,但有必要装饰我们想要具有此行为的任何功能。
我看到this 并尝试像这样实现它:
def on_failure_callback(context):
operator = PythonOperator(
python_callable=failure)
return operator.execute(context=context)
def failure():
return 'Failure in the failure func'
dag_args = {
"retries": 2,
"retry_delay": timedelta(minutes=2),
'on_failure_callback': on_failure_callback
}
然后在 DAG 定义中,我使用[...] default_args=dag_args [...],但是这个选项不起作用。
实现这一目标的最佳方法是什么?
谢谢
【问题讨论】:
-
你想要什么样的通知?邮件、slack 还是其他?
-
@gokmust Slack,但我配置任何东西都会更好
标签: python etl decorator airflow