EKS环境下Jupyter Notebook与Spark Operator交互式会话配置及故障排查
在EKS中实现Jupyter Notebook与Spark Operator的交互式会话及执行器自动清理
一、环境前置检查
- 确认EKS集群版本与Spark Operator兼容(建议Spark Operator 1.10+搭配K8s 1.21+)
- 集群节点预留足够CPU/内存资源,避免执行器因资源不足被驱逐
- 配置IAM权限:
- Spark Operator需要创建Pod、Service、ConfigMap等资源的权限
- Jupyter Pod需要提交
SparkApplication资源的权限 - 若涉及S3等云存储,需通过IRSA为Spark服务账号绑定对应访问角色
二、部署并配置Spark Operator
用Helm部署官方Spark Operator,开启交互式会话支持:
helm repo add spark-operator https://googlecloudplatform.github.io/spark-on-k8s-operator helm install spark-operator spark-operator/spark-operator \ --namespace spark-operator \ --create-namespace \ --set webhook.enable=true \ --set enableBatchScheduler=false \ --set sparkJobNamespace=spark-jobs
- 验证部署:
kubectl get pods -n spark-operator,确保Operator Pod正常运行 - 为Spark应用配置默认资源模板:在
spark-jobs命名空间创建ConfigMap,指定执行器的CPU/内存请求、镜像等,避免每次提交都重复配置
三、配置Jupyter Notebook集成Spark
使用带Spark客户端的镜像
自定义Jupyter镜像,预装Spark Operator客户端依赖:FROM jupyter/pyspark-notebook:latest RUN pip install spark-on-k8s-operator构建后推送到ECR镜像仓库,部署到EKS的
spark-jobs命名空间。在Notebook中配置SparkSession
启动SparkSession时指定K8s集群参数,确保与Spark Operator联动:from pyspark.sql import SparkSession # 替换<EKS-API-SERVER-URL>为你的EKS集群API地址 spark = SparkSession.builder \ .appName("Jupyter-Interactive-Spark") \ .master("k8s://https://<EKS-API-SERVER-URL>") \ .config("spark.kubernetes.namespace", "spark-jobs") \ .config("spark.kubernetes.container.image", "apache/spark:3.4.0") \ .config("spark.kubernetes.authenticate.driver.serviceAccountName", "spark-operator-spark") \ .config("spark.executor.memory", "2g") \ .config("spark.executor.cores", "1") \ .getOrCreate()注意:
spark-operator-spark是Helm部署时默认创建的服务账号,需确保它在spark-jobs命名空间有足够权限。
四、配置执行器自动终止
通过Spark配置实现任务完成后自动清理驱动和执行器Pod:
- 在SparkSession中添加动态分配与自动清理参数:
# 开启动态分配,空闲执行器自动销毁 spark.conf.set("spark.dynamicAllocation.enabled", "true") spark.conf.set("spark.dynamicAllocation.shuffleTracking.enabled", "true") spark.conf.set("spark.dynamicAllocation.minExecutors", "0") spark.conf.set("spark.dynamicAllocation.maxExecutors", "5") spark.conf.set("spark.dynamicAllocation.executorIdleTimeout", "60s") # 任务结束后自动删除驱动和执行器Pod spark.conf.set("spark.kubernetes.driver.pod.deleteOnTermination", "true") spark.conf.set("spark.kubernetes.executor.pod.deleteOnTermination", "true") - 任务完成后主动停止SparkSession:
执行后驱动Pod和所有执行器Pod会被自动删除。spark.stop()
五、执行器启动失败排查
- 查看Spark Operator日志:
kubectl logs -n spark-operator deployment/spark-operator,排查资源创建失败的错误(比如权限不足、配额超限) - 查看驱动Pod日志:
kubectl logs -n spark-jobs <driver-pod-name>,检查是否有镜像拉取失败、K8s API连接错误、存储访问权限问题 - 检查资源配额:确认
spark-jobs命名空间的CPU/内存配额满足执行器的资源请求 - 检查网络策略:确保驱动Pod与执行器Pod之间能正常通信,没有被网络策略阻断
- 验证IRSA配置:若使用S3等云存储,确认Spark服务账号的IAM角色绑定正确,能正常访问存储资源
内容的提问来源于stack exchange,提问作者Abdul
相关产品推荐
相关产品推荐

