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

如何在Kubernetes的Flink Runner上提交Beam Python作业并部署镜像?

flink master container就是你集群中JobManager Pod内的核心容器,它负责:

  • 整个Flink集群的资源调度与作业生命周期管理
  • 提供你已通过Service和Ingress暴露的Flink Web UI
  • 接收Beam作业的提交请求

对应你的Pipeline配置,需要把本地调试用的localhost:8081替换为Kubernetes集群内JobManager Service的集群内部地址,示例调整后配置:

options = PipelineOptions([
    "--runner=FlinkRunner",
    "--flink_version=1.10",
    "--flink_master=flink-jobmanager.default.svc.cluster.local:8081",  # 替换为实际Service地址
    "--environment_type=EXTERNAL",
    "--environment_config=localhost:50000"
])

注:flink-jobmanager是JobManager Service的默认名称,default为命名空间,需根据你的实际集群配置修改。

二、Python代码镜像的部署方案

必须将包含Python作业代码的镜像部署在每个Task Manager Pod的Beam worker pool容器中,原因如下:

  • Beam的EXTERNAL环境模式下,Flink TaskManager会通过localhost:50000调用同Pod内的Beam Worker Pool执行Python算子,同一个Pod内的容器共享网络命名空间,localhost通信直接可达
  • 如果单独部署Python代码Pod,多个TaskManager无法精准绑定对应的Worker Pool,会出现调度不匹配、通信失败的问题

具体操作步骤:

  1. 构建自定义Beam Worker镜像:
    以官方Beam Python Worker镜像为基础,添加你的作业代码和依赖,示例Dockerfile:
    FROM apache/beam_python3.7_sdk:2.17.0  # 选择与Flink1.10兼容的Beam版本
    COPY your_beam_job.py /app/
    COPY requirements.txt /app/
    RUN pip install -r /app/requirements.txt
    
  2. 修改Task Manager的Pod模板:
    将原来的Beam worker pool容器镜像替换为你构建的自定义镜像,确保容器暴露50000端口,保持与Flink TaskManager容器在同一个Pod内。

三、作业提交方式

你有两种可行的作业提交方式:

  • 集群内提交:在Kubernetes中创建临时客户端Pod,安装Beam及依赖,运行作业代码(使用集群内的JobManager地址)
  • 本地提交:通过kubectl port-forward service/flink-jobmanager 8081:8081将JobManager端口转发到本地,再在本地运行作业代码(此时--flink_master可设为localhost:8081)

四、关键注意事项

  • 版本兼容性:确保Beam版本与Flink版本匹配(比如Flink1.10对应Beam2.17.x,避免版本不兼容引发异常)
  • 资源配置:根据作业负载调整Task Manager Pod的CPU、内存资源,以及Beam Worker Pool的进程数
  • 日志排查:通过kubectl logs查看Flink TaskManager和Beam Worker Pool的日志,定位作业运行中的问题

内容的提问来源于stack exchange,提问作者John

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 04:17:03