这些年在公司做实时数仓,经常要面对的一个问题就是 Flink 和 Hadoop 的版本配套关系。Flink 1.13 这个版本用得人不少,但当你拿着官方下载的 Flink 1.13 包去对接一个 Hadoop 3.x 集群时,十有八九会碰一鼻子灰,报错信息千奇百怪,核心原因只有一个:默认发行包根本不是按照 Hadoop 3 给你准备的。这篇文章就专门说清楚 Flink 1.13 集成 Hadoop 3.x 到底该怎么解决,从版本冲突的原理解析到两条可落地的操作路线,再到运行期会踩的那些坑,我尽量按实际项目里能直接抄作业的方式来写。
我当时的场景比较典型:生产环境 HDFS 已经升级到 Hadoop 3.3.x,YARN 也换了新版本,但实时任务的代码还跑在 Flink 1.13 上。整理出来的解决方案总共有两条主流路线,一条是从源码编译一个适配 Hadoop 3 的 Flink 发行包,另一条是使用官方无 Hadoop 依赖的包,再让 Flink 加载集群自带的 Hadoop 客户端环境。两条路线分别适合不同的运维条件,后文我会把关键步骤和踩坑点都摊开说。
1. 先搞明白:Flink 1.13 和 Hadoop 3.x 的版本代沟在哪里
1.1 默认包用的是 Hadoop 2 协议栈
Flink 1.13 发布的时候,官方二进制包默认内置的是 Hadoop 2.10.x 的 shaded 客户端。这个“内置”不是你写代码时引入一个依赖那么简单,而是 Flink 的 lib 目录下直接躺着flink-shaded-hadoop-2-uber打头的 jar 包,整个分布式文件系统访问、提交 YARN 应用、读 HDFS 上 checkpoint 的行为,全部默认按 Hadoop 2 的这套类库来走。
Hadoop 2 和 Hadoop 3 之间的差异不是换个版本号就能糊弄过去的。HDFS 客户端代码里有相当一部分类的包名变了,比如 HDFS 文件系统 provider、一部分 RPC 协议类,还涉及到 UGI、token 的底层处理方式。Hadoop 3 的 RPC 协议本身也升级了,NameNode 和 DataNode 的通信协议版本号不同。你拿着 Hadoop 2 的客户端去连 Hadoop 3 的 NameNode,双方做 RPC 握手的时候,NameNode 一看协议版本对不上,直接拒绝连接。
同时,Hadoop 3 默认启用的服务级别和端口配置也跟 Hadoop 2 不完全一样。比如 HDFS NameNode 的 RPC 地址相关配置项在 Hadoop 3 里更强调dfs.namenode.rpc-address的显式配置,HA 场景下 Nameservice 的写法也有变化。如果 Flink 进程里的 Hadoop 客户端是 2.x,它可能无法正确解析 Hadoop 3 的配置语义,最后表现出来的就是一个非常让人困惑的连接超时,或者 “Failed to connect to /...:8020” 这种错误。很多人一开始以为是自己网络配错了,其实根子在客户端协议栈版本不对。
1.2 集成失败时的典型报错面孔
我在多个环境里见过集成失败时出现的报错,把它们归类一下,基本就三种面孔:
第一类是ClassNotFoundException。比如org.apache.hadoop.hdfs.DistributedFileSystem找不到,或者org.apache.hadoop.hdfs.protocol.HdfsFileStatus找不到。这个很容易理解,lib 目录下那个 Hadoop 2 的 uber jar 里就没有 Hadoop 3 新引入的类,或者包名对不上。
第二类是NoSuchMethodError或者NoClassDefFoundError。例如某段代码调用org.apache.hadoop.fs.FileSystem.get(...)时,运行期发现方法签名不一致,或者加载到了两个不同版本的 FileSystem 类。这种报错比 ClassNotFoundException 更隐蔽,因为类名看起来是存在的,只是编译期和运行期的 class 不是同一个。
第三类是在 YARN 客户端提交阶段直接报错。常见的有org.apache.hadoop.yarn.exceptions.YarnRuntimeException,或者是跟 YARN ResourceManager 握手失败。这时候要立刻怀疑 Flink 发行包内的 YARN 客户端库版本不匹配,不要先怀疑网络或者权限。
如果看到的是下面这种典型信息:
java.lang.RuntimeException: java.lang.NoClassDefFoundError: org/apache/hadoop/hdfs/provider/HdfsFileSystemProvider那基本可以确定问题就出在 Flink 内置 Hadoop 客户端与集群 Hadoop 版本不一致上。你不需要继续深挖其他原因,直接把注意力放到“让 Flink 使用正确版本的 Hadoop 客户端”这件事上来。
2. 集成路线怎么选:源码编译、替换 jar 还是 HADOOP_CLASSPATH
2.1 三条路线的成本对比
网上搜 Flink 1.13集成 Hadoop 3.x,方案五花八门,但剥掉外壳,核心就三条路线。下面这个表格是我对照自己维护过的几套环境整理的,简单直接:
| 方案 | 操作难度 | 对集群影响 | 维护成本 | 适合场景 |
|---|---|---|---|---|
| 从源码编译 Flink 发行包 | 中等,需要 Maven 环境和网络 | 几乎无,纯客户端替换 | 低,一次编译长期使用 | 生产环境长期使用,有统一发版规范 |
| 替换 lib 下的 shaded jar | 低,但容易遗漏 | 几乎无 | 中,后续升级 Flink 还得再来一遍 | 临时验证,内部测试环境 |
| 无 Hadoop 依赖包 + HADOOP_CLASSPATH | 低,环境变量一配就行 | 无,依赖每个节点的 hadoop 命令 | 中,依赖节点环境稳定性 | 测试环境,或者客户端环境高度可控的小集群 |
很多人第一反应是选第二个方案,觉得把 lib 下的flink-shaded-hadoop-2-uber换成 Hadoop 3 的对应包就行。这个思路方向对,但实际操作有个麻烦:官方发布包里并不总会捆绑 Hadoop 3 版本的 shaded uber jar,你经常要去自己找对应 Flink 版本的flink-shaded-hadoop-3-uber,或者其他仓库里编译好的产物。组件多的时候,依赖版本很容易对不上,遇到奇奇怪怪的冲突反而更难排查。
第三个方案是目前官方文档里也认可的做法,而且对很多团队来说非常省事。它的核心思想是:Flink 只负责自己的流计算引擎,Hadoop 相关的类全部从系统环境的 Hadoop 安装目录里去拿。让 Flink 启动时的 classpath 包含$(hadoop classpath)的输出就行。这个方式的坑在于依赖每个提交节点的 Hadoop 环境必须完整且版本一致,一旦某台节点环境变量不对,任务提交表现就很不稳定。
2.2 为什么我优先推荐源码编译路线
如果有条件,我建议优先走源码编译。原因不是这条路最时髦,而是因为它把 Flink 和 Hadoop 的版本关系彻底锁死在了构建产物里。编译时指定 Hadoop 3.x 版本,Flink 源码里的 Hadoop 相关模块会按照这个版本来编译和打包,最终生成的发行包中内置的 shaded jar 就是 Hadoop 3 的客户端。这一步做完,后面所有节点的部署、提交、运行,都不需要额外担心客户端版本漂移。
我理解很多团队的顾虑:编译 Flink 从源码跑起来太耗时,而且怕搞坏依赖。其实 Flink 本身不依赖 Hadoop 集群环境,只要 Maven 能从中央仓库拉到对应版本的 Hadoop jar,编译就能顺利完成。第一次编译半小时到一小时很常见,之后如果再要编译其他组件版本,速度会快很多。而且用源码编译还能顺手把flink-shaded-hadoop-3-uber这个产物放到公司内部私有仓库,给其他同事用,后续收益很大。
替换 jar 或者无 Hadoop 包方案更适合应急。我见过一些团队用无 Hadoop 包的方式把任务跑起来了,看起来非常简单,但后面升级组件或者新增节点时经常出幺蛾子。所以本文后面对源码编译路线做更详细的展开,无 Hadoop 包路线也会给出完整实操步骤。
3. 路线一实战:从 Flink 1.13 源码编译适配 Hadoop 3
3.1 环境准备和版本取舍
开始编译之前,先把基础环境准备好。Flink 1.13 这个版本建议 JDK 8,虽然 JDK 11 也能编译运行,但 JDK 8 是官方测试最充分的。Maven 版本 3.6 以上,Git 自然是必须的。
编译机不需要安装 Hadoop,也不需要什么特殊权限,只要 maven 仓库能访问外网,或者你配置了公司内部的 Maven 镜像,拉取依赖没问题就够用。
版本取舍这块要特别说一下。Flink 1.13 源码里其实内置了一个hadoop3profile,如果你直接激活这个 profile,默认拉取的 Hadoop 版本是 3.2.0,Hive 相关模块默认版本是 3.1.2。如果你集群是 Hadoop 3.2.x,直接用内置 profile 最省心,不用改任何 pom 文件。如果你的集群是 Hadoop 3.3.x,那就需要在编译命令里额外指定-Dhadoop.version=3.3.1之类的版本号。
我一般会先确认集群上部署的 Hadoop 具体版本,比如hadoop version命令返回的信息,再用完全一致的客户端版本去编译。这里是强制建议按照客户端和集群版本保持一致来配,不要只图省事用默认的 3.2.0 去连接 3.3.x 集群。虽然大多数情况下能跑通,但 RPC 协议在 3.3 以后有小版本演进,本地文件系统接口也有调整,真踩到协议上不兼容的问题,排查成本比编译成本高得多。
3.2 具体编译命令
从 GitHub 拉取 Flink 1.13 的 release 分支。注意不要拉 master,而是拉release-1.13这个分支,版本号相对干净,也不会夹带后续版本的未发布特性。命令如下:
git clone -b release-1.13 https://github.com/apache/flink.git cd flink接下来执行编译。如果集群是 Hadoop 3.2.x,直接用内置 hadoop3 profile:
mvn clean package -DskipTests -Phadoop3 -Pinclude-hadoop这里解释一下参数含义。-Phadoop3是激活 Flink 源码中预定义的 Hadoop 3 版本属性集,-Pinclude-hadoop表示生成的发行包内要带上 Hadoop 客户端库。如果不带后者,编出来的 Flink 包不会在 lib 目录下生成 shaded Hadoop jar,那这个包就得走 HADOOP_CLASSPATH 路线了。
如果集群是 Hadoop 3.3.x,例如 3.3.1,那就在上面命令后面追加:
mvn clean package -DskipTests -Phadoop3 -Pinclude-hadoop -Dhadoop.version=3.3.1这里-Dhadoop.version=3.3.1的含义你一看就懂,就是覆盖源码 profile 里默认的 3.2.0。
编译过程中大概率会遇到某些依赖下载失败,或者编译到某个模块报错,这些多数是网络问题或 Maven 仓库源的问题。可以在~/.m2/settings.xml里配置阿里云或公司私有 Maven 镜像。如果报的是某个插件版本号解析失败,可以试着先执行一次mvn clean install -DskipTests -T 8 -Dcheckstyle.skip,跳过一些校验把依赖先拉到本地,再重新执行完整打包。
我这里再补充一个可选的 Scala 版本参数。Flink 1.13 默认 Scala 2.12,如果你团队内部有标准约定要用 Scala 2.11,那就需要指定 profile。我建议生产环境直接跟随官方默认 2.12,没必要在 Scala 版本上增加额外变量,否则后续写 UDF 或者依赖 Flink 周边组件时,很容易出现 Scala 二进制版本冲突。
3.3 编译后的完整性检查
编译结束后,主要看两个地方。
第一,确认flink-dist/target/flink-1.13.x-bin/目录下的发行包已经生成。这个目录里就是可以直接用的 Flink 目录,内部结构和官方打包出来的很像,包括 bin、lib、conf、plugins 等目录。
第二,重点检查lib目录下的 Hadoop 相关 jar。正常情况下,你会看到类似flink-shaded-hadoop-3-uber-xxx.jar的文件,而不是默认发行包里那个flink-shaded-hadoop-2-uber-xxx.jar。用一句话验证命令检查即可,比如:
ls -lh flink-1.13.*/lib/ | grep hadoop如果输出里flink-shaded-hadoop-3-uber存在,说明编译产物已经内置了 Hadoop 3 客户端。这时候把这个 Flink 目录拷贝到所有需要提交作业的节点,进入新任务上线阶段。
还要顺带检查opt目录下的其他组件,比如flink-sql-connector-hive相关包,因为这个包也需要和你集群的 Hive/Hadoop 版本对应。如果 SQL 任务比较多,编译时最好把 Hive 连接器也一并编译进去,后面省事。
4. 路线二实战:无 Hadoop 发布包配合集群客户端类路径
4.1 客户端节点准备工作
如果你暂时不能编译 Flink,或者只想在测试环境快速验证,这条路也完全可行。它的思路是:下载官方提供的flink-1.13.x-bin-scala_2.12.tgz,注意文件名里没有hadoop2,这个就是不带 Hadoop 依赖的版本。接着在运行 Flink 的节点上,利用系统自带的 Hadoop 安装目录来提供所有 Hadoop 类。
这里有个前置条件,就是运行 Flink 的节点上必须安装了 Hadoop 客户端,或者说至少有hadoop命令、完整的 Hadoop 配置文件目录HADOOP_CONF_DIR和配套客户端 jar。生产环境如果是纯 Flink 节点,没有装 Hadoop 客户端,那这条路走不了。
我在实际环境中最稳妥的做法是在conf/flink-conf.yaml顶部加一个环境变量引用,不直接改系统全局配置,把影响面缩小到 Flink 内部。比如在 conf/flink-conf.yaml 最前面写入:
env.java.opts: -Djava.library.path=$HADOOP_HOME/lib/native然后更关键的是把 Hadoop 的 classpath 传给 Flink。方式很多,我建议在 Flink 目录下的bin/flink调用命令前,使用 export 的方式设置:
export HADOOP_CLASSPATH=$(hadoop classpath) export HADOOP_CONF_DIR=/etc/hadoop/conf export FLINK_CLASSPATH=$HADOOP_CLASSPATH bin/flink run -m yarn-cluster ...hadoop classpath这个命令会把整个 Hadoop 安装目录下所有客户端 jar、配置文件地址、第三方依赖全部拼成一条 classpath 输出。把这条内容放进 Flink 的 classpath 之后,Flink 进程加载 Hadoop 类时就不再依赖 lib 目录下的 shaded jar 了。
4.2 本地模式任务验证
拿到无 Hadoop 依赖包之后,不要一上来就直接提 YARN 作业,先跑一个本地模式的 HDFS 读写任务,验证 classpath 是否正确。
最简单的测试是提交一个能读取 HDFS 文件路径的任务。比如先准备一个文本文件放到 HDFS:
hadoop fs -put /tmp/test.txt /tmp/flink_test/然后找一个 Flink 自带的、可以触发文件系统访问的 example,或者直接写一个每几秒钟读一次 HDFS 的小任务。如果 classpath 有问题,最常见的情况是报org.apache.hadoop.fs.FileSystem找不到,或者找不到 HDFS 对应的 FileSystem 实现类,这种报错说明hadoop classpath没生效。
运行时也可以加-t yarn-session等参数,但我建议先跑 local 模式确认基础网络和 RPC 没问题,再上 YARN。本地模式跑的顺不代表 YARN 就顺,因为 YARN 模式还涉及资源调度、应用提交这些环节,但至少可以隔离掉大量底层类的缺失问题。
4.3 YARN 模式提交验证
YARN 模式提交时,重点检查两个环境变量:HADOOP_CLASSPATH和HADOOP_CONF_DIR。
HADOOP_CONF_DIR指向的目录里必须包含core-site.xml、hdfs-site.xml、yarn-site.xml。Flink 在向 YARN 集群申请资源的时候,需要从这些配置里找到 ResourceManager 地址、HDFS 地址,以及各种安全认证相关配置。如果配置目录不对,最常见的报错是 Connection refused,指向的地址是 localhost 或错误的 IP。
HADOOP_CLASSPATH一定要在启动 Flink 的 shell 环境里明确 export,因为很多 Flink 版本不会自动去执行hadoop classpath。即使某些版本的脚本会尝试自动读取,也经常因为权限或 PATH 问题失败,所以显式设置是最稳妥的。
做完这些,看任务是否能在 YARN 上申请到容器:
bin/flink run -m yarn-cluster -ys 2 -ytm 2048 -yjm 1024 \ -c org.apache.flink.streaming.examples.wordcount.WordCount \ examples/streaming/WordCount.jar --input hdfs:///tmp/flink_test/test.txt --output hdfs:///tmp/flink_test/out-ys 2表示申请 2 个 TaskManager 槽位,-ytm 2048表示每个 TaskManager 内存 2048MB,-yjm 1024表示 JobManager 内存 1024MB。如果任务能正常启动,看到 “Job has been submitted successfully” 并且 Web UI 上能看到 Running 状态,那这条路线就算通了。
5. 集成后运行期易踩的坑
5.1 类加载冲突与 Guava/Protobuf 引发的 NoSuchMethodError
Flink 和 Hadoop 在运行期都要用到一些非常基础的第三方库,最典型的是 Guava 和 Protobuf,以及 Netty。Hadoop 3 内部某些组件使用的 Guava 版本比 Flink 1.13 内置的更高或者更低,如果这两个版本同时出现在一个 classloader 里,轻则日志漂漂亮亮,运行到某个方法时突然抛NoSuchMethodError。
Flink 1.13 的类加载策略默认是 parent-first,也就是父加载器能加载的类不会让子加载器重新加载。这个策略有时会把 Hadoop 带进来的某个过期 Guava 类先生效,导致 Flink 内部调用时找不到方法。面对这种问题,我一般先检查 Flink Web UI 的日志,确认报错的类是哪个方法名,再用mvn dependency:tree或者jar tf去确认这个类存在于哪些 jar 包中,找出冲突源头。
如果确认是 Guava 冲突,一个处理思路是在conf/flink-conf.yaml中设置类加载顺序:
classloader.resolve-order: child-firstchild-first的意思是让用户 jar 和 Flink 自身提供的类优先被加载,父 classpath 里的同类则靠后。但这个策略要小心,设置完以后,某些依赖父加载器提供的类时可能会引发新的问题。我一般在测试环境试稳定后再推到生产。
更稳妥的思路其实是确保 Flink 的lib目录只有一份 Hadoop 客户端相关的 uber jar,不要同时存在多个版本。很多人习惯把集群上散落的各种 Hadoop jar 一股脑拷贝进 Flink lib 目录,这样特别容易整出多个相同的类,且加载顺序不确定。
5.2 YARN 提交时的资源与权限问题
Flink 任务提交到 YARN 时,还经常遇到两类问题,一类是资源队列权限,一类是 YARN 上可用资源不足。
提交时如果不指定-yqu参数,默认使用default队列。如果账号没有访问 default 队列的权限,报错信息是AccessControlException或者 “Queue default does not exist”。解决方式就是在提交命令里显式指定队列:
bin/flink run -m yarn-cluster -yqu realtime另外一个非常隐蔽的问题,是 YARN 集群的yarn.scheduler.maximum-allocation-mb和yarn.scheduler.minimum-allocation-mb配置。如果 TaskManager 申请的内存超过了最大分配限制,提交会一直卡在 “Waiting for AM container to be allocated” 这种状态。如果申请的内存小于最小分配限制,资源管理器又可能起不来容器。一般来说,先确认 YARN 集群资源管理页面上的可用资源和最大值,再决定-ytm填多少。
HDFS 权限上也要未雨绸缪。Flink 作业运行时的系统用户如果对 HDFS 的 checkpoint 目录或输出目录没有写权限,任务运行起来后 checkpoint 会一直失败。这个不一定会即时导致作业挂掉,但会让任务状态一直变不到 RUNNING,或者频繁重启。所以提交前手动测试一下当前用户在 HDFS 上的读写能力,非常有必要。
5.3 与 Hive Metastore 结合时要额外留意的点
实时数仓场景里,用 Flink SQL 去读 Hive 表或者同步 Hive 元数据是常有的事。Flink 1.13 集成 Hadoop 3.x 时,如果还要连 Hive,那就得额外注意 Hive 连接器的版本匹配。
Hive 2.x 生态里的很多 jar 是按照 Hadoop 2 编译的,直接放到 Hadoop 3 环境里容易报UnsupportedOperationException或者各种 MethodError。我建议如果集群已经升级到 Hadoop 3,那 Hive 至少得是 3.1.x,同时在 Flink 的lib目录或者opt目录里放对应版本的 Hive 连接器。
比如 Flink 1.13 有flink-sql-connector-hive-3.1.2这样的预编译包,它对 Hive 3.1.x 和 Hadoop 3.x 的兼容性相对稳定。如果你自己从源码把 Flink 编译成适配 Hadoop 3.3.1 的版本,那 Hive 连接器也最好一起参与编译,或者确认这个连接器的 shaded 版本里没有和 Hadoop 3.3.1 冲突的内容。
检查 Hive 连接器是否正常,最直接的方式是在 Flink SQL Client 里建一个 Hive Catalog,然后执行一条简单的SHOW TABLES。如果这一步能出结果,说明 Flink、Hive、Hadoop 三者的底层 RPC 链路都是通的。如果卡在org.apache.thrift相关的报错上,通常是 Hive 版本与客户端 thrift 库不匹配,需要再去调整 Hive 连接器的版本。
6. 生产落地的稳定经验
6.1 把版本组合固定成标准模板
这次集成做完以后,我最大的体会是:不要在每次部署的时候现去搜 “Flink 1.13 Hadoop 3.x 怎么配”,而是把验证过的版本组合固定下来,形成公司内部的标准模板。
我自己的模板大概是:Flinkrelease-1.13分支 + Hadoop 3.3.1 + Hive 3.1.2 + Scala 2.12。这个组合我在测试环境反复验证过,本地模式能读写 HDFS,YARN 模式能稳定跑流任务,SQL Client 连接 Hive Catalog 也没有问题。后续如果再有新同事入职,我直接把这个编译好的 Flink 发行包让他拿走,比让他自己折腾要高效得多。
编译过一版之后,也可以顺手把flink-shaded-hadoop-3-uberjar 放进公司私有 Maven 仓库。因为 Flink 任务里如果直接依赖 Hadoop 的某些 API,在编写 Flink 程序时,Maven 坐标里的hadoop-client依赖也可以固定成和集群一致的版本,避免编译环境和运行环境不一致。
6.2 升级替换的先后顺序
大概率的量产环境升级路径是:先升级 Hadoop 集群到 3.x,然后 Flink 任务再切换客户端。但如果条件允许,我建议反过来,先在测试环境用编译好的 Flink 1.13 + Hadoop 3 客户端包跑一段时间的模拟流量,再观察 Hadoop 3 集群是否有异常日志。因为 Flink 客户端和 Hadoop 3 服务端的交互问题并不总是一开始就暴露,很多问题会在运行几小时后,比如 checkpoint 周期触发时,才体现出来。
还有一个经验是流量切的时候不要一把梭。可以挑两三个非核心的任务先切到新的 Flink 发行包上,运行两三天,重点看两点:HDFS 的写入 QPS 是否正常,YARN 上的 container 日志里有没有周期性抛出的异常。排掉这些异常后,再把剩余任务分批切过去,这样出问题影响面可控。
6.3 程序里如何正确声明依赖
Flink 任务工程的 pom.xml 里,如果要声明 Hadoop 3 相关依赖,我建议尽量标记为provided,或者代码里只用flink-shaded-hadoop-3-uber里已有的类。原因是 Flink 运行时已经通过发行包把 Hadoop 客户端带到了 classpath 中,如果任务 jar 里再重复打入一套 Hadoop 类,很容易触发类冲突。
我一般会在 pom 里这样写核心依赖:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-shaded-hadoop-3-uber</artifactId> <version>${flink.shaded.version}</version> <scope>provided</scope> </dependency>这里flink.shaded.version要和实际编译 Flink 时使用的版本保持一致。如果找不到精确版本,或者不想维护这个依赖,还有一种做法是只在代码里引用 Flink 公开的文件系统 API,不显式依赖 Hadoop 类,运行时让 Flink 自己去解析 HDFS 路径。比如用StreamingFileSink、FileSource这类内部封装好的 API,不用手动直接创建FileSystem实例,能有效减少编译期和运行期的依赖纠缠。
做完整套集成之后,我自己的感受是 Flink 1.13 集成 Hadoop 3.x 的难度并不在配置本身,而在你是不是真的理解了 Flink 发行包内置什么版本的 Hadoop 客户端,以及运行时 classpath 里到底有哪几份 Hadoop 类。把这两件事搞明白,不管是编译路线还是无 Hadoop 包路线,判断起来都很快,遇到报错也不会手足无措。希望这篇实际踩坑记录能让你少走点弯路。