You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow 2.2.2 K8s Executor任务被SIGKILL XCom返回None问题

根因定位

两个叠加问题导致对应现象:

  1. Airflow 2.2.x版本的ShortCircuitOperator默认将do_xcom_push设为False,python_callable的返回值默认不会被推送到XCom,下游任务自然拉取不到值。
  2. 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为默认返回值keyreturn_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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.31 20:36:21