如何在Kubernetes已部署的Spark应用上运行Apache Beam应用?
问题描述
我已在Kubernetes上部署了Spark,同时安装了GCP Spark Operator,相关配置如下:
部署文件 deployment.yaml
apiVersion: "sparkoperator.k8s.io/v1beta2" kind: SparkApplication metadata: name: pyspark-pi namespace: default spec: type: Python pythonVersion: "3" mode: cluster image: "user/pyspark-app:1.0" imagePullPolicy: Always mainApplicationFile: local:///app/pyspark-app.py sparkVersion: "3.1.1" restartPolicy: type: OnFailure onFailureRetries: 3 onFailureRetryInterval: 10 onSubmissionFailureRetries: 5 onSubmissionFailureRetryInterval: 20 driver: cores: 1 coreLimit: "1200m" memory: "512m" labels: version: 3.1.1 serviceAccount: spark executor: cores: 1 instances: 2 memory: "512m" labels: version: 3.1.1
服务文件 service.yaml
apiVersion: "rbac.authorization.k8s.io/v1" kind: ClusterRoleBinding metadata: name: spark-operator-role namespace: default roleRef: apiGroup: "rbac.authorization.k8s.io" kind: ClusterRole name: edit subjects: - kind: ServiceAccount name: spark namespace: default
当前运行的服务
pyspark-pi-84dad9839f7f5f43-driver-svc ClusterIP None <none> 7078/TCP,7079/TCP,4040/TCP 2d14h
我希望在该Spark Driver上运行Apache Beam应用,示例代码如下:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions options = PipelineOptions([ "--runner=PortableRunner", "--job_endpoint=http://127.0.0.1:4040/", "--environment_type=DOCKER", "--environment_config=docker.io/apache/beam_python3.7_sdk:2.33.0" ]) # lets have a sample string data = ["this is sample data", "this is yet another sample data"] # create a pipeline pipeline = beam.Pipeline(options=options) counts = (pipeline | "create" >> beam.Create(data) | "split" >> beam.ParDo(lambda row: row.split(" ")) | "pair" >> beam.Map(lambda w: (w, 1)) | "group" >> beam.CombinePerKey(sum)) # lets collect our result with a map transformation into output array output = [] def collect(row): output.append(row) return True counts | "print" >> beam.Map(collect) # Run the pipeline result = pipeline.run() # lets wait until result a available result.wait_until_finish() # print the output print(output)
运行上述代码时出现以下错误:
$ python beam2.py WARNING:root:Make sure that locally built Python SDK docker image has Python 3.7 interpreter. Traceback (most recent call last): File "beam2.py", line 31, in <module> result = pipeline.run() File "C:\Users\eapasnr\Anaconda3\envs\oden2\lib\site-packages\apache_beam\pipeline.py", line 565, in run return self.runner.run_pipeline(self, self._options) File "C:\Users\eapasnr\Anaconda3\envs\oden2\lib\site-packages\apache_beam\runners\portability\portable_runner.py", line 438, in run_pipeline job_service_handle = self.create_job_service(options) File "C:\Users\eapasnr\Anaconda3\envs\oden2\lib\site-packages\apache_beam\runners\portability\portable_runner.py", line 317, in create_job_service return self.create_job_service_handle(server.start(), options) File "C:\Users\eapasnr\Anaconda3\envs\oden2\lib\site-packages\apache_beam\runners\portability\job_server.py", line 54, in start grpc.channel_ready_future(channel).result(timeout=self._timeout) File "C:\Users\eapasnr\Anaconda3\envs\oden2\lib\site-packages\grpc\_utilities.py", line 139, in result self._block(timeout) File "C:\Users\eapasnr\Anaconda3\envs\oden2\lib\site-packages\grpc\_utilities.py", line 85, in _block raise grpc.FutureTimeoutError() grpc.FutureTimeoutError
我认为问题出在Pipeline Options的配置上,尤其是job_endpoint参数。未配置Pipeline Options时应用可正常运行并输出结果。请问我应该配置什么IP和主机地址到job_endpoint才能让Beam应用在Spark上正常运行?
解决方案
1. 核心问题分析
你遇到的超时是因为job_endpoint指向本地127.0.0.1,但Beam代码的运行环境和Spark Driver不在同一个网络空间:如果是本地运行,无法直接访问K8s集群内部服务;如果是集群内其他Pod运行,127.0.0.1指向的是当前Pod而非Spark Driver。
2. 正确的job_endpoint配置
根据你提供的Spark Driver服务信息,正确的地址需要遵循K8s内部服务的DNS规则:服务名.命名空间.svc.cluster.local:端口,所以你的job_endpoint应该设置为:
http://pyspark-pi-84dad9839f7f5f43-driver-svc.default.svc.cluster.local:4040
3. 分场景调整配置
场景A:Beam代码在K8s集群内的Pod中运行
直接使用上述完整DNS地址即可,K8s集群内部会自动解析这个地址到Spark Driver服务。
场景B:Beam代码在本地机器运行
本地无法直接访问K8s内部服务,需要先通过端口转发把Spark Driver的4040端口映射到本地:
kubectl port-forward service/pyspark-pi-84dad9839f7f5f43-driver-svc 4040:4040 -n default
端口转发成功后,job_endpoint保持http://127.0.0.1:4040即可,本地请求会被转发到K8s内的Spark Driver。
4. 额外注意事项
- 确保Spark Driver开启了Portable Runner支持:Spark 3.x默认支持,但需确认镜像中包含Beam Spark Runner的依赖包,或启动时通过
--conf spark.driver.extraClassPath加载对应jar。 - 保持Beam SDK版本(2.33.0)和Spark版本(3.1.1)兼容,避免版本不匹配导致的隐性问题。
environment_config中的Docker镜像Python版本(3.7)要和Spark应用的Python版本(3)保持一致。
内容的提问来源于stack exchange,提问作者user19930511
相关产品推荐
相关产品推荐

