news 2026/10/10 2:55:58

Flink 2.3.0 从理论到实践 —— 第 3 章 环境准备与集群部署

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink 2.3.0 从理论到实践 —— 第 3 章 环境准备与集群部署

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 版本兼容矩阵

本系列涉及多组件协同,版本组合是踩坑重灾区。下表是实测可用的版本基线:

组件版本说明
Flink2.3.0必须 JDK 17+
JDK17 (LTS)推荐 Eclipse Temurin / OpenJDK 17
Paimon1.4.x使用paimon-flink-2.3-1.4.x.jar
Doris4.1Doris Connector 适配 Flink 2.3
Kafka3.xflink-connector-kafka内置
MySQL8.0 / 8.4flink-connector-jdbc+ MySQL Driver
Flink CDC3.6.x适配 Flink 2.3
Hadoop3.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 硬件与操作系统要求

维度最低推荐生产
CPU4 核8 核16+ 核
内存8 GB16 GB32+ GB
磁盘50 GB200 GB SSD1+ TB SSD + 冷数据 HDD
网络千兆万兆万兆 + 专用存储网
OSLinux 7+Rocky Linux 9 / CentOS 8同推荐
文件描述符655366553665536+

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 chronyd

3.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-la

3.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/profile

3.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-*.out

3.2.4 停止集群

./bin/stop-cluster.sh

3.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 EOF

3.3.3 SSH 免密登录

# Master 节点生成密钥ssh-keygen-trsa-b4096# 分发到所有节点ssh-copy-id flink-master ssh-copy-id flink-worker1 ssh-copy-id flink-worker2

3.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|grepTaskManagerRunnerdone

3.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_456

3.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 2

3.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-operator

3.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.yaml

3.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/# 不依赖 ZooKeeper

3.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:4

3.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:4

3.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:true

3.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: 10s

3.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 ManagersTM 列表:内存、Slot、线程栈
JobManagerJM 状态:内存、配置、日志

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用途
Kafkaflink-sql-connector-kafka-2.3.0.jarKafka 读写(内置,无需额外)
JDBCflink-sql-connector-jdbc-2.3.0.jarMySQL 读写
MySQL Drivermysql-connector-java-8.0.33.jarJDBC 驱动
Paimonpaimon-flink-2.3-1.4.x.jarPaimon 湖存储
Dorisdoris-flink-connector-1.0.3.jarDoris Sink
Flink CDCflink-sql-connector-mysql-cdc-3.6.x.jarMySQL 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.jar

3.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$'done

3.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/
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/10 2:55:13

EtherCAT 的 PDO 映射到底怎么工作的?从对象字典到过程数据内存

在前面的文章中&#xff0c;我们已经从 EtherCAT 主站的角度理解了 IgH EtherCAT Master 的基本架构&#xff0c;也进一步分析了 Datagram、PDO 和 Domain。到了这里&#xff0c;一个非常关键的问题自然会出现&#xff1a;EtherCAT 从站明明已经被 IgH 识别了&#xff0c;为什么…

作者头像 李华
网站建设 2026/10/10 2:54:45

【YOLO】YOLOv5原理

===== 【相关资源获取途径】 ===== 1.CSDN:YOLO专栏 2.---------: XiaoJ1234567 目录 1. YOLOv5 概述 1)目标检测与 YOLO 2)检测网络三段式(Backbone,Neck,Head) 3)YOLOv5 介绍 2. YOLOv5 结构及参数 1)YOLOv5 模型整体结构 2)YOLOv5 整体参数(yaml 里) 3)YOLOv5 …

作者头像 李华
网站建设 2026/10/10 2:54:05

海外仓错发率居高不下?从WMS系统设计看二次复核与防错机制

海外仓错发率居高不下&#xff1f;从 WMS 系统设计角度看“二次复核”与防错机制的实现跨境电商圈子里有个说法&#xff1a;发错一件货&#xff0c;轻则赔运费、赔货款&#xff0c;重则丢账号、丢信任。我见过不少做海外仓的朋友&#xff0c;月错发率从0.5%涨到2%就急得睡不着&…

作者头像 李华
网站建设 2026/10/10 2:53:42

【前端】TypeScript学习中…(着重与JS区别)

TS是JS静态类型系统&#xff0c;完全兼容JS&#xff0c;TS通过类型注解、接口、泛型等特性在编写时就能检测类型错误&#xff0c;提前规避bug&#xff08;JS变量类型动态&#xff0c;可随时变更&#xff0c;运行时才能发现报错&#xff09;1、类型1&#xff09;基础&#xff1a…

作者头像 李华