导读:本期聚焦于高宇创作的《如何在 Docker 中部署 Apache Flink 集群并提交流处理作业?》,敬请观看详情。要在本地快速验证 Flink 作业的逻辑,手动安装 JDK、配置环境变量、下载发行包再启动集群往往耗时且容易出错。借助 Docker 可以把 JobManager 和 TaskManager 分别封装成容器,用一条 docker-compose 命令拉起完整集群,作业提交和依赖管理也会更加干净。本文基于 Flink 官方镜像,演示如何用 Docker Compose 编排一个包含 JobManager 与 TaskManager 的最小集群,说明 Flink Web UI 的访问方式,并通过命令行客户端提交一个简单的 WordCount 作业。同时讨论容器网络、状态持久化、日志查看以及镜像版本选择等集成细节,帮助读者在开发与测试环境中快速复用这套方案。

Apache Flink 是一个面向分布式流处理和批处理的统一计算框架,它的集群通常由 JobManager 和多个 TaskManager 组成。JobManager 负责协调分布式执行、调度任务和管理检查点,而 TaskManager 则是真正执行数据流算子并缓存数据的工作节点。手动部署这套架构需要对齐 JDK 版本、配置主机名和端口、分发密钥等,过程比较繁琐。Docker 的出现让这些组件可以被标准化封装,配合 Docker Compose 可以一键启动完整的 Flink 集群,非常适合本地开发、联调测试以及小规模演示场景。

如何在 Docker 中部署 Apache Flink 集群并提交流处理作业?

Flink 官方提供了 flink 镜像,该镜像内部已经包含启动脚本、Java 运行时以及 Web UI 所需的静态资源。使用 Docker 集成 Flink 的核心思路是让 JobManager 容器负责协调,TaskManager 容器通过设置 jobmanager.rpc.address 指向 JobManager 的服务名来加入集群。由于容器之间可以通过 Docker 网络互相解析服务名,因此不需要手动配置 IP 地址,这大幅降低了部署门槛。

为什么选择 Docker 来运行 Flink

在传统部署模式下,每新增一个测试环境都需要重复安装 JDK、解压 Flink 发行包、修改 flink-conf.yaml 并启动多个进程。如果团队成员使用的操作系统不一致,还可能遇到路径分隔符、权限以及依赖库差异带来的问题。Docker 将这些依赖打包进镜像,保证每次启动的运行时行为一致,排除了环境差异导致的作业执行不一致。

另一个重要优势是资源隔离和快速回收。Flink 作业在开发和调试阶段经常需要反复修改代码并重新提交,手动清理 TaskManager 进程和临时目录比较麻烦。使用容器方案后,只需执行 docker compose down 就能释放全部容器资源,必要时还可以通过 docker compose up -d --scale taskmanager=3 动态扩展 TaskManager 数量,模拟不同并行度下的资源分配情况。这种弹性对测试作业的扩展行为非常有帮助。

此外,Docker 镜像自带版本标签管理,例如 flink:1.18.1-scala_2.12 可以明确锁定 Flink 版本和 Scala 版本。当需要验证某个作业在不同 Flink 版本下的兼容性时,只修改镜像标签即可切换环境,无需重新安装整个发行包。这种版本切换能力在升级验证和回归测试中非常实用。

使用 Docker Compose 搭建 Flink 集群

下面是一个最小化的 docker-compose.yml 配置,它定义了一个 JobManager 和一个 TaskManager 服务。JobManager 暴露 8081 端口用于访问 Flink Web UI,而 TaskManager 不对外暴露端口,只通过内部网络与 JobManager 通信。

version: "3"
services:
  jobmanager:
    image: flink:1.18.1-scala_2.12
    container_name: flink-jobmanager
    ports:
      - "8081:8081"
    command: jobmanager
    environment:
      - |
        FLINK_PROPERTIES=
        jobmanager.rpc.address: jobmanager
        parallelism.default: 2
    networks:
      - flink-network

  taskmanager:
    image: flink:1.18.1-scala_2.12
    container_name: flink-taskmanager
    depends_on:
      - jobmanager
    command: taskmanager
    environment:
      - |
        FLINK_PROPERTIES=
        jobmanager.rpc.address: jobmanager
        taskmanager.numberOfTaskSlots: 4
        parallelism.default: 2
    networks:
      - flink-network

networks:
  flink-network:

上面的配置中,jobmanager.rpc.address 被设置为服务名 jobmanager。Docker Compose 会创建一个名为 flink-network 的桥接网络,并把两个服务加入到该网络中。TaskManager 启动时会通过这个地址向 JobManager 注册自身,从而形成完整的 Flink 集群。若需要使用外部配置文件,可以将本地的 flink-conf.yaml 挂载到容器内的 /opt/flink/conf 目录,但要注意容器内路径与本地路径的映射关系。

执行 docker compose up -d 后,可以访问 http://localhost:8081 打开 Flink Web UI。默认情况下,页面会显示已注册的 TaskManager 数量、可用任务槽位以及运行中的作业。如果页面没有显示 TaskManager,通常是因为 jobmanager.rpc.address 写成了 localhost 或容器自身的 IP,导致 TaskManager 无法正确连接。检查容器日志 docker logs flink-taskmanager 可以看到详细的注册信息。

提交作业与查看运行结果

Flink 提供了 Web UI 上传作业包和命令行客户端两种提交方式。使用命令行时,需要先进入 JobManager 容器内部,再调用 flink run 命令。以下命令进入容器并提交一个示例作业,示例作业通常位于 /opt/flink/examples/streaming 目录下。

docker exec -it flink-jobmanager /bin/bash
flink run /opt/flink/examples/streaming/WordCount.jar

如果作业需要读取外部数据源,例如 Kafka,则需要额外准备包含 Kafka 连接器的作业 JAR 包,并将其复制到容器中或挂载到容器目录。容器内作业的依赖管理遵循 Flink 的标准机制,可以使用 flink run -c com.example.MainClass /path/to/your-job.jar 指定入口类。对于本地开发环境,更推荐把作业打包成可执行 JAR 后通过卷挂载的方式提供给容器,避免频繁构建镜像。

作业提交成功后,Web UI 的 Running Jobs 列表会显示该作业,点击进入可以看到数据流图、算子并行度、背压状态以及检查点统计。在 TaskManager 的日志中也能看到算子输出的记录。如果使用 WordCount 示例,它会从内置文本中读取数据并输出统计结果,这些结果通常会打印到 TaskManager 的标准输出中。使用 docker logs -f flink-taskmanager 可以持续跟踪输出,方便调试。

状态持久化与容器网络细节

Flink 作业的状态默认存储在 TaskManager 的内存中,一旦容器重启就会丢失。对于需要保存检查点和保存点的场景,必须配置持久化存储。常见的做法是将 state.checkpoints.dir 指向一个共享文件系统,例如挂载到所有容器内的卷或对象存储。Docker Compose 中可以定义一个命名卷,并同时挂载到 JobManager 和 TaskManager 的 /opt/flink/checkpoints 目录,然后在 flink-conf.yaml 中设置该目录为检查点路径。

容器网络方面,默认的桥接网络已经能够满足 JobManager 与 TaskManager 之间的通信。但如果作业需要访问宿主机上的服务,例如本地的 Kafka 或 Redis,不能使用 localhost,因为容器内的 localhost 指向容器自身。此时可以使用 host.docker.internal 作为宿主机地址,或者把外部服务也放入同一 Docker 网络。需要注意的是,在 Linux 系统上 host.docker.internal 需要 Docker 版本 20.10 以上且手动添加 extra_hosts 配置。

日志查看是排查问题的关键。容器化之后,日志不再写入宿主机上的 log 目录,而是通过容器的标准输出捕获。可以使用 docker compose logs -f 同时查看两个服务的日志,也可以单独指定服务名。如果日志量较大,建议将日志驱动配置为 json-file 并设置大小轮转,或者直接将 Flink 的日志目录挂载出来,方便使用宿主机的日志分析工具处理。

镜像版本选择同样值得注意。Flink 官方镜像同时提供 flink:latest 和带版本号的标签,生产环境或需要重现结果时应避免使用 latest。另外,不同 Scala 版本的镜像内部依赖不同,如果作业使用 Scala 2.12 编写,却选择了 scala_2.11 标签的镜像,可能会在运行时抛出序列化相关的异常。选定镜像后,可以通过 docker inspect flink:1.18.1-scala_2.12 查看镜像的详细信息和环境变量,确认与作业兼容。

总结来说,Docker 与 Apache Flink 的集成主要解决环境一致性和快速部署问题。通过 Compose 编排 JobManager 与 TaskManager,配合正确的网络配置和状态目录挂载,可以在几分钟内获得一个可用的流处理集群。后续如果需要迁移到 Kubernetes,也可以参考相同的容器化思路,将 Compose 服务定义转换为 Deployment 和 Service 资源,实现更高级别的调度与弹性伸缩。

DockerApache Flink流处理修改时间:2026-10-01 21:57:11

免责声明:已尽一切努力确保本网站所含信息的准确性。网站作品多为原创整理与精心创作,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们进行处理Email:chomcom@qq.com。
引用或转载本作品时,请注明当前出处:https://www.ipipp.com/html/1001/64402.html,基于非商业用途的前提下,欢迎转载或二创本作品。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。