ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

使用Docker-compose快速部署Apache Flink集群:从环境搭建到生产调优

2026/8/13 6:16:59 拓冰建站 浏览量
使用Docker-compose快速部署Apache Flink集群:从环境搭建到生产调优 1. 从单体部署到容器编排为什么选择Docker-compose部署Flink如果你正在处理实时数据流无论是电商的实时推荐、物联网的设备状态监控还是金融交易的风控分析Apache Flink 大概率已经进入了你的技术选型视野。作为一个强大的流处理框架Flink 以其高吞吐、低延迟和精确一次Exactly-Once的状态一致性保证成为了实时计算领域的核心引擎之一。然而当我们从开发测试走向生产部署时一个现实的问题就摆在了面前如何高效、一致地管理 Flink 集群的各个组件特别是对于中小型团队或项目初期直接上马 Kubernetes 可能显得过于沉重而手动部署 JobManager、TaskManager 又容易陷入配置繁琐、环境不一致的泥潭。这正是 Docker-compose 可以大显身手的地方。它不是一个生产级的集群编排工具但对于搭建一个功能完整、可用于开发、测试甚至小规模生产的 Flink 集群原型来说它几乎是完美的选择。想象一下你只需要一个docker-compose.yml文件就能一键拉起包含 JobManager、TaskManager、甚至 Web UI 的完整集群并且能确保在任何一台安装了 Docker 的机器上集群的行为完全一致。这极大地简化了环境搭建的复杂度让开发者能更专注于 Flink 应用逻辑本身而不是纠结于端口冲突、依赖缺失或者配置文件路径错误。我经历过手动部署 Flink 的“痛苦”需要分别启动 JobManager 和多个 TaskManager 进程管理它们的日志处理网络互通每次换台机器都要重新检查一遍。而转向 Docker-compose 后整个部署过程变成了一个可版本化、可重复的“配方”。无论是新同事加入快速搭建环境还是需要在本地复现一个线上问题这个docker-compose.yml文件就是最可靠的蓝图。接下来我将带你从零开始一步步拆解如何用 Docker-compose 部署一个功能完备的 Flink 集群并深入其中几个关键配置背后的逻辑以及我在实际使用中积累的一些避坑经验。2. 环境准备与核心镜像选择不只是拉取镜像那么简单在动手编写docker-compose.yml之前我们需要确保基础环境就绪并做出第一个关键决策选择哪个 Flink 镜像。2.1 基础环境检查与安装首先你需要确保你的机器上已经安装了 Docker 和 Docker-compose。对于 Linux 系统可以通过包管理器安装。这里以 Ubuntu 为例但原理相通# 安装 Docker sudo apt-get update sudo apt-get install docker.io sudo systemctl start docker sudo systemctl enable docker # 安装 Docker-compose # 注意较新版本的 Docker Desktop 已包含 compose 插件可通过 docker compose 命令使用。 # 如需独立安装可下载特定版本 sudo curl -L https://github.com/docker/compose/releases/download/v2.23.0/docker-compose-$(uname -s)-$(uname -m) -o /usr/local/bin/docker-compose sudo chmod x /usr/local/bin/docker-compose注意生产环境建议使用特定版本而非latest标签以保证稳定性。同时确保当前用户拥有执行 Docker 命令的权限通常需要加入docker用户组。2.2 Flink 官方镜像的版本与变体选择访问 Docker Hub 上的flink镜像仓库你会发现有多个标签。选择哪一个直接决定了你集群的基础特性。主要分为两大类Scala 版本如1.17.2-scala_2.12。Flink 本身是用 Java 编写的但其 API 为 Scala 也提供了支持。如果你的作业是用 Scala 编写的或者依赖的某些连接器Connector需要特定 Scala 版本就必须选择对应的 Scala 变体。2.12是目前最主流和稳定的 Scala 版本。Java 版本如1.17.2-java11。从 Flink 1.15 开始官方推荐使用 Java 11 或更高版本。Java 8 镜像已逐渐被弃用。选择与你的开发环境和依赖兼容的 Java 版本。对于大多数使用 Java API 或 Flink SQL 的用户选择flink:1.17.2-java11这样的标签就足够了。它是最通用、问题最少的版本。我个人的经验是除非有强制的 Scala 依赖否则优先选择纯 Java 版本可以减少因 Scala 版本冲突带来的潜在麻烦。此外镜像还分-slim和普通版本。-slim版本体积更小但可能缺少一些调试工具如telnet,vim。对于生产倾向的部署普通版本更稳妥。对于本次部署我们选择flink:1.17.2-java11。2.3 网络规划容器间通信的基石Docker-compose 默认会为所有服务创建一个独立的网络服务间可以使用服务名作为主机名互相访问。这非常适合 Flink 集群JobManager 需要知道 TaskManager 的地址来分发任务TaskManager 需要向 JobManager 注册心跳。我们不需要手动创建网络Docker-compose 会处理好。但需要理解在docker-compose.yml中定义的服务名如jobmanager,taskmanager在容器内部就是有效的主机名。例如TaskManager 的配置中jobmanager.rpc.address就可以直接设置为jobmanager。3. 编写 docker-compose.yml逐行解析集群定义这是最核心的部分。我们将创建一个docker-compose.yml文件定义一个包含一个 JobManager、两个 TaskManager 的集群。我会对每个关键配置进行解释。version: 2.1 # 使用 2.1 或更高版本以支持健康检查等特性 services: jobmanager: image: flink:1.17.2-java11 container_name: flink-jobmanager hostname: jobmanager ports: - 8081:8081 # Flink Web UI 端口 - 6123:6123 # JobManager RPC 端口用于客户端提交作业 expose: - 6123 command: jobmanager environment: - | FLINK_PROPERTIES jobmanager.rpc.address: jobmanager taskmanager.numberOfTaskSlots: 2 parallelism.default: 1 state.backend: filesystem state.checkpoints.dir: file:///opt/flink/checkpoints state.savepoints.dir: file:///opt/flink/savepoints volumes: - ./checkpoints:/opt/flink/checkpoints - ./savepoints:/opt/flink/savepoints - ./job-artifacts:/opt/flink/job-artifacts healthcheck: test: [CMD, curl, -f, http://localhost:8081] interval: 30s timeout: 10s retries: 3 start_period: 60s taskmanager: image: flink:1.17.2-java11 container_name: flink-taskmanager-1 hostname: taskmanager-1 depends_on: jobmanager: condition: service_healthy # 等待 JobManager 健康后再启动 expose: - 6121 - 6122 command: taskmanager environment: - | FLINK_PROPERTIES jobmanager.rpc.address: jobmanager taskmanager.numberOfTaskSlots: 2 volumes: - ./checkpoints:/opt/flink/checkpoints - ./savepoints:/opt/flink/savepoints - ./job-artifacts:/opt/flink/job-artifacts deploy: replicas: 2 # 启动两个 TaskManager 实例 healthcheck: test: [CMD, curl, -f, http://localhost:8081] interval: 30s timeout: 10s retries: 3 start_period: 60s现在我们来拆解这个配置文件的关键部分版本与服务定义version: 2.1确保了我们对健康检查等功能的支持。在services下我们定义了两个服务jobmanager和taskmanager。注意taskmanager服务通过deploy.replicas: 2启动了2个实例Docker-compose 会为它们生成不同的容器名如flink-taskmanager-1,flink-taskmanager-2但主机名需要特殊处理见下文。网络与主机名我们没有显式定义网络因此 Docker-compose 会使用默认网络。jobmanager容器的主机名被设置为jobmanager。对于taskmanager这里有个关键技巧由于我们使用了replicas每个副本都会有相同的配置。如果都设置相同的hostname会导致冲突。因此上面的配置中taskmanager服务的hostname: taskmanager-1只对第一个副本生效。实际上在 Docker-compose v3 中更推荐的做法是不设置hostname让 Docker 自动分配然后在 Flink 配置里使用服务名taskmanager进行通信因为 Flink 的 TaskManager 动态注册机制不依赖固定的主机名。但为了清晰我们也可以在命令中动态设置不过这会增加复杂度。对于入门部署使用服务名通信是最简单的。端口映射我们将宿主机的8081端口映射到 JobManager 容器的8081端口这样就能通过http://localhost:8081访问 Flink 的 Web 仪表盘。6123端口是 JobManager 的 RPC 端口用于接收flink run命令提交的作业。expose指令声明容器内部暴露的端口供其他服务访问但不映射到宿主机。环境变量与 Flink 配置这是核心。我们通过FLINK_PROPERTIES环境变量来覆盖 Flink 的默认配置conf/flink-conf.yaml。这里采用了 YAML 的多行字符串格式|。jobmanager.rpc.address: jobmanager告诉 TaskManagerJobManager 的地址是服务名jobmanager。taskmanager.numberOfTaskSlots: 2每个 TaskManager 提供 2 个任务槽Task Slot。一个 Slot 是资源调度的基本单位可以运行一个算子子任务。假设你启动2个 TaskManager集群总 Slot 数就是4。parallelism.default: 1作业的默认并行度。提交作业时如果不指定就使用这个值。state.backend: filesystem状态后端设置为文件系统。这是最简单的后端将状态快照Checkpoint/Savepoint保存到磁盘。对于生产环境通常会考虑rocksdb更高效或配置外部存储如 HDFS, S3。state.checkpoints.dir和state.savepoints.dir分别指定 Checkpoint 和 Savepoint 的存储路径。我们将其挂载到宿主机实现数据持久化。数据卷挂载通过volumes将宿主机的目录./checkpoints,./savepoints,./job-artifacts挂载到容器内的固定路径。这样做有两个巨大好处一是数据不会随着容器销毁而丢失二是方便我们在宿主机上查看和管理 Checkpoint/Savepoint 文件或者预先放置需要提交的作业 JAR 包。健康检查healthcheck配置让 Docker 可以感知服务的健康状态。这里使用curl检查 Web UI 端口是否可达。depends_on中的condition: service_healthy确保了 TaskManager 会等待 JobManager 完全启动就绪后才启动避免了启动顺序问题导致的连接失败。这是一个非常实用的稳定性增强配置。4. 启动集群、提交作业与日常操作实战配置文件就绪后我们就可以操作这个容器化的 Flink 集群了。4.1 启动与停止集群在包含docker-compose.yml的目录下执行以下命令# 启动集群后台运行 docker-compose up -d # 查看集群运行状态 docker-compose ps # 查看 JobManager 的日志实时跟踪 docker-compose logs -f jobmanager # 查看特定 TaskManager 的日志 docker-compose logs -f flink-taskmanager-1 # 停止并移除集群会删除容器 docker-compose down # 停止并移除集群同时删除数据卷慎用会丢失 Checkpoint 数据 docker-compose down -v启动后打开浏览器访问http://localhost:8081你应该能看到 Flink 的 Web UI。在 “Task Managers” 标签页下应该能看到两个已注册的 TaskManager每个提供 2 个 Slot总共 4 个 Slot。4.2 提交作业的几种方式作业如何提交到容器内的集群这里提供三种最常用的方法方法一通过 Web UI 提交这是最简单直观的方式。在 Web UI 的 “Submit New Job” 页面直接上传你的作业 JAR 包并填写入口类名和参数即可。但是JAR 包需要在你本地浏览器可访问的位置。对于容器环境更推荐下面两种方式。方法二使用flink run命令提交推荐这是最标准的方式。你需要进入 JobManager 容器内部执行命令。# 1. 将你的作业 JAR 包复制到共享的挂载目录 cp your-flink-job.jar ./job-artifacts/ # 2. 进入 JobManager 容器 docker-compose exec jobmanager bash # 3. 在容器内部使用 flink run 提交作业 # 注意 JAR 包路径是容器内的挂载路径 ./bin/flink run /opt/flink/job-artifacts/your-flink-job.jar --input topic1 --output topic2 # 4. 提交后退出容器 exit方法三通过 REST API 提交Flink 提供了 REST API允许你通过 HTTP 请求提交作业。这便于自动化脚本集成。# 假设 JAR 包已在 ./job-artifacts/ 目录下 JAR_FILEyour-flink-job.jar JAR_PATH_ON_HOST./job-artifacts/${JAR_FILE} # 使用 curl 通过 REST API 提交 # 首先上传 JAR 包到集群 UPLOAD_RESPONSE$(curl -X POST -H Expect: -F jarfile${JAR_PATH_ON_HOST} http://localhost:8081/jars/upload) # 从响应中提取 jarid具体解析取决于响应格式通常是 JSON JAR_ID$(echo $UPLOAD_RESPONSE | grep -oP filename:${JAR_FILE},id:\K[^]) # 然后触发 JAR 包中作业的执行 curl -X POST http://localhost:8081/jars/${JAR_ID}/run?entry-classcom.example.YourMainClassprogram-args--input%20topic1%20--output%20topic2注意REST API 提交方式需要仔细处理响应和错误码对于复杂参数方法二更直接可靠。4.3 管理作业状态Checkpoint 与 Savepoint在 Web UI 的 “Running Jobs” 或 “Completed Jobs” 页面你可以管理作业。触发 Savepoint可以对运行中的作业手动触发 Savepoint用于有状态的作业升级或迁移。从 Savepoint 恢复提交新作业时可以通过-s参数指定一个 Savepoint 路径作业会从该状态恢复。查看 Checkpoint在作业详情页可以查看 Checkpoint 的历史记录、配置和统计信息这是监控作业稳定性的重要依据。由于我们将 Checkpoint/Savepoint 目录挂载到了宿主机你可以在./checkpoints和./savepoints目录下找到对应的文件。一个重要经验定期清理旧的 Checkpoint 目录因为 Flink 默认不会自动清理长期运行可能占满磁盘。可以通过配置state.checkpoints.num-retained来保留最近 N 个 Checkpoint。5. 配置调优与生产就绪考量用 Docker-compose 能快速搭起集群但要让其更健壮、更适合准生产环境还需要一些调优。5.1 资源限制与调优默认情况下容器可以使用宿主机的所有资源。这可能导致单个容器耗尽资源影响其他服务。我们应该为容器设置资源限制。services: jobmanager: # ... 其他配置 ... deploy: resources: limits: memory: 2048M cpus: 1.0 reservations: memory: 1024M cpus: 0.5 taskmanager: # ... 其他配置 ... deploy: replicas: 2 resources: limits: memory: 4096M # 每个 TaskManager 内存限制 cpus: 2.0 reservations: memory: 2048M cpus: 1.0这里设置了内存和 CPU 的限制limits和预留reservations。同时你需要对应地调整 Flink 的配置使 Flink 感知到的内存与容器限制对齐否则可能因内存超出限制被 Docker 杀死。关键配置在FLINK_PROPERTIES中environment: - | FLINK_PROPERTIES jobmanager.memory.process.size: 1600m # 应小于容器内存限制 taskmanager.memory.process.size: 3600m # 应小于容器内存限制 taskmanager.memory.managed.size: 800m # 托管内存用于排序、哈希表等 taskmanager.numberOfTaskSlots: 2taskmanager.memory.process.size是 TaskManager 的总内存它必须小于 Docker 容器的内存限制为操作系统和其他进程留出余地。taskmanager.memory.managed.size是 Flink 管理的堆外内存用于缓存状态等根据作业特点调整。5.2 状态后端与高可用配置我们之前使用了filesystem状态后端它简单但不适合高可用场景因为状态文件在单个节点的本地磁盘。对于需要容错的生产环境应考虑RocksDB 状态后端更节省内存支持增量 Checkpoint适合大状态作业。配置state.backend: rocksdb并设置state.backend.rocksdb.localdir为一个持久化卷路径。外部化 Checkpoint 存储即使使用 filesystem也应配置为共享存储如 NFS、HDFS 或 S3这样 JobManager 故障恢复后还能找到 Checkpoint。配置state.checkpoints.dir: hdfs://namenode:8020/flink/checkpoints。高可用HA模式Docker-compose 部署单个 JobManager 存在单点故障。Flink 支持基于 ZooKeeper 的 HA 模式可以部署多个 JobManager 实例一个为主Leader其余为备。这需要引入 ZooKeeper 服务并配置high-availability相关参数。在 Docker-compose 中实现相对复杂通常这标志着需要向 Kubernetes 等更成熟的编排平台迁移了。5.3 日志与监控集成默认日志会输出到容器的标准输出可以通过docker-compose logs查看。为了持久化和集中管理可以将日志目录挂载出来或者配置日志框架如 log4j将日志发送到 ELKElasticsearch, Logstash, Kibana栈。监控方面Flink 提供了丰富的 Metrics可以对接 Prometheus 和 Grafana。在FLINK_PROPERTIES中启用 Prometheus Reportermetrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter metrics.reporter.prom.port: 9250-9260在docker-compose.yml中为 JobManager 和 TaskManager 暴露额外的端口范围如9250-9260:9250-9260或者使用expose。在同一个docker-compose.yml中添加 Prometheus 和 Grafana 服务配置 Prometheus 抓取 Flink 容器的 Metrics 端口。这样你就能在 Grafana 中创建丰富的仪表盘监控作业的吞吐量、延迟、背压、Checkpoint 时长等关键指标。6. 常见问题排查与实战经验分享即使配置得当在实际运行中也可能遇到问题。下面分享几个我踩过的坑和解决方法。6.1 TaskManager 无法注册到 JobManager现象Web UI 中看不到 TaskManager或者 TaskManager 日志中不断报连接拒绝的错误。排查思路检查网络确保docker-compose.yml中 JobManager 的服务名正确并且 TaskManager 的jobmanager.rpc.address配置指向了这个服务名。在 TaskManager 容器内执行ping jobmanager看是否能通。检查端口确认 JobManager 的 RPC 端口默认6123在容器网络内是暴露的expose并且没有被防火墙阻挡。检查启动顺序使用depends_on和healthcheck确保 TaskManager 在 JobManager 就绪后才启动。早期我忽略了这一点TaskManager 启动时 JobManager 的 RPC 服务还没起来导致注册失败。查看日志仔细查看 JobManager 和 TaskManager 的日志。docker-compose logs --tail50 jobmanager taskmanager可以快速查看最近日志。6.2 作业提交失败或卡住现象通过flink run或 Web UI 提交作业后作业一直处于CREATED或SCHEDULED状态不运行。排查思路检查资源最常见的原因是集群没有足够的 Slot。在 Web UI 的 “Task Managers” 页查看总 Slot 数在 “Running Jobs” 页查看作业申请的并行度。如果作业并行度或默认并行度大于可用 Slot 总数作业就无法调度。检查 JAR 包依赖如果作业 JAR 包缺少依赖如 Kafka 连接器TaskManager 在加载用户代码时会抛出ClassNotFoundException。确保使用maven-shade-plugin或maven-assembly-plugin打好包含所有依赖的 “uber jar”。或者在docker-compose.yml中将包含依赖的目录挂载到容器的lib/目录下不推荐易冲突。查看 JobManager 日志提交作业时的异常通常会在 JobManager 日志中体现。6.3 容器内内存不足导致进程被 Kill现象TaskManager 或 JobManager 容器突然消失docker-compose ps显示状态为Exited (137)或Exited (1)。137 通常表示因内存超限被系统终止OOM Killer。解决方案调整 Docker 资源限制如上文所述在docker-compose.yml中增加deploy.resources.limits.memory。调整 Flink 内存配置确保jobmanager.memory.process.size和taskmanager.memory.process.size的值小于Docker 容器的内存限制。建议预留至少 10%-20% 的内存给容器内的其他进程如 JVM 本身、Native 库。监控内存使用可以通过docker stats命令实时查看容器的内存和 CPU 使用情况辅助定位问题。6.4 状态后端路径权限问题现象作业可以运行但无法完成 Checkpoint日志显示IOException: Permission denied。原因与解决Docker 容器内的进程通常以非 root 用户如flink用户运行。如果挂载的宿主机目录如./checkpoints对容器用户不可写就会报错。方案一推荐在宿主机上确保挂载目录对 Docker 容器用户可写。一个简单粗暴但有效的方法是赋予 777 权限仅限开发环境chmod -R 777 ./checkpoints ./savepoints。方案二在docker-compose.yml中以 root 用户身份运行容器user: root但这会降低安全性不推荐。6.5 宿主机端口冲突现象执行docker-compose up时报错Bind for 0.0.0.0:8081 failed: port is already allocated。解决这意味着你宿主机上的 8081 端口已被其他进程占用。找到并停止占用端口的进程sudo lsof -i :8081。或者在docker-compose.yml中修改端口映射例如- 8082:8081然后通过http://localhost:8082访问 Web UI。7. 进阶集成外部系统与自定义镜像基本的集群运行起来后你可能需要连接 Kafka、MySQL、HDFS 等外部系统或者安装自定义的依赖包。7.1 连接 Kafka 作为 Source/Sink这是非常常见的场景。Flink 容器默认不包含 Kafka 连接器。有两种方式解决方式一将连接器 JAR 包放入挂载目录从 Maven 仓库下载 Flink Kafka 连接器 JAR 包如flink-connector-kafka-1.17.2.jar及其依赖如kafka-clients-xxx.jar。将这些 JAR 包放入宿主机./job-artifacts/目录或专门创建一个./lib/目录。在提交作业时通过-C或--classpath参数指定额外的 JAR 包路径比较麻烦。更简单的方法是在docker-compose.yml中将这个目录挂载到 Flink 容器的lib/目录下但要注意版本冲突。方式二构建自定义 Docker 镜像推荐这是更干净、可复用的方式。创建一个DockerfileFROM flink:1.17.2-java11 # 将 Kafka 连接器 jar 包添加到 Flink 的 lib 目录 # 注意下载的 jar 包需要与 Flink 版本兼容 COPY flink-connector-kafka-1.17.2.jar /opt/flink/lib/ COPY kafka-clients-3.4.0.jar /opt/flink/lib/ # 可以继续添加其他依赖如 JDBC 驱动 # COPY mysql-connector-java-8.0.33.jar /opt/flink/lib/ USER flink然后在docker-compose.yml中将image: flink:1.17.2-java11替换为build: .假设 Dockerfile 在当前目录。这样构建的镜像就自带了所需连接器。7.2 在 Flink SQL 中使用 Hive Catalog如果你想在 Flink SQL 中直接查询 Hive 表需要配置 Hive Catalog。准备 Hive 依赖将 Flink 的 Hive 连接器 JAR 包flink-sql-connector-hive-3.1.2_2.12-1.17.2.jar和 Hive 相关的依赖包放入自定义镜像的/opt/flink/lib/目录。配置 Hive Metastore在FLINK_PROPERTIES中增加配置并确保 Flink 容器能访问 Hive Metastore 服务可能需要将 Metastore 服务也定义在docker-compose.yml中或使用外部服务地址。在 SQL Client 或程序中创建 Catalog这步通常在你的作业代码或 SQL 脚本中完成。这个过程涉及较多细节但它展示了 Docker-compose 的灵活性你可以通过自定义镜像和网络配置将 Flink 集群与一整套大数据生态服务如 Kafka、Hive、HDFS集成在同一个编排文件中形成一个完整的、本地可用的实时数据处理微服务栈。虽然这离真正的生产环境还有距离但对于集成测试和概念验证PoC来说其价值是巨大的。