十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink 1.13集成Hadoop 3.x:版本冲突解决与源码编译实战

Flink 1.13集成Hadoop 3.x:版本冲突解决与源码编译实战 这些年在公司做实时数仓经常要面对的一个问题就是 Flink 和 Hadoop 的版本配套关系。Flink 1.13 这个版本用得人不少但当你拿着官方下载的 Flink 1.13 包去对接一个 Hadoop 3.x 集群时十有八九会碰一鼻子灰报错信息千奇百怪核心原因只有一个默认发行包根本不是按照 Hadoop 3 给你准备的。这篇文章就专门说清楚 Flink 1.13 集成 Hadoop 3.x 到底该怎么解决从版本冲突的原理解析到两条可落地的操作路线再到运行期会踩的那些坑我尽量按实际项目里能直接抄作业的方式来写。我当时的场景比较典型生产环境 HDFS 已经升级到 Hadoop 3.3.xYARN 也换了新版本但实时任务的代码还跑在 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_CLASSPATH2.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 33.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.0Hive 相关模块默认版本是 3.1.2。如果你集群是 Hadoop 3.2.x直接用内置 profile 最省心不用改任何 pom 文件。如果你的集群是 Hadoop 3.3.x那就需要在编译命令里额外指定-Dhadoop.version3.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 profilemvn 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.version3.3.1这里-Dhadoop.version3.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 文件路径的任务。比如先准备一个文本文件放到 HDFShadoop 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 引发的 NoSuchMethodErrorFlink 和 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。这个组合我在测试环境反复验证过本地模式能读写 HDFSYARN 模式能稳定跑流任务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 groupIdorg.apache.flink/groupId artifactIdflink-shaded-hadoop-3-uber/artifactId version${flink.shaded.version}/version scopeprovided/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 包路线判断起来都很快遇到报错也不会手足无措。希望这篇实际踩坑记录能让你少走点弯路。
返回列表