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" ]
2. 调整docker-compose.yml的pyflink-job服务
如果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") # 后续代码保持不变...
验证步骤
- 重新构建镜像:
docker-compose build - 启动集群和作业:
docker-compose up - 访问Flink UI(localhost:8081),就能在运行中作业或已完成作业板块看到
word_count_job的状态了
内容的提问来源于stack exchange,提问作者Vishal
相关产品推荐
相关产品推荐

