【问题标题】:Slack API request, limiting to 1 request per DAG failure (Airflow)Slack API 请求,每个 DAG 故障限制为 1 个请求 (Airflow)
【发布时间】:2020-12-10 15:38:15
【问题描述】:

您好,这里是小数据工程师!

由于某种奇怪的原因,我的 task_fail_slack_alert 模块多次触发 Slack API 请求,然后多次出现在我们的 Slack 频道中,真的很烦人。我的模块应该只在 Slack 频道中运行并显示与失败的任务数量相同的数量。

我错过了什么?

import os
from airflow.models 
import Variable 
import json import requests

def get_channel_name():
  channel = '#airflow_alerts_local'
  env = Variable.get('env', None)
  if env == 'prod':
       channel = '#airflow_alerts'
  elif env == 'dev':
       channel = '#airflow_alerts_dev'
  return channel

def task_fail_slack_alert(context):
   webhook_url = os.environ.get('SLACK_URL')

   slack_data = {
         'channel': get_channel_name(),
         'text':
             """ :red_circle: Task Failed.
            *Task*: {task}
            *Dag*: {dag}
            *Execution Time*: {exec_date}
            *Log Url*: {log_url}
            """.format(
            task=context.get('task_instance').task_id,
            dag=context.get('task_instance').dag_id,
            ti=context.get('task_instance'),
            exec_date=context.get('execution_date'),
            log_url=context.get('task_instance').log_url,
          )}

    response = requests.post(webhook_url, data=json.dumps(slack_data),
           headers={'Content-Type': 'application/json'})

    if response.status_code != 200: 
        raise ValueError( 'Request to slack returned an error %s, 
                    the response is:\n%s'(response.status_code, response.text))
    task_fail_slack_alert(context)

这就是我在每个 dag 的参数中显示它的方式:

default_args = {
    'on_failure_callback': task_fail_slack_alert,
}

【问题讨论】:

    标签: python-3.x python-requests airflow slack-api directed-acyclic-graphs


    【解决方案1】:

    你提供的代码是递归的:

    def task_fail_slack_alert(context):
        ......
        task_fail_slack_alert(context)
    

    删除不需要的递归。

    【讨论】:

    • 我只是想确保我理解正确,你是说我应该从我的函数中删除函数调用 task_fail_slack_alert(context) 吗?完全喜欢?
    • @KristiLuna task_fail_slack_alert 在任务失败时由 Airflow 自动调用。因此它会执行您在该函数中编写的代码。在它向 slack requests.post 发布消息后,您还显式调用了 task_fail_slack_alert,因此该函数是递归的,并一次又一次地执行整个代码。
    • @kristiluna 留下这个电话会实现什么?说“部分”与“完全”删除呼叫是什么意思?我不明白...
    • @acushner 我已经删除了 task_fail_slack_alert(context),现在我没有收到任何警报
    • @acushner 是的!只是一个警报,我觉得我现在被困在中间了。我要么收到永无止境的警报,因为我有一个递归问题,要么当我删除函数调用时我没有收到任何信息,哈哈:-/
    猜你喜欢
    • 1970-01-01
    • 2020-12-23
    • 2018-03-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-10-29
    • 2015-02-17
    • 1970-01-01
    相关资源
    最近更新 更多