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

Docker Compose部署PyFlink作业后Flink UI无法查看已完成作业求助

问题描述

我通过Docker Compose运行一个简单的PyFlink作业,日志显示作业运行成功,但在Flink UI(localhost:8081)的已完成作业板块中看不到该作业。我能正常访问Flink UI,JobManager和TaskManager都处于正常运行状态,但PyFlink作业的状态始终没在UI里展示。我推测作业没有关联到JobManager,是作为独立进程在运行,想知道需要修改哪些配置才能解决这个问题。


相关代码文件

PyFlink作业代码

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.common.typeinfo import Types
from pyflink.datastream.functions import MapFunction

class Splitter(MapFunction):
    def map(self, value):
        return value.split()

def word_count():
    env = StreamExecutionEnvironment.get_execution_environment()

    # 创建包含示例语句的DataStream
    text = env.from_collection([
        'hello world',
        'hello flink',
        'flink is fun'
    ])

    # 应用转换操作实现单词计数
    word_stream = text.map(Splitter(), output_type=Types.STRING()) \
                      .flat_map(lambda words: [(word, 1) for word in words], output_type=Types.TUPLE([Types.STRING(), Types.INT()])) \
                      .key_by(lambda x: x[0]) \
                      .sum(1)

    word_stream.print()

    env.execute('word_count_job')

if __name__ == '__main__':
    word_count()

Dockerfile

# Flink基础镜像
FROM apache/flink:1.18.0-scala_2.12

# 安装Python 3.10和pip
RUN apt-get update && apt-get install -y python3.10 python3-pip

# 设置Python 3.10为`python`和`pip`的默认版本
RUN update-alternatives --install /usr/bin/python python /usr/bin/python3.10 1 \
    && update-alternatives --install /usr/bin/pip pip /usr/bin/pip3 1

# 安装PyFlink
RUN pip install apache-flink==1.18.0

# 设置工作目录
WORKDIR /opt/flink/python_job

# 复制Python Flink作业脚本到容器中
COPY word_count.py /opt/flink/python_job/

# 设置入口点,用Flink CLI将作业提交到Flink集群
# ENTRYPOINT [ "flink", "run", "-m", "jobmanager:8081", "/opt/flink/python_job/word_count.py" ]
ENTRYPOINT [ "python", "/opt/flink/python_job/word_count.py" ]

docker-compose.yml

version: '3.8'
services:
  jobmanager:
    build:
        context: .
        dockerfile: Dockerfile
    image: flink:1.18
    container_name: jobmanager
    hostname: jobmanager
    networks:
      - flink-network
    ports:
      - "8081:8081"  # Flink UI端口
      # - "6123:6123"  # RPC端口
    expose:
      - "6123"
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager
    command: jobmanager

  taskmanager:
    build:
      context: .
      dockerfile: Dockerfile
    image: flink:1.18
    container_name: taskmanager
    hostname: taskmanager
    # expose:
    #   - "6121"
    #   - "6122"
    networks:
      - flink-network
    depends_on:
      - jobmanager
    links:
      - jobmanager:jobmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager
      - TASK_MANAGER_NUMBER_OF_TASK_SLOTS=5
    command: taskmanager


  pyflink-job:
    build:
      context: .  # Dockerfile所在路径
      dockerfile: Dockerfile
    container_name: pyflink-job
    hostname: pyflink-job
    networks:
      - flink-network
    depends_on:
      - jobmanager
      - taskmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager
    # entrypoint: flink run -m jobmanager:8081 /opt/flink/python_job/word_count.py

networks:
  flink-network:
    driver: bridge

解决方案

核心原因

你当前直接用python命令运行PyFlink脚本,默认会启动本地迷你集群(MiniCluster),完全独立于你部署的Flink集群,所以JobManager无法追踪到这个作业的状态。要让作业提交到你搭建的集群运行,需要修改以下配置:

1. 修改Dockerfile的启动入口

把原来直接运行Python脚本的ENTRYPOINT改成用Flink CLI提交作业到JobManager:

# 注释掉原Python启动命令
# ENTRYPOINT [ "python", "/opt/flink/python_job/word_count.py" ]
# 改用Flink CLI提交作业
ENTRYPOINT [ "flink", "run", "-m", "jobmanager:8081", "/opt/flink/python_job/word_count.py" ]

如果Dockerfile已经修改,这一步可以跳过;如果不想改Dockerfile,也可以在docker-compose.yml中覆盖entrypoint:

pyflink-job:
  # 其他配置保持不变...
  entrypoint: flink run -m jobmanager:8081 /opt/flink/python_job/word_count.py

3. 可选:在PyFlink代码中显式指定集群模式

虽然用Flink CLI提交会自动切换到集群模式,但在代码中显式配置JobManager地址可以避免误触发本地模式:

def word_count():
    env = StreamExecutionEnvironment.get_execution_environment()
    # 显式指定连接到集群的JobManager
    env.set_jobmanager_address("jobmanager:8081")
    
    # 后续代码保持不变...

验证步骤

  1. 重新构建镜像:docker-compose build
  2. 启动集群和作业:docker-compose up
  3. 访问Flink UI(localhost:8081),就能在运行中作业或已完成作业板块看到word_count_job的状态了

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 15:23:09