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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 13:25:16