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