Airflow Kubernetes Operator调用Spark-Submit失败求助
核心问题概述
- 环境: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

