简介:面向大数据平台工程师与实时计算开发者,这份资源提供在CDH 6.3.2环境中部署Flink 1.13.1所需的完整文件集合,旨在帮助团队解决在Cloudera生态中快速接入实时计算能力的难题。压缩包为zip格式,共包含5个文件,整体大小约299.52MB,内有Flink核心库JAR、YARN客户端JAR、适配Scala 2.11与CentOS 7的Parcel分发包、Parcel元数据JSON以及SHA校验文件,这些文件分别承担程序运行、集群提交、节点分发、配置管理和完整性校验等职责,覆盖了从部署准备到作业提交的各个关键环节。该资源目前已被3492人学习下载,实用性与关注度兼备。通过这套资料,读者可以直接获得与CDH 6.3.2精确匹配的预编译二进制文件,无需自行处理编译与适配问题,并且能够借助Parcel机制简化多节点部署,结合YARN完成资源统一调度,从而显著缩短Flink的上线周期。对于需要在大数据集群中快速启用实时计算框架的团队,这是一套结构清晰、开箱即用的基础素材,同时也为理解Flink与CDH的集成细节提供了很好的参考。
1. 为什么要在 CDH 6.3.2 上折腾 Flink 1.13.1
先说个背景。如果你所在的公司还在用 CDH 6.3.2,又想上实时计算,那几乎逃不开这个问题:CDH 6.3.2 自带的 Flink 组件版本老得可怜,而社区版的 Flink 1.13.1 又是一个非常稳定、用得最广的版本,两者之间的兼容性就成了绕不过去的坎。
我最初接到这个任务时,第一反应是“直接下载 flink-1.13.1-bin-scala_2.12.tgz 解压就能跑”,结果被 YARN 上的一堆 ClassNotFound 异常教做人了。后来查了一圈才发现,问题出在 Hadoop 客户端版本的匹配上。CDH 6.3.2 底层是 Apache Hadoop 3.0.0 的发行版,但 Cloudera 对很多内部实现做了修改,直接把 Apache 原版 Flink 扔上去,会因为缺少对应版本的 shaded Hadoop 依赖而翻车。
这篇文章就是把我自己从“下载 Flink 包”到“在 YARN 上稳定跑通任务”的完整过程记录下来。所有内容都基于 CDH 6.3.2 + CentOS 7 + Flink 1.13.1 这套组合,适合正在搞实时数仓、需要把 Flink 任务提交到公司已有 CDH 集群上的同学参考。折腾过一遍之后你会发现,真正的坑不在 Flink 本身,而在 Hadoop 依赖的取舍上。
2. 版本兼容性:这套组合的底层逻辑
2.1 CDH 6.3.2 的 Hadoop 到底是什么版本
CDH 6.3.2 对应的底层 Hadoop 版本是 3.0.0,这一点很多文章提过,但实际环境里远没那么简单。Cloudera 的发行版对 Hadoop 做了大量 patch,包括 YARN 的调度器改进、HDFS 的 NameNode 优化、以及各种安全相关的修复。也就是说,你面对的不是一个“干净的 Apache Hadoop 3.0.0”,而是一个被深度定制过的版本。
这种定制化带来的直接影响是:Flink 提交作业时需要用 Hadoop 客户端与 YARN ResourceManager、HDFS NameNode 通信,如果客户端依赖跟集群的 RPC 协议对不上,就会出现各种奇怪的序列化异常、协议不匹配的报错。这也是为什么网上很多教程强调“一定要用 CDH 自带的 hadoop-client”,而不是随意引一个 Maven 里的 hadoop-client 依赖。
2.2 Flink 1.13.1 官方包的隐藏问题
Flink 1.13.1 官方提供了两个关键下载包:flink-1.13.1-bin-scala_2.11.tgz和flink-1.13.1-bin-scala_2.12.tgz。如果你只是做离线计算或者纯 DataStream API 开发,选哪个都行;但如果你用到了 Table API 或者 SQL,建议统一用 Scala 2.12 版本,因为社区后续的很多连接器(比如 flink-connector-kafka 的某些版本)对 2.12 支持更好。
真正的问题在这两个包的 lib 目录下。默认情况下,官方包里带的是flink-shaded-hadoop-uber-2.10.2.jar(有些版本是 3.x 的 uber jar),这个 jar 是给“裸机部署 Flink”用的,里面打进去了完整的 Hadoop 2.10.2 客户端。但在 CDH 6.3.2 上,这个 uber jar 反而成了灾难的源头,因为它内部的 Hadoop 版本跟 CDH 的 YARN 通信时,会出现org.apache.hadoop.yarn.exceptions.YarnRuntimeException这类异常。
所以第一步就必须把 lib 目录下默认的flink-shaded-hadoop-uber-*.jar删掉,换成 CDH 自带的 hadoop-client。
2.3 明确自己的部署模式再选版本
Flink 在 YARN 上有两种主流部署模式:Session 模式和 Application 模式。Session 模式是预先启动一个常驻的 Flink 集群,然后往里面提交多个作业,适合小作业多、启动频繁的场景。Application 模式是每个作业启动一个独立的 Flink 集群,作业结束集群自动释放,适合大作业、长任务、资源隔离要求高的场景。
我推荐的生产环境方案是 Application 模式,原因后面实操部分会细说。这里先记住一个结论:无论哪种模式,你都要把 CDH 的 Hadoop 客户端依赖准备齐全,而不是依赖 Flink 自带的 uber jar。
3. 部署前置工作:环境检查与包下载
3.1 集群环境硬性条件确认
动手部署前,先把环境摸清楚。我用的是三台 CentOS 7.6 节点组成的 CDH 6.3.2 集群,节点角色分配是:两个 Master 节点(一个跑 NameNode + ResourceManager,另一个跑 Standby NameNode + JobHistory),三台 Worker 节点(跑 DataNode + NodeManager)。
以下是关键版本的核对清单,每一条都要确认。
| 组件 | 版本要求 | 说明 |
|---|---|---|
| JDK | 1.8(最好是 CDH 自带的) | Flink 1.13.1 官方要求 Java 8 或 11,实测 1.8 最稳 |
| Hadoop | 3.0.0(CDH 6.3.2 内置) | 不要用其他版本的 hadoop-client 替换 |
| Zookeeper | 3.4.5+(CDH 6.3.2 内置 3.4.5) | Flink 的 HA 模式和部分连接器依赖 |
| YARN | 已启用且 NodeManager 正常 | 用yarn node -list确认节点都在线 |
| HDFS | 已配置且 NameNode 处于 Active 状态 | 用hdfs dfsadmin -report看容量 |
有一个容易被忽略的点:确保 Flink 部署的机器(通常叫 client 节点)能跟 YARN ResourceManager 的 8032 端口、HDFS NameNode 的 8020 端口正常通信。排查方式就是telnet或nc -vz测一下,别等跑任务时才报 Connection refused。
3.2 选择 Flink 包:两个安装包怎么选
Flink 1.13.1 的下载页提供两个二进制包:flink-1.13.1-bin-scala_2.11.tgz和flink-1.13.1-bin-scala_2.12.tgz。如果你不想踩坑,直接选 Scala 2.12 版本。
原因不只是社区推荐,更实际的一点是:Flink 1.13 之后,很多第三方连接器(比如较新的 Kafka connector、Hive connector)的预编译包都优先基于 Scala 2.12。你如果选了 2.11,后面想扩展 API 时容易因为 Scala 版本不兼容而找不到对应依赖。
3.3 下载并解压到指定目录
这里我习惯把 Flink 放在/opt/flink下,并创建一个软链接指向当前版本,方便后续升级。
# 下载(这里以 2.12 版本为例) wget https://archive.apache.org/dist/flink/flink-1.13.1/flink-1.13.1-bin-scala_2.12.tgz # 解压 tar -zxvf flink-1.13.1-bin-scala_2.12.tgz -C /opt/ # 创建软链接 ln -s /opt/flink-1.13.1 /opt/flink解压完成后,先别急着改配置。进入lib目录看一眼,确认默认的 Hadoop uber jar 是否存在。这一步决定了你后面会不会踩坑。
3.4 关键准备:构建一个“干净”的 Flink lib 目录
这一步是整个部署中最容易踩坑的地方,我单独拿出来说。默认lib目录下会有flink-shaded-hadoop-uber-2.10.2.jar,在 CDH 上必须移走。同时,为了能往 YARN 上提交作业,你还需要把 CDH 的 Hadoop 客户端相关 jar 放进来。
具体做法是:在集群的任意一台机器上,找到cloudera/parcels/CDH/lib/hadoop/client目录,把里面的所有 jar 复制到 Flink 的 lib 目录下,再额外补充hadoop-mapreduce-client-core、hadoop-mapreduce-client-common等几个 jar,因为 Flink 提交 YARN 作业时需要它们。
如果你不想手动拼 jar,还有一个取巧的办法:直接用 CDH 自带的hadoop classpath命令输出完整依赖列表,然后把这些依赖全部放进 Flink 的 lib 目录。但这个做法会导致 lib 目录非常臃肿,我实测下来虽然能跑,但启动速度会慢很多。
4. 核心配置:flink-conf.yaml 的每一项都是经验
4.1 基础配置清单
解压完成并整理好 lib 目录后,接下来就是核心配置环节。文件在${FLINK_HOME}/conf/flink-conf.yaml,我直接给出我在 CDH 6.3.2 上稳定运行的关键配置项,并解释每一行的原因。
# JobManager 内存:建议至少 1G,如果任务多可以按 2G 起步 jobmanager.memory.process.size: 1600m # TaskManager 内存:需要根据节点剩余内存和并发度来衡量 taskmanager.memory.process.size: 4096m # 每个 TaskManager 的 slot 数:经验值 = CPU 核数 / 每个任务所需核数 taskmanager.numberOfTaskSlots: 4 # 默认并行度:建议先从 2 开始跑通,后续再根据集群规模调整 parallelism.default: 2 # 指定 YARN 上运行的 Flink 集群名称 yarn.application.name: FlinkOnCDH # 队列:如果你有专门的实时任务队列就写上,否则用 default yarn.application.queue: default所有内存配置都建议使用process.size而不是heap.size,因为前者会把堆外内存、Metaspace、网络缓冲等全部计算在内,更符合 YARN 容器内存限制的判断逻辑,避免频繁被 NodeManager 杀掉。
4.2 必须开启的 HA 配置
如果集群有多个 Master 节点,强烈建议开启 Flink 的 HA 模式。Flink 1.13.1 的 HA 依赖 ZooKeeper,CDH 6.3.2 自带 ZooKeeper 3.4.5,可以直接复用。
# 启用 HA high-availability: zookeeper # ZK 连接串:多个节点用逗号分隔 high-availability.zookeeper.quorum: cdh-master-01:2181,cdh-master-02:2181,cdh-worker-01:2181 # 在 ZK 中存储 Flink 集群信息的根路径 high-availability.zookeeper.path.root: /flink-ha # JobManager 的元数据存放在 HDFS 上 high-availability.storageDir: hdfs:///flink/ha/ # 集群 ID:每个独立 Flink 集群必须唯一 high-availability.cluster-id: /flink-on-cdh这里注意一个问题:high-availability.cluster-id在 Flink 1.13 版本中变得非常重要。多个 Flink 集群如果共享同一个 ZK 根路径,cluster-id 就是它们在 ZK 中的唯一标识,不唯一会导致集群互相干扰。
4.3 历史记录保存与日志级别
Flink 1.13.1 跑完作业后,Web UI 默认不会自动同步历史任务数据。如果你想要类似 Spark HistoryServer 的效果,可以把作业的 Web 访问地址写入 HDFS:
jobmanager.archive.fs.dir: hdfs:///flink/completed-jobs/ historyserver.web.address: cdh-master-01:8082 historyserver.web.port: 8082日志级别默认是 INFO,排查问题时通常够用,但如果任务经常被 YARN kill,可能需要临时调成 DEBUG,定位具体是内存超了还是心跳超时。
4.4 配置环境变量
在/etc/profile.d/flink.sh中写入:
export FLINK_HOME=/opt/flink-1.13.1 export PATH=$FLINK_HOME/bin:$PATH export HADOOP_CLASSPATH=/opt/cloudera/parcels/CDH/lib/hadoop/client/*第三行很有用。因为 Flink 启动 YARN 作业时,需要拿HADOOP_CLASSPATH里的 YARN 和 HDFS 相关类。直接把 CDH 的 client lib 目录通配符设进去,省去手动拼 jar。
4.5 配置 workers 和 masters
如果你用 Session 模式,还要在conf/slaves(1.13 版本仍叫 slaves)文件里填 Worker 节点的主机名。但如果是 Application 模式,这个文件基本用不上,因为 TaskManager 的分配由 YARN 调度决定。
conf/masters文件在 Standalone 模式下才需要,跑在 YARN 上时可以不管。
5. 提交前的最后校验:验证环境是否真的通
配置完成后,别急着提交正式任务,先跑几个简单的验证命令。
# 检查 Flink 能否正确识别 Hadoop 版本和依赖 ${FLINK_HOME}/bin/flink --version # 查看 YARN 队列和资源情况 yarn node -list -all yarn queue -status default # 验证 HDFS 可选 hdfs dfs -mkdir -p /flink/ha hdfs dfs -ls /flink6. 在 YARN 上跑通第一个任务的完整流程
6.1 Session 模式的快速验证
先试 Session 模式,因为它最简单,适合验证环境配置是否正确。
${FLINK_HOME}/bin/yarn-session.sh -n 2 -tm 2048 -s 2 -d参数说明:
-n 2:申请 2 个 TaskManager(1.13 里可能提示废弃,改用-D yarn.taskmanager.node.label或直接用yarn-session.sh -d)-tm 2048:每个 TaskManager 的内存是 2048MB-s 2:每个 TaskManager 有 2 个 slot-d:后台运行
启动后,控制台会打印一行 URL,类似http://cdh-master-01:34912,这个就是 Flink Web 界面地址。打开后你能看到集群的资源使用情况、已提交作业列表。
6.2 提交一个 Batch 作业测试
这里直接用一个统计文本行数的简单作业做验证。
${FLINK_HOME}/bin/flink run \ -m yarn-cluster \ -ytm 2048 \ -p 2 \ ${FLINK_HOME}/examples/batch/WordCount.jar \ --input hdfs:///tmp/test.txt \ --output hdfs:///tmp/flink_output如果输出目录正常生成,说明整个链路(Flink -> YARN -> HDFS)已经打通,可以进入下一阶段。
6.3 Application 模式:生产推荐的部署方式
Session 模式的好处是启动快、可以复用集群,但它有一个很大的隐患:多个作业共享同一个 TaskManager 的 JVM,如果某个作业里有恶意的静态变量或大对象,很容易把别的作业也拖垮。生产环境我推荐用 Application 模式。
Application 模式的命令很明确,每个作业单独跑一套 Flink 集群:
${FLINK_HOME}/bin/flink run-application \ -t yarn-application \ -Djobmanager.memory.process.size=1024m \ -Dtaskmanager.memory.process.size=2048m \ -Dtaskmanager.numberOfTaskSlots=2 \ -Dparallelism.default=2 \ -Dyarn.application.name="MyFlinkJob" \ -c com.example.MyJob \ /path/to/your-job.jar优点突出:作业之间完全隔离,JobManager 挂掉后 YARN 会自动重启(需要配置yarn.application-attempts),用完资源自动释放。缺点是要多敲几个-D参数,但可以在脚本里封装成一个模板。
6.4 从 Web UI 或命令行验证任务状态
登录 Flink Web UI 后,重点看这几个指标:
- Job Status:必须是 RUNNING
- TaskManagers 数量:和申请的
-n或 taskmanager 数量一致 - Memory Usage:堆内存不能频繁触发 GC
- BackPressure:如果长时间高,说明数据倾斜或下游处理不过来
命令行验证用 REST API 也可以:
curl http://cdh-master-01:34912/jobs/overview7. 常见问题与排查技巧实录
7.1 问题速查表
| 报错信息 | 根因 | 解决方案 |
|---|---|---|
java.lang.ClassNotFoundException: org.apache.hadoop.yarn.exceptions.YarnRuntimeException | Flink 默认带的 shaded hadoop jar 与 CDH 不兼容 | 删除 lib 目录下的 shaded uber jar,替换为 CDH hadoop client 的 jar |
Caused by: org.apache.hadoop.security.AccessControlException: Permission denied | HDFS 目录权限不足 | 用 hdfs 用户创建目录并授权:hdfs dfs -chmod -R 755 /flink |
YarnApplicationMaster: Exception in thread "main" java.lang.IllegalArgumentException: Required executor memory exceeds max memory allocated | TaskManager 内存申请超过容器最大限制 | 检查yarn.nodemanager.resource.memory-mb和yarn.scheduler.maximum-allocation-mb |
java.net.ConnectException: Connection refusedon ResourceManager host | 网络不通或端口未开放 | 检查 8032/8030/8088 等端口 |
NoClassDefFoundError: org/apache/flink/streaming/api/environment/StreamExecutionEnvironment | 作业 jar 缺少 Flink 相关依赖 | 用 Maven 构建 fat jar,确保包含所有依赖 |
7.2 实际踩坑记录:TaskManager 频繁被杀
我有一次在 CDH 节点上跑 Flink 任务,发现 YARN 反复 kill 掉 TaskManager,日志显示 Container 超过物理内存限制。排查过程:
第一步,先看 YARN NodeManager 日志,确认是物理内存超限还是虚拟内存超限。
第二步,用-Dtaskmanager.memory.process.size=2048m代替原先的-Xmx方式,因为 process.size 会把堆外内存也算进去,YARN 校验内存时会更准确。
第三步,如果虚拟内存总是超,修改yarn.nodemanager.vmem-check-enabled=false(不推荐生产用,但可以临时验证问题是不是虚拟内存引起)。
我最后发现是默认的 JVM Metaspace 设定太大,每个 TaskManager 都占了额外一部分内存,调低taskmanager.memory.jvm-metaspace.size后稳定跑通了。
7.3 一个容易被忽略的细节:ZK 会话超时
CDH 集群通常网络环境复杂,Flink 的 HA 使用 ZooKeeper 时,如果 ZK 会话超时时间过短,JobManager 主备切换时容易出现误判。建议在flink-conf.yaml中加一个参数:
high-availability.zookeeper.client.session-timeout: 60000这里时间单位是毫秒,我设成 60 秒。默认值偏小,在 GC 停顿或网络抖动时容易触发重新选举,导致作业重启。
7.4 调试技巧:用 log4j 快速定位问题
Flink 的日志配置文件在conf/log4j-cli.properties和conf/log4j-yarn-session.properties。遇到难题时,把这些文件里的rootLogger.level临时调成 DEBUG,跑一个最小复现任务,日志会详细打印出 YARN 提交过程中的每个步骤,排查效率能提升很多。但线上一定记得调回 INFO,否则日志量太大,反而把关键信息淹没。
8. 写在最后:从这里还能再往下走几步
我在实际部署过程中发现一个很有用的做法:把整个部署步骤写成一个 Ansible 脚本,包括 lib 目录整理、flink-conf.yaml 生成、环境变量配置、HA 目录初始化,以后新加节点或重建环境时,一条命令搞定,省去手动操作的重复劳动。
如果你准备进一步做实时数仓,下一步建议研究 Flink SQL 与 Hive Catalog 的集成,这里有个天然优势:CDH 6.3.2 已经内置了 Hive 3.1.0,Flink 1.13.1 对 Hive 3.x 支持的适配度不错,只是需要额外准备一个flink-sql-connector-hive-3.1.0_2.12的 jar。等这一整套跑通之后,你就能在一个相对统一的平台上同时处理流式任务和离线数仓的数据同步了。
眼下这套 Flink 1.13.1 on CDH 6.3.2 的组合,虽然版本不算最新,但在很多公司的生产环境里依然占据不少份额。关键是把底层依赖关系理清楚,后面的扩展和排障都会顺畅很多。
本文还有配套的精品资源,点击获取