【问题标题】:Airflow, on failure, make all dags do something in specific气流,失败时,让所有 dags 做一些特定的事情
【发布时间】: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


【解决方案1】:

如果 DAG 失败,IMO 最简单的方法是将其定义为默认参数。

default_args = { '所有者':'气流', 'depends_on_past':错误, 'start_date': 日期时间(2015, 6, 1), '电子邮件':['airflow@example.com'], 'email_on_failure':错误, 'email_on_retry':错误, “重试”:1, 'retry_delay': timedelta(分钟=5), # '队列': 'bash_queue', # 'pool': '回填', # 'priority_weight': 10, # 'end_date': datetime(2016, 1, 1), }

如果您想根据您的任务依赖性指定发送电子邮件的行为,您也可以使用 sendgrid 运算符。 https://github.com/apache/airflow/blob/master/airflow/contrib/utils/sendgrid.py

【讨论】:

    【解决方案2】:

    最简单的方法:如果 BaseOperator 中的 email_on_retryemail_on_failure 属性为 true(默认为 true),并且设置了 airflow 邮件配置,则airflow 会在重试时发送邮件并失败。

    使用自定义运算符:

    def on_failure_callback(context):
        # with mail:
        error_mail = EmailOperator(
            task_id='error_mail',
            to='user@example.com',
            subject='Fail',
            html_content='a task failed',
            mime_charset='utf-8')
        error_mail.execute({})  # no need to return, just execute
    
        # with slack:
        error_message = SlackAPIPostOperator(
            task_id='error_message',
            token=getSlackToken(),
            text='a task failed',
            channel=SLACK_CHANNEL,
            username=SLACK_USER)
        error_message.execute({})  # no need to return, just execute
    
    dag_args = {
        "retries": 2,
        "retry_delay": timedelta(minutes=2),
        'on_failure_callback': on_failure_callback
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-10-17
      • 1970-01-01
      • 2019-06-18
      • 1970-01-01
      • 2022-01-18
      • 1970-01-01
      相关资源
      最近更新 更多