Flink 2.3.0 从理论到实践 —— 第 3 章 环境准备与集群部署
课程定位:本章动手搭建 Flink 2.3.0 运行环境,从本地模式到 Standalone、YARN、Kubernetes 三种部署形态,覆盖
flink-conf.yaml核心参数、Web UI 与 REST API 使用,以及与本系列后续章节配套的 Paimon / Doris / Kafka / MySQL 连接器准备。所有命令与配置可直接复用于生产。版本基线:Flink 2.3.0 + JDK 17 + Paimon 1.4.x + Doris 4.1 + Kafka 3.x + MySQL 8.0/8.4
章节导读
- 3.1 版本与依赖准备
- 3.2 本地模式安装与冒烟测试
- 3.3 Standalone 集群部署
- 3.4 YARN 部署(Session / Per-Job / Application)
- 3.5 Kubernetes 部署(Flink Operator)
- 3.6 核心配置:flink-conf.yaml
- 3.7 Web UI 与 REST API
- 3.8 连接器与依赖准备
- 3.9 本章小结与下章预告
3.1 版本与依赖准备
3.1.1 版本兼容矩阵
本系列涉及多组件协同,版本组合是踩坑重灾区。下表是实测可用的版本基线:
| 组件 | 版本 | 说明 |
|---|---|---|
| Flink | 2.3.0 | 必须 JDK 17+ |
| JDK | 17 (LTS) | 推荐 Eclipse Temurin / OpenJDK 17 |
| Paimon | 1.4.x | 使用paimon-flink-2.3-1.4.x.jar |
| Doris | 4.1 | Doris Connector 适配 Flink 2.3 |
| Kafka | 3.x | flink-connector-kafka内置 |
| MySQL | 8.0 / 8.4 | flink-connector-jdbc+ MySQL Driver |
| Flink CDC | 3.6.x | 适配 Flink 2.3 |
| Hadoop | 3.3+ | YARN 部署时需要 |
关键提醒:
- Flink 2.x 已移除
addSink(SinkFunction)API,必须使用sinkTo(Sink)(影响 DataStream 程序写法)- Doris 4.1 的 Arrow 协议对 DATE 类型存在时区兼容问题(详见第 15 章),读 Doris 推荐用 JDBC Connector
- Paimon 1.4 的 Flink Connector JAR 必须严格匹配 Flink 大版本(
paimon-flink-2.3),用错版本会报not a subtype
3.1.2 硬件与操作系统要求
| 维度 | 最低 | 推荐 | 生产 |
|---|---|---|---|
| CPU | 4 核 | 8 核 | 16+ 核 |
| 内存 | 8 GB | 16 GB | 32+ GB |
| 磁盘 | 50 GB | 200 GB SSD | 1+ TB SSD + 冷数据 HDD |
| 网络 | 千兆 | 万兆 | 万兆 + 专用存储网 |
| OS | Linux 7+ | Rocky Linux 9 / CentOS 8 | 同推荐 |
| 文件描述符 | 65536 | 65536 | 65536+ |
3.1.3 系统参数调优
生产环境部署前需调整以下系统参数(一键执行脚本):
#!/bin/bash# flink-system-tune.sh - Flink 生产环境系统调优# 1. 文件描述符cat>>/etc/security/limits.conf<<EOF * soft nofile 65536 * hard nofile 65536 * soft nproc 65536 * hard nproc 65536 EOF# 2. 关闭 swap(避免被换出导致 GC 抖动)swapoff-ased-i'/swap/s/^/#/'/etc/fstab# 3. 关闭 THP(Transparent Huge Pages)echonever>/sys/kernel/mm/transparent_hugepage/enabledechonever>/sys/kernel/mm/transparent_hugepage/defrag# 4. 调整虚拟内存参数cat>>/etc/sysctl.conf<<EOF vm.max_map_count=262144 vm.swappiness=0 net.ipv4.tcp_retries2=5 net.core.somaxconn=32768 EOFsysctl-p# 5. NTP 时间同步(必装!Checkpoint 强依赖时钟)yuminstall-ychrony systemctlenablechronyd&&systemctl start chronyd3.2 本地模式安装与冒烟测试
本地模式用于开发与调试,单进程启动 JM + TM。
3.2.1 下载与解压
# 下载 Flink 2.3.0wgethttps://archive.apache.org/dist/flink/flink-2.3.0/flink-2.3.0-bin-scala_2.12.tgz# 解压tar-xzfflink-2.3.0-bin-scala_2.12.tgzcdflink-2.3.0# 查看目录结构ls-la3.2.2 JDK 17 安装与验证
# 安装 OpenJDK 17yuminstall-yjava-17-openjdk java-17-openjdk-devel# 验证java-version# openjdk version "17.0.x" 2026-xx-xx# 设置 JAVA_HOMEecho'export JAVA_HOME=/usr/lib/jvm/java-17-openjdk'>>/etc/profileecho'export PATH=$JAVA_HOME/bin:$PATH'>>/etc/profilesource/etc/profile3.2.3 启动与冒烟测试
# 启动本地集群./bin/start-cluster.sh# Starting cluster.# Starting standalonesession daemon on host xxx.# Starting taskexecutor daemon on host xxx.# 查看进程jps# 12345 StandaloneSession# 12346 TaskManagerRunner# 访问 Web UI: http://localhost:8081# 提交 WordCount 示例作业./bin/flink run examples/streaming/WordCount.jar# 查看作业输出tail-flog/flink-*-taskexecutor-*.out3.2.4 停止集群
./bin/stop-cluster.sh3.3 Standalone 集群部署
Standalone 模式是 Flink 自带的集群模式,不依赖 YARN 或 K8s,适合小规模生产或测试环境。
3.3.1 部署架构
┌──────────────────────────────────────────────────────────────┐ │ Standalone 集群架构 │ └──────────────────────────────────────────────────────────────┘ Master 节点 Worker 节点 ┌──────────────────┐ ┌──────────────────┐ │ JobManager │ 心跳 + 任务 │ TaskManager 1 │ │ (StandaloneSession)│ ◄──────► │ 4 Slots │ └──────────────────┘ └──────────────────┘ │ ┌──────────────────┐ │ TaskManager 2 │ │ 4 Slots │ └──────────────────┘3.3.2 配置masters与workers
# masters 文件:JM 节点cat>conf/masters<<EOF flink-master:8081 EOF# workers 文件:TM 节点列表cat>conf/workers<<EOF flink-worker1 flink-worker2 EOF3.3.3 SSH 免密登录
# Master 节点生成密钥ssh-keygen-trsa-b4096# 分发到所有节点ssh-copy-id flink-master ssh-copy-id flink-worker1 ssh-copy-id flink-worker23.3.4 分发与启动
# 分发 Flink 安装包到所有节点forhostinflink-worker1 flink-worker2;doscp-rflink-2.3.0$host:/opt/done# 在 Master 节点启动集群./bin/start-cluster.sh# 验证:所有节点的 TM 进程forhostinflink-master flink-worker1 flink-worker2;dossh$hostjps|grepTaskManagerRunnerdone3.4 YARN 部署(Session / Per-Job / Application)
YARN 是 Flink 在 Hadoop 生态中最常用的部署模式,支持三种形态。
3.4.1 三种模式对比
| 模式 | 特点 | 适用场景 | Flink 2.x 状态 |
|---|---|---|---|
| Session | 共享 JM,多作业复用 | 短作业、开发测试 | 推荐 |
| Per-Job | 每作业独占 JM + TM | 长作业隔离 | 已废弃(2.x 不再支持) |
| Application | 每应用独占,main 在 JM 运行 | 生产长作业 | 推荐 |
重要变更:Flink 2.x 已废弃 Per-Job 模式,推荐使用 Application 模式替代。原因:Per-Job 模式的 JM 与客户端在同一进程,客户端故障会导致作业丢失。
3.4.2 Session 模式部署
# 启动 YARN Session./bin/yarn-session.sh\-d\# 后台运行-nmflink-session\# YARN 应用名-qudefault\# 队列-tm4096\# 每个 TM 内存-s4\# 每个 TM 的 Slot 数-jm2048# JM 内存# 输出会包含 YARN Application ID,后续提交作业时使用# JobManager Web Interface: http://xxx:8081提交作业到 Session:
# 方式 1: 命令行(自动发现 YARN Session)./bin/flink run\-d\-p4\# 并行度-yidapplication_123_456\# 显式指定 YARN Application IDexamples/streaming/WordCount.jar# 方式 2: SQL Client(指定 Session)./bin/sql-client.sh embedded\-Dyarn.application.id=application_123_4563.4.3 Application 模式部署
# Application 模式:main 函数在 JM 运行./bin/flink run-application\-d\-tyarn-application\-Djobmanager.memory.process.size=2048m\-Dtaskmanager.memory.process.size=4096m\-Dtaskmanager.numberOfTaskSlots=4\-Dpipeline.name=my-flink-app\# 作业名(WebUI 显示)-ccom.example.MyFlinkJob\/path/to/my-flink-job.jar项目硬约束:Flink 任务必须用
-D pipeline.name命名,便于在 WebUI 与 YARN ResourceManager 中识别作业。
3.4.4 YARN 部署两铁律
基于项目实战经验,YARN 部署有两铁律(违反即事故):
| 铁律 | 命令 | 背景 |
|---|---|---|
① 显式yarn.application.id | -Dyarn.application.id=<appId> | /tmp/.yarn-properties-<user>自动发现会指向"最近一次"的 Session,可能把作业投到错误的 Session 抢 slot |
②table.dml-sync=true | -Dtable.dml-sync=true | 默认false时 INSERT 提交后客户端立即返回,调度平台会把"提交成功"当"跑成功" |
# 标准提交模板(生产推荐)./bin/flink run-application\-d\-tyarn-application\-Dyarn.application.id=${YARN_APP_ID}\-Dtable.dml-sync=true\-Dpipeline.name=${JOB_NAME}\-ccom.example.MyFlinkJob\${JAR_PATH}3.5 Kubernetes 部署(Flink Operator)
Kubernetes 是 Flink 2.x 推荐的生产部署形态,通过Flink Kubernetes Operator实现声明式部署。
3.5.1 架构
┌──────────────────────────────────────────────────────────────┐ │ Flink Kubernetes Operator │ └──────────────────────────────────────────────────────────────┘ ┌──────────────────┐ │ K8s API Server │ └────────┬─────────┘ │ 监听 FlinkDeployment CRD ▼ ┌──────────────────┐ │ Flink Operator │ ◄── 声明式管理 Flink 集群生命周期 └────────┬─────────┘ │ ┌──────┼──────┐ ▼ ▼ ▼ JM Pod TM Pod 1 TM Pod 23.5.2 安装 Operator
# 添加 Helm 仓库helm repoaddflink-operator https://downloads.apache.org/flink/flink-kubernetes-operator-1.0.0/ helm repo update# 安装 Operatorhelminstallflink-operator flink-operator/flink-kubernetes-operator3.5.3 部署 Flink 作业
# flink-deployment.yamlapiVersion:flink.apache.org/v1beta1kind:FlinkDeploymentmetadata:name:flink-jobnamespace:flinkspec:image:flink:2.3.0flinkVersion:v2_3serviceAccount:flinkjobManager:replicas:1resource:memory:2048mcpu:1taskManager:replicas:2resource:memory:4096mcpu:2job:jarURI:local:///opt/flink/usrlib/my-job.jarparallelism:4entryClass:com.example.MyFlinkJobargs:["--config","/etc/config/app.properties"]flinkConfiguration:pipeline.name:my-flink-jobexecution.checkpointing.interval:30sstate.checkpoints.dir:s3://my-bucket/checkpoints/kubectl apply-fflink-deployment.yaml3.5.4 K8s HA(无 ZooKeeper)
spec:flinkConfiguration:kubernetes.cluster-id:flink-jobhigh-availability:org.apache.flink.kubernetes.highavailability.KubernetesHaServicesFactoryhigh-availability.storageDir:s3://my-bucket/ha/# 不依赖 ZooKeeper3.6 核心配置:flink-conf.yaml
flink-conf.yaml是 Flink 的核心配置文件,位于conf/目录。
3.6.1 基础配置
# ==================== 基础配置 ====================jobmanager.rpc.address:flink-masterjobmanager.rpc.port:6123jobmanager.bind-host:0.0.0.0taskmanager.bind-host:0.0.0.0taskmanager.host:#{自动解析}# Web UIrest.address:0.0.0.0rest.port:8081# 默认并行度parallelism.default:43.6.2 内存配置
# ==================== JobManager 内存 ====================jobmanager.memory.process.size:2048mjobmanager.memory.flink.size:1024m# Flink 内存# 其余为 JVM Metaspace + JVM Overhead# ==================== TaskManager 内存 ====================taskmanager.memory.process.size:4096mtaskmanager.memory.flink.size:3072m# Task 堆内存(算子状态、用户对象)taskmanager.memory.task.heap.size:1024m# 托管内存(RocksDB/ForSt 状态、排序、批处理)taskmanager.memory.managed.size:1024m# 网络缓冲(TM 间数据传输)taskmanager.memory.network.min:256mbtaskmanager.memory.network.max:512mbtaskmanager.memory.network.fraction:0.1# Slot 数taskmanager.numberOfTaskSlots:43.6.3 Checkpoint 配置
# ==================== Checkpoint ====================execution.checkpointing.interval:30sexecution.checkpointing.timeout:10minexecution.checkpointing.min-pause:30sexecution.checkpointing.max-concurrent-checkpoints:1execution.checkpointing.externalized-checkpoint-retention:RETAIN_ON_CANCELLATION# 状态后端(Flink 2.x 默认 HashMap,推荐 ForSt)state.backend:forststate.checkpoints-dir:hdfs:///flink/checkpoints/state.savepoints-dir:hdfs:///flink/savepoints/# 增量 Checkpoint(仅 ForSt/RocksDB 支持)state.backend.incremental:true3.6.4 容错与重启
# ==================== 重启策略 ====================restart-strategy.type:fixed-delayrestart-strategy.fixed-delay.attempts:2147483647# 无限重试restart-strategy.fixed-delay.delay:10s# 或使用失败率策略# restart-strategy.type: failure-rate# restart-strategy.failure-rate.max-failures-per-interval: 3# restart-strategy.failure-rate.failure-rate-interval: 5min# restart-strategy.failure-rate.delay: 10s3.6.5 Flink 2.x 特定配置
# ==================== Flink 2.x 特性 ====================# 执行模式execution.runtime-mode:streaming# 自适应调度器(2.x 默认)jobmanager.scheduler:adaptive# Watermark 对齐(2.3 新特性)# 在作业级别启用,详见第 5 章# pipeline.watermark-alignment.enabled: true# Unaligned Checkpoint(背压场景)execution.checkpointing.unaligned.enabled:false# 默认关闭,按需开启3.7 Web UI 与 REST API
Flink Web UI 是作业监控的入口,默认端口8081。
3.7.1 Web UI 功能
| 页面 | 功能 |
|---|---|
| Overview | 集群总览:可用 Slot、运行作业数、已完成作业数 |
| Running Jobs | 运行中作业列表,点击查看作业详情 |
| Job Detail | 作业详情:DAG 图、算子并行度、背压、Checkpoint 统计 |
| Task Managers | TM 列表:内存、Slot、线程栈 |
| JobManager | JM 状态:内存、配置、日志 |
3.7.2 REST API 常用接口
# 集群总览curlhttp://flink-master:8081/overview# {# "taskmanagers": 2,# "slots-total": 8,# "slots-available": 4,# "jobs-running": 1,# "flink-version": "2.3.0"# }# 作业列表curlhttp://flink-master:8081/jobs/overview# 作业详情(含异常)curlhttp://flink-master:8081/jobs/<job-id>/exceptions# 作业执行计划(DAG)curlhttp://flink-master:8081/jobs/<job-id>/plan# 取消作业curl-XPATCH http://flink-master:8081/jobs/<job-id># 触发 Savepointcurl-XPOST http://flink-master:8081/jobs/<job-id>/savepoints\-H"Content-Type: application/json"\-d'{"target-directory": "hdfs:///flink/savepoints/"}'实战技巧:REST API 是自动化运维的基础。在生产中,可以通过 REST API 定期采集作业状态、背压指标、Checkpoint 时长,构建监控告警系统(详见第 16 章)。
3.8 连接器与依赖准备
本系列后续章节会用到以下连接器,建议在部署阶段一次性准备好放入${FLINK_HOME}/lib/。
3.8.1 连接器清单
| 连接器 | JAR | 用途 |
|---|---|---|
| Kafka | flink-sql-connector-kafka-2.3.0.jar | Kafka 读写(内置,无需额外) |
| JDBC | flink-sql-connector-jdbc-2.3.0.jar | MySQL 读写 |
| MySQL Driver | mysql-connector-java-8.0.33.jar | JDBC 驱动 |
| Paimon | paimon-flink-2.3-1.4.x.jar | Paimon 湖存储 |
| Doris | doris-flink-connector-1.0.3.jar | Doris Sink |
| Flink CDC | flink-sql-connector-mysql-cdc-3.6.x.jar | MySQL CDC |
3.8.2 下载脚本
#!/bin/bash# download-connectors.sh - 下载 Flink 2.3 连接器cd${FLINK_HOME}/lib# Kafka (内置,无需下载)# ls -la flink-sql-connector-kafka-*.jar# JDBC Connectorwgethttps://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-jdbc/2.3.0/flink-sql-connector-jdbc-2.3.0.jar# MySQL Driverwgethttps://repo1.maven.org/maven2/mysql/mysql-connector-java/8.0.33/mysql-connector-java-8.0.33.jar# Paimon (1.4.x)wgethttps://repo1.maven.org/maven2/org/apache/paimon/paimon-flink-2.3/1.4.0/paimon-flink-2.3-1.4.0.jar# Doris Connector (适配 Flink 1.20+,Flink 2.x 兼容)wgethttps://repo1.maven.org/maven2/org/apache/doris/flink-doris-connector-1.0.3/flink-doris-connector-1.0.3.jar# Flink CDC (MySQL)wgethttps://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-mysql-cdc/3.6.0/flink-sql-connector-mysql-cdc-3.6.0.jar3.8.3 jar 加载铁律
关键约束(来自项目实测):所有连接器 JAR必须只放在
${FLINK_HOME}/lib/,提交作业时绝不传-j参数。
违规后果:-j与 lib 双份加载同一连接器类,导致双 ClassLoader 加载,建 Catalog 时报:
ServiceConfigurationError: org.apache.paimon.factories.Factory not a subtype自检命令:
# 确认每个连接器 JAR 只有一份find${FLINK_HOME}/lib-name"paimon-flink-*.jar"|wc-l# 应为 1find${FLINK_HOME}/lib-name"flink-doris-*.jar"|wc-l# 应为 1# 确认无重复类forjarin${FLINK_HOME}/lib/*.jar;doecho"===$jar==="unzip-l$jar|grep-c'\.class$'done3.9 本章小结与下章预告
本章小结
┌────────────────────────────────────────────────────────────────┐ │ 第 3 章 要点回顾 │ └────────────────────────────────────────────────────────────────┘ ✓ 版本基线: Flink 2.3.0 + JDK 17 + Paimon 1.4 + Doris 4.1 注意: Flink 2.x 已移除 addSink(SinkFunction),必须用 sinkTo(Sink) ✓ 部署形态: 本地模式(开发) → Standalone(小规模) → YARN/K8s(生产) Flink 2.x 已废弃 Per-Job 模式,推荐 Application 模式 ✓ YARN 部署两铁律: ① 显式 -Dyarn.application.id ② -Dtable.dml-sync=true ✓ 核心配置 flink-conf.yaml: 内存: process.size / task.heap / managed / network Checkpoint: interval 30s / ForSt 后端 / 增量 重启: fixed-delay 无限重试 ✓ K8s 部署: Flink Kubernetes Operator(声明式) K8s HA(不依赖 ZooKeeper) ✓ Web UI: http://<host>:8081 REST API: /overview /jobs /jobs/<id> 支持 ✓ 连接器铁律: 只放 ${FLINK_HOME}/lib,绝不传 -j 否则: not a subtype(双 ClassLoader)下章预告
第 4 章 状态管理(State):深入 Flink 的有状态计算核心。讲解 Keyed State vs Operator State、状态后端(HashMap / ForSt)、状态 TTL 与过期策略、状态大小估算与监控、状态迁移与 Schema Evolution。状态是 Flink 区别于普通流处理引擎的关键能力,也是后续 Checkpoint 容错的基础。
官方参考资料
- Flink 部署概览:https://nightlies.apache.org/flink/flink-docs-stable/docs/deployment/overview/
- YARN 部署:https://nightlies.apache.org/flink/flink-docs-stable/docs/deployment/resource-providers/yarn/
- Kubernetes Operator:https://nightlies.apache.org/flink/flink-docs-stable/docs/deployment/resource-providers/kubernetes/
- 配置参数:https://nightlies.apache.org/flink/flink-docs-stable/docs/deployment/config/