Airflow的on_failure_callback回调函数未执行函数内全部代码问题求助
问题根因
你的自定义回调函数在执行第二次HTTP POST请求时触发了未捕获的异常,导致函数直接中断退出,因此后续的FAILED TASK日志没有输出。
常见触发原因
- 参数拼写错误:第二次
requests.post调用时你传入的参数是header=headers,而requests库接收请求头的参数名是headers(末尾多一个s),参数不匹配会直接触发TypeError异常。 - 网络连通异常:Airflow Worker节点无法访问第二个请求的目标地址
http://myerrorappli:7180,可能是DNS解析失败、端口未放通、防火墙拦截导致。 - SSL证书校验失败:第二个请求开启了
verify=True,如果目标服务的SSL证书没有被Airflow运行环境的信任根证书库收录,会触发证书校验异常。 - 请求无超时卡住:你没有给requests请求设置超时时间,默认会无限等待,若请求长时间无响应,可能会触发Airflow回调的执行超时阈值被强制终止进程。
修复方案
- 首先修正参数拼写错误,将第二个post调用的
header=headers改为headers=headers。 - 给所有requests请求添加
timeout参数,避免无限等待,建议设置为5~30秒,根据接口实际响应时长调整。 - 给回调函数整体添加try-except异常捕获逻辑,打印完整异常栈方便定位问题,参考代码如下:
import traceback import requests import json import logging def custom_failure_function(context): try: logging.error("These task instances ahhh") to_json= json.loads(t_teams) var1= json.dumps(to_json) print(var1) r = requests.post('https://myteamschannel/teams', data=var1,verify=False, timeout=10) logging.error("hello") runID='OPERATION_CONTEXT .OCV8.TEST2 alarm_object 193' headers = {'Content-Type':'text/xml'} alarmRequest='<soapenv:Envelope xmlns:soapenv="http://schemas.xmlsoap.org/soap/envelope/" xmlns:oper="http://172.19.146.147:7180/TeMIP_WS/services/OPERATION_CONTEXT-alarm_object"><soapenv:Header xmlns:wsa="http://www.w3.org/2005/08/addressing"><wsu:Timestamp xmlns:wsu="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-wssecurity-utility-1.0.xsd"><wsu:Created>2014-05-22T11:57:38.267Z</wsu:Created><wsu:Expires>2014-05-22T12:02:38.000Z</wsu:Expires></wsu:Timestamp><wsse:Security xmlns:wsse="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-wssecurity-secext-1.0.xsd" soapenv:mustUnderstand="1"><wsse:UsernameToken xmlns:wsu="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-wssecurity-utility-1.0.xsd"><wsse:Username>girws</wsse:Username><wsse:Password Type="http://docs.oasis-open.org/wss/2004/01/oasis-200401-wss-username-token-profile-1.0#PasswordText">Temip</wsse:Password></wsse:UsernameToken></wsse:Security></soapenv:Header> <soapenv:Body> <oper:Set_Request xmlns:oper="http://172.19.146.147:7180/TeMIP_WS/services/OPERATION_CONTEXT-alarm_object"><EntitySpec><Natural> '+ runID + '</Natural></EntitySpec><Arguments> <Attribute_Values><Filtering_Type>' + 'AUTOFAIL' + '</Filtering_Type></Attribute_Values></Arguments></oper:Set_Request> </soapenv:Body> </soapenv:Envelope>' # 修正参数、添加超时 r = requests.post( 'http://myerrorappli:7180/TeMIP_WS/services/OPERATION_CONTEXT-alarm_object', headers=headers, data=alarmRequest, verify=True, timeout=10 ) logging.error ("FAILED TASK") logging.error("============================================") except Exception as e: logging.error(f"回调执行失败,异常信息:{str(e)}") logging.error(f"完整异常栈:{traceback.format_exc()}")
- 登录到Airflow Worker节点,手动测试第二个接口的连通性:可通过curl命令模拟发送SOAP请求,确认网络、证书、接口返回都正常。
内容的提问来源于stack exchange,提问作者SrikanthR
相关产品推荐
相关产品推荐

