Spark 2.3集成K8s后,如何通过AWS Lambda向K8s集群提交任务?
如何通过AWS Lambda向K8s集群提交Spark任务
嘿,我刚好有过类似的实践经验,结合你提到的Spark 2.3集成K8s的场景,咱们可以通过以下步骤来实现用AWS Lambda触发Spark任务提交到K8s集群:
一、先搞定Lambda的权限和运行环境
这一步是基础,没做好后面肯定踩坑:
- IAM角色权限配置:给Lambda的执行角色加这些权限:
- 访问K8s集群的权限:如果是AWS EKS,直接附加
AmazonEKSClusterPolicy和AmazonEKSServicePolicy就行,还要确保角色能通过IAM认证访问集群(后续kubeconfig会用到这个)。 - S3权限:如果你的Spark作业包、配置存在S3,得加
S3:GetObject这类权限,方便拉取资源。 - CloudWatch Logs权限:必须加!不然调试的时候连日志都看不到,给角色配
logs:CreateLogGroup、logs:CreateLogStream、logs:PutLogEvents这三个权限就够。
- 访问K8s集群的权限:如果是AWS EKS,直接附加
- 运行环境选Python或Java:推荐Python,因为代码写起来更灵活,而且Spark有pyspark支持;如果是Java的话,要注意打包依赖的大小,别超过Lambda的限制。
二、两种实现方式任你选
方式1:直接在Lambda里跑spark-submit命令
这种最贴近你之前用控制台操作的习惯,不过得给Lambda装Spark客户端:
- 把Spark客户端做成Lambda层:Spark 2.3的客户端压缩后大概几百MB,刚好在Lambda层的大小限制内,上传成层后,Lambda运行时会把它挂载到
/opt目录。 - Python示例代码:
注意哈:Spark镜像得和2.3版本匹配,还要确保K8s集群能拉到这个镜像;另外kubeconfig一定要加密存储,别直接写在代码里。import subprocess import os import boto3 def lambda_handler(event, context): # 从Secrets Manager拉取加密的kubeconfig(别硬编码!) secrets_manager = boto3.client('secretsmanager') kubeconfig_secret = secrets_manager.get_secret_value(SecretId='k8s-kubeconfig') with open('/tmp/kubeconfig', 'w') as f: f.write(kubeconfig_secret['SecretString']) os.environ['KUBECONFIG'] = '/tmp/kubeconfig' # 构造spark-submit命令,替换成你的集群和作业信息 spark_submit_cmd = [ '/opt/spark/bin/spark-submit', '--master', 'k8s://https://<你的EKS集群API地址>', '--deploy-mode', 'cluster', '--name', 'lambda-triggered-spark-job', '--class', 'com.yourcompany.YourSparkMainClass', '--conf', 'spark.kubernetes.container.image=你的Spark镜像:2.3', '--conf', 'spark.kubernetes.authenticate.driver.serviceAccountName=spark-sa', 's3://你的存储桶路径/你的spark作业.jar' ] # 执行命令并捕获输出 result = subprocess.run(spark_submit_cmd, capture_output=True, text=True) if result.returncode == 0: return { 'statusCode': 200, 'body': f"任务提交成功!输出: {result.stdout}" } else: return { 'statusCode': 500, 'body': f"任务提交失败,错误信息: {result.stderr}" }
方式2:用Spark的Python API直接提交作业
这种适合把作业逻辑直接写在Lambda里,不用单独打包Jar包:
- Python示例(pyspark):
先把pyspark打包成Lambda层,然后写代码:
这种方式的好处是不用维护单独的Jar包,适合简单的作业;如果是复杂的大数据作业,还是方式1更靠谱。from pyspark.sql import SparkSession import os import boto3 def lambda_handler(event, context): # 同样先配置kubeconfig secrets_manager = boto3.client('secretsmanager') kubeconfig_secret = secrets_manager.get_secret_value(SecretId='k8s-kubeconfig') with open('/tmp/kubeconfig', 'w') as f: f.write(kubeconfig_secret['SecretString']) os.environ['KUBECONFIG'] = '/tmp/kubeconfig' # 初始化SparkSession,配置K8s参数 spark = SparkSession.builder \ .master("k8s://https://<你的EKS集群API地址>") \ .appName("lambda-spark-api-job") \ .config("spark.kubernetes.container.image", "你的Spark镜像:2.3") \ .config("spark.kubernetes.authenticate.driver.serviceAccountName", "spark-sa") \ .getOrCreate() # 这里写你的Spark作业逻辑,比如读取S3数据然后处理 df = spark.read.csv("s3://你的存储桶/输入数据.csv") df.write.parquet("s3://你的存储桶/输出数据") spark.stop() return { 'statusCode': 200, 'body': "Spark作业执行完成!" }
三、几个必注意的坑
- K8s服务账号权限:一定要给Spark的driver服务账号(比如上面的
spark-sa)绑定足够的ClusterRole,不然它在K8s里创建pod、service的时候会被拒。 - Lambda超时限制:Lambda最长只能跑15分钟,所以如果你的Spark作业是长时间运行的,别让Lambda等着作业完成,只让Lambda负责提交任务,然后用K8s的API或者Spark的REST API去监控状态。
- Spark镜像的依赖:镜像里要包含AWS SDK,不然访问S3会出问题;如果用的是官方Spark镜像,可能得自己二次打包加依赖。
- kubeconfig的有效期:如果用的是IAM认证的kubeconfig,注意它的有效期,最好在Lambda运行时动态生成,而不是用固定的kubeconfig文件。
四、测试调试小技巧
- 先在本地用Docker模拟Lambda环境跑代码,确保
spark-submit或者Spark API能正常提交任务到K8s。 - 部署Lambda后,一定要去CloudWatch Logs看执行日志,权限问题、路径问题、配置错误都能在日志里找到线索。
内容的提问来源于stack exchange,提问作者amza
相关产品推荐
相关产品推荐

