Airflow 2.2.2 K8s Executor任务被SIGKILL XCom返回None问题
根因定位
两个叠加问题导致对应现象:
- Airflow 2.2.x版本的
ShortCircuitOperator默认将do_xcom_push设为False,python_callable的返回值默认不会被推送到XCom,下游任务自然拉取不到值。 - Airflow 2.2.2版本搭配Kubernetes Executor存在已知bug:任务被标记为SUCCESS状态后,executor会立刻触发worker Pod删除流程,不会等待任务完成XCom推送、日志上传等收尾操作,导致进程先后收到SIGTERM、SIGKILL被强制终止。此时抛出的 Job
was killed before it finished (likely due to running out of memory) 是通用报错文案,和实际内存不足无关,属于误报。
两个问题叠加后,就算手动修改配置尝试推送XCom,也会因为Pod被提前回收导致XCom写入流程被打断,下游始终拉取到None,且下游任务如果存在收尾流程慢的情况,会触发相同的误报日志。
解决方案
按以下顺序调整即可修复:
- 显式开启ShortCircuitOperator的XCom推送能力,初始化任务时传入
do_xcom_push=True覆盖默认配置:
TASK_1 = ShortCircuitOperator( task_id="task_1", python_callable=test_script_1, executor_config=EXECUTOR_CONFIG, do_xcom_push=True )
- 修复Kubernetes Executor提前回收Pod的问题,二选一即可:
- 优先选择升级Airflow到2.3.0及以上稳定版本,该版本已官方修复K8s Executor提前删除worker Pod的逻辑缺陷,会等待所有任务收尾流程完成后再回收Pod资源。
- 若暂时无法升级版本,修改Airflow配置文件
airflow.cfg中[kubernetes]段的delete_worker_pods_on_success参数为False,同时给worker Pod配置ttlSecondsAfterFinished=60,由K8s本身在Pod运行结束1分钟后自动回收资源,避免executor主动删Pod打断进程。
- 优化下游XCom拉取逻辑,显式指定XCom的key为默认返回值key
return_value,避免隐式匹配导致的拉取失败:
def test_script_2(**context) -> List[str]: task_instance = context["task_instance"] return_value = task_instance.xcom_pull(task_ids="task_1", key="return_value") print("logging return value of first task ", return_value)
调整完成后重新触发DAG运行,即可正常拉取到上游任务的返回值,SIGKILL和误报OOM的日志也会消失。
内容的提问来源于stack exchange,提问作者Vineet
相关产品推荐
相关产品推荐

