Airflow PythonSensor仅执行一次poke导致任务无限运行问题
问题根源分析
你的PythonSensor仅执行一次python_callable且任务持续处于Running状态,核心原因有两个:
1. 请求未设置超时,导致函数阻塞
你的requests.post(...)未配置timeout参数,当目标服务无响应或响应极慢时,请求会一直卡在等待环节,无法执行后续的判断、异常处理或返回逻辑。这会让第一次poke调用持续阻塞,既无法返回结果结束任务,也不会触发下一次poke。
2. 函数返回逻辑不符合Sensor预期
即便请求正常返回,你的函数逻辑也存在问题:无论校验是否通过(甚至捕获异常后),最终都会return True。按照PythonSensor的规则,只要python_callable返回True,任务就会立即标记为成功结束,不会再进行后续的poke调用。但你提到任务一直Running,说明大概率是第一个原因导致请求卡住,根本没走到return语句。
修复方案
给请求添加超时限制:修改
requests.post,添加timeout参数避免无限等待,示例:validation_output = requests.post(..., timeout=30) # 设置30秒超时调整函数返回逻辑:PythonSensor需要
python_callable返回False来触发下一次poke。修改函数,在校验不通过或异常时返回False,校验通过时返回True:def my_python_callable(): try: validation_output = requests.post(..., timeout=30) validation_output_json = json.loads(validation_output.content) if not validation_output_json['is_valid']: send_error(validation_output_json['error']) return False # 校验不通过,触发下一次poke return True # 校验通过,结束任务 except Exception as e: send_error(str(e)) return False # 异常时返回False,继续重试
调整后,当校验不通过或发生异常时,函数返回False,PythonSensor会每隔60秒重新调用一次,直到返回True或触发timeout(15分钟后任务失败)。
内容的提问来源于stack exchange,提问作者Nacho
相关产品推荐
相关产品推荐

