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

Airflow Kubernetes Operator调用Spark-Submit失败求助

排查Airflow DAG中KubernetesPodOperator触发Spark Driver Pod启动失败问题

核心问题概述

  • 环境:Kubernetes 1.25(客户端/服务端),Airflow官方Helm Chart部署
  • 场景:通过Airflow的KubernetesPodOperator执行spark-submit,触发Spark Driver和Executor Pod,但Driver Pod始终无法进入Running状态最终失败
  • 关键线索:镜像内直接执行spark-submit命令可正常运行,所有依赖Jar包已存在于镜像的/opt/spark/connectors/路径

优先排查步骤

1. 获取Driver Pod的详细状态与日志

先定位Pod启动失败的直接原因,执行以下命令:

# 查看Driver Pod的事件(替换为实际Pod名称)
kubectl describe pod <driver-pod-name> -n development

# 查看Pod启动失败的日志(若Pod已终止)
kubectl logs <driver-pod-name> -n development --previous

重点关注:

  • 镜像拉取是否成功(ImagePullBackOff/ErrImagePull事件)
  • PVC挂载是否失败(FailedMount事件)
  • 资源请求是否超出节点可用资源(Pending状态下的Insufficient CPU/Memory)
  • ServiceAccount权限是否不足(Forbidden事件)

2. 修复DAG中的明显配置错误

(1)重复的Task ID

当前代码中load_table函数返回的KubernetesPodOperator使用固定task_id="dry_run_demo",循环生成任务时会导致重复Task ID,Airflow不允许DAG内存在重复Task ID,会引发解析或执行异常。修改为动态Task ID:

def load_table(table_name, application_args, **kwargs):
    # ... 其他代码 ...
    return KubernetesPodOperator(
        dag=dag_daily,
        name=f"spark-submit-{table_name}",
        task_id=f"spark_submit_for_{table_name}",  # 动态生成唯一Task ID
        # ... 其他参数 ...
    )

(2)Spark Master地址问题

硬编码外部IP的K8s Master地址可能导致Pod内部无法访问API,替换为集群内部服务地址:

k8s_arguments = [
    # ... 其他参数 ...
    '--master=k8s://https://kubernetes.default.svc:443',  # 集群内部API地址
    # ... 其他参数 ...
]

(3)重复的PVC挂载配置

代码中同时通过KubernetesPodOperator挂载了air-connectors PVC,又在spark-submit参数里配置了spark.kubernetes.driver.volumes.persistentVolumeClaim,导致重复挂载。由于Jar包已存在于镜像中,可移除Spark配置中的PVC挂载参数:

k8s_arguments = [
    # ... 移除以下重复配置 ...
    # '--conf=spark.kubernetes.driver.volumes.persistentVolumeClaim.air-connectors.mount.path=/air-connectors/',
    # '--conf=spark.kubernetes.driver.volumes.persistentVolumeClaim.air-connectors.mount.readOnly=false',
    # '--conf=spark.kubernetes.driver.volumes.persistentVolumeClaim.air-connectors.options.claimName=nfspvc-airconnectors',
    # ... 其他参数 ...
]

(4)动态分配与固定Executor数量冲突

同时开启spark.dynamicAllocation.enabled=true和配置固定num-executors会导致冲突,保留一种配置即可:

k8s_arguments = [
    # ... 移除固定executor配置,保留动态分配 ...
    # '--num-executors=1',
    '--conf=spark.dynamicAllocation.enabled=true',
    '--conf=spark.dynamicAllocation.shuffleTracking.enabled=true',
    # ... 其他参数 ...
]

(5)Driver资源配置缺失

注释掉Driver的资源配置可能导致Pod无法申请到足够资源,添加合理的资源参数:

k8s_arguments = [
    # ... 添加以下配置 ...
    '--driver-cores=2',
    '--driver-memory=4G',
    # ... 其他参数 ...
]

3. 验证ServiceAccount权限

确保air-airflow-sa在development namespace下有足够权限创建Executor Pod、访问PVC等:

# 验证创建Pod权限
kubectl auth can-i create pods --as=system:serviceaccount:development:air-airflow-sa -n development

# 验证获取PVC权限
kubectl auth can-i get persistentvolumeclaims --as=system:serviceaccount:development:air-airflow-sa -n development

若权限不足,需为该ServiceAccount绑定对应的Role/RoleBinding。

4. 检查DAG解析阶段的文件读取问题

当前代码在DAG解析阶段直接执行get_tables()读取/csv-directory/success-dag.csv,若该文件不存在于Airflow调度器容器中,会导致DAG解析失败。建议:

  • 确保调度器容器中存在该CSV文件
  • 或改用Airflow Variables存储表列表,在PythonOperator执行阶段读取

总结

优先通过Pod日志和事件定位直接失败原因,再逐一修复DAG中的配置错误(重复Task ID、Master地址、资源配置、权限等),结合镜像内可正常运行的线索,重点排查Airflow与K8s交互的配置问题。

内容的提问来源于stack exchange,提问作者SAGE

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 04:15:38