尧图精选

CDH6.3.2上部署Flink1.13.1:版本适配、YARN提交与踩坑全解析

🕒 发布时间:2026/10/2 16:06:07 📁 来源:尧图网络
简介面向 CDH 6.3.2 大数据平台运维与实时计算开发人员用于解决 Flink 1.13.1 与 CDH 的版本兼容及集群分发部署问题是一份可直接落地的安装资源。资源共 5 个文件、约 299.52MB其中两个 Jar 分别对应 Flink 核心库和 YARN 客户端用于程序运行与任务提交Parcel 格式的二进制发行版可被 Cloudera Manager 直接识别配合 SHA 校验文件和 JSON 元数据能够安全地完成集群分发、版本校验与依赖匹配。已有 3495 人学习下载适合需要快速在 CDH 环境中部署 Flink、接入 YARN 资源调度并提交实时作业的工程师。资料包按 CDH Parcel 机制组织文件分工明确有助于读者理解 Flink 与 CDH 集成时的版本对应关系、校验流程及资源配置要点从而避免部署中常见的版本冲突和资源分配问题更稳妥地完成配置与运行验证。结合 YARN 日志与 Flink Web UI还能快速定位任务异常提升实时作业的运维效率。1. 为什么在 CDH6.3.2 上跑 Flink 1.13.1 要先解决版本之争做过 CDH 集群实时接入的人都会遇到同一个场景YARN 和 HDFS 每天都在正常跑 Spark、Hive 任务团队想上 Flink 做流式计算下载了 flink-1.13.1-bin-scala_2.12.tgz 解压后直接yarn-session.sh启动结果要么报NoClassDefFoundError: org/apache/hadoop/conf/Configuration要么任务起了一半卡在申请 container 上。这不是 Flink 坏了而是 CDH6.3.2 的 Hadoop 是 Cloudera 基于 3.0.0 深度定制的版本Flink 1.13.1 官方预编译包默认不携带任何 Hadoop 发行版依赖两者之间缺一层版本适配。这篇笔记会从版本匹配讲起手把手带你在 CDH6.3.2 上部署 Flink 1.13.1覆盖 Hadoop 依赖注入、YARN 提交、Hive 集成和排查手段适合正在做落地选型或已经被类冲突折磨了一周的工程师。2. 版本匹配与前置检查CDH6.3.2 到底缺了什么2.1 CDH6.3.2 组件版本清单与 Flink 1.13.1 的兼容性边界先在 CM 页面上核对你的集群版本。CDH6.3.2 的核心组件版本大致如下我列的是最常见的一版具体以你的 CM 显示为准组件CDH6.3.2 版本Flink 1.13.1 的期望兼容性结论Hadoop3.0.0-cdh6.3.2官方支持 Hadoop 2.8/2.9/3.1关键矛盾点API 兼容但类名可能冲突Hive2.1.1-cdh6.3.2Flink Hive 集成需要 Hive 2.x可以配合但要留意 Hive 依赖导入顺序JavaOpenJDK 8必须 Java 8完全兼容Scala2.11Spark 自带Flink 1.13.1 需要 Scala 2.12互不影响Flink 单独使用 Scala 2.12 版本Zookeeper3.4.5Flink 不直接依赖 ZK无冲突Flink 1.13.1 的源码根 pom 里默认使用 Hadoop 2.8.3但你拿到的是 CDH6.3.2 的 Hadoop 3.0.0。这带来两个直接问题第一Flink 的flink-shaded-hadoop官方 jar 不包含 CDH 的定制补丁比如某些 RPC 协议、HA 配置解析直接拿通用 Hadoop 3.1.1 的 shaded jar 去连 CDH 的 NameNode可能在 RPC 握手时报NoSuchMethodError或AbstractMethodError。第二CDH 的 HDFS 客户端依赖很多 Cloudera 修改过的类但 Maven 官方仓库没有这些 artifact只能从 Cloudera 仓库拉取。所以我们不能直接拿官方发行包里的任何flink-shaded-hadoop-*-uber去替代正确的做法是构建一个CDH 定制版 hadoop-client fat jar放到 Flink 的 lib 目录下。这是后续所有步骤的地基。2.2 用 Maven 生成一个兼容 CDH6.3.2 的 hadoop-uber jar常见做法是创建一个最小 Maven 工程把hadoop-client的 CDH 版本依赖打成一个 uber jar然后替换 Flink lib 下缺失的 Hadoop 支撑。这一步可以避免重新编译整个 Flink省去大量时间。我们需要先确认本机有 Maven 3.5然后新建目录并放入下面这个pom.xml。project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion groupIdcom.example/groupId artifactIdcdh-hadoop-uber/artifactId version1.0/version packagingjar/packaging repositories repository idcloudera/id urlhttps://repository.cloudera.com/content/repositories/releases//url /repository /repositories dependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.0.0-cdh6.3.2/version /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.2.4/version executions execution phasepackage/phase goals goalshade/goal /goals configuration filters filter artifact*:*/artifact excludes excludeMETA-INF/*.SF/exclude excludeMETA-INF/*.DSA/exclude excludeMETA-INF/*.RSA/exclude /excludes /filter /filters /configuration /execution /executions /plugin /plugins /build /project这段 pom 的核心逻辑有两个一是把 Cloudera 仓库加进依赖源确保能找到3.0.0-cdh6.3.2版本的hadoop-client二是用 shade 插件在打包时将所有传递依赖合并到一个 jar 里。过滤掉 META-INF 下的签名文件是因为合并 jar 时会遇到多个 jar 的数字签名冲突导致 JVM 报SecurityException。执行mvn package -DskipTests后在target目录下会生成cdh-hadoop-uber-1.0.jar。这个 jar 包含了 HDFS、YARN、MapReduce 的客户端类以及 CDH 的配置文件解析逻辑。复制到 Flink 的 lib 目录这一步放在第 3 章先不急着操作。2.3 环境变量与基础配置检查在动手改 Flink 配置之前先在 CDH 节点上确认三件事。用下面的命令逐条检查任何一条不满足都可能让后续任务起不来。java -version # 必须是 Java 8CDH 6.3.2 修改过某些 JVM 参数Java 11 会直接失败 echo $JAVA_HOME # 需要指向 JDK8 安装路径比如 /usr/local/jdk1.8.0_281 which hadoop # 确认 hadoop 命令来自 /opt/cloudera/parcels/CDH/bin/hadoop hadoop classpath # 打印 CDH 的完整 classpath后面诊断时会用到如果which hadoop显示的是/usr/bin/hadoop说明你当前环境里装的是 Apache Hadoop 客户端而不是 CDH 的 parcels后续 Flink 任务会把配置解析成默认的本地文件系统而不是 CDH 的 HDFS。解决办法是修改~/.bashrc把 CDH 的 bin 目录加到 PATH 前面。export JAVA_HOME/usr/local/jdk1.8.0_281 export PATH/opt/cloudera/parcels/CDH/bin:$PATH export HADOOP_CONF_DIR/etc/hadoop/conf export HADOOP_CLASSPATH$(hadoop classpath)这里的HADOOP_CONF_DIR尤其重要。Flink 1.13.1 的 YARN 客户端在启动时会通过它读取core-site.xml和hdfs-site.xml才能知道 NameNode 的地址和 HA 配置。CDH 的配置文件通常在/etc/hadoop/conf你可以在 CM 页面的 HDFS 服务实例里找到确认。3. 把 Flink 1.13.1 装进 CDH6.3.2三步部署与配置项解析3.1 解压 Flink 与 lib 目录准备找一个专门的部署目录比如/opt/flink。下载flink-1.13.1-bin-scala_2.12.tgz后解压然后把你构建的 CDH hadoop-uber jar 放进去。这里有一个需要特别注意的细节不要把官方flink-shaded-hadoop-2-uber和 CDH jar 同时放进 lib否则会出现两个版本类加载顺序谁在前完全看 JVM 心情这类问题很难排查。mkdir -p /opt/flink tar -zxvf flink-1.13.1-bin-scala_2.12.tgz -C /opt/flink cp cdh-hadoop-uber-1.0.jar /opt/flink/flink-1.13.1/lib/复制完成后看一下lib目录的文件列表确认里面有一个cdh-hadoop-uber-1.0.jar并且没有flink-shaded-hadoop-2-uber-*.jar。Flink 1.13.1 的预编译包本身不带 Hadoop jar如果你之前加过什么别的 Hadoop 依赖一并删掉。这一步我见过很多人翻车他们觉得官方 jar 更稳定留着多个 Hadoop jar结果任务启动时 HDFS 客户端加载了旧版本写入文件直接抛EOFException。3.2 flink-conf.yaml 必须调的参数Flink 1.13.1 的配置在conf/flink-conf.yaml。不要全用默认值我根据 CDH6.3.2 的环境整理了几个必调项。先用一个内存充足的中型节点为例。jobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 4096m taskmanager.numberOfTaskSlots: 2 parallelism.default: 2 classloader.resolve-order: parent-first env.java.opts: -Djava.library.path/opt/cloudera/parcels/CDH/lib/hadoop/lib/native逐项说明jobmanager.memory.process.sizeJobManager 的堆外内存。CDH 节点上往往还跑着 DataNode 和 NodeManager不要贪大我先给 2G。taskmanager.memory.process.sizeTaskManager 的进程总内存。这里 4G 是一个能跑通 WordCount 的最小值实际生产按每任务的常驻内存来估算。taskmanager.numberOfTaskSlots每个 TaskManager 的 slot 数默认是 1。我习惯设成 2这样单流任务可以同时跑 source 和 sink也不会因为 slot 太少导致 idle。parallelism.default没在任务里显式指定并行度时兜底值。如果集群有 10 个节点每个节点 2 slot可以设成 20但前期调通建议先 2。classloader.resolve-order默认是child-first在 CDH 上会踩坑因为 Flink 自身用了一部分旧 Guava而 Hadoop 3.0 使用的是新 Guava。改成parent-first可以优先加载 lib 下 CDH 的类避免冲突。env.java.opts-Djava.library.path是给 HDFS 客户端加载 native 库用的不配的话 HDFS 读写会走纯 Java 实现性能差很多而且某些 HA 场景会报Unable to load native-hadoop library。如果你后续要接 Kerberos还需要追加一行-Djava.security.krb5.conf/etc/krb5.conf。不过我在第 5 章会单独讲坑这里先不给。3.3 验证 HDFS 与 YARN 连通性配置完成后不要急着起 Flink先手动测一下 CDH 的 HDFS 和 YARN 是否真的认这套配置。用下面的命令创建目录并跑一个简单的 distcp确认没有异常。hadoop fs -mkdir -p /tmp/flink-test hadoop fs -put /etc/hosts /tmp/flink-test/hosts hadoop jar /opt/cloudera/parcels/CDH/lib/hadoop-mapreduce/hadoop-mapreduce-examples.jar wordcount \ /tmp/flink-test/hosts /tmp/flink-test/out最后一条命令如果报Could not find or load main class不用慌这是 CDH 6.3.2 偶尔会出现的 examples jar 路径问题。你只需要确认前两条命令能正常执行并且hadoop fs -ls /tmp/flink-test能看到hosts文件。同时用yarn node -list确认 NodeManager 都在运行。这里还要做一次物理隔离验证在部署 Flink 的节点上执行hadoop fs -get /tmp/flink-test/hosts /dev/null不要报错。有些环境把 Flink 部署在了没有 HDFS 客户端库的纯工具节点上导致后面提交任务时一堆java.net.UnknownHostException提前能发现。4. 在 YARN 上跑通第一个 Flink 作业Per-Job 模式与 HDFS 数据检查4.1 选择 yarn-session 还是 per-job 模式Flink 1.13.1 在 CDH6.3.2 上提交任务有两种主流模式yarn-session和per-job。yarn-session先在 YARN 上启动一个常驻的 Flink 集群然后向这个集群提交多个作业per-job则是每次提交时单独申请一个 YARN Application跑完自动释放。我在这套 CDH6.3.2 环境上强烈建议用per-job。原因是 CDH 的 YARN 队列通常已经给 Spark 和 Hive 划分好了资源常驻 session 会一直占着队列不放容易引发资源死锁。而且 Flink 1.13.1 的yarn-session模式需要单独维护一个 ApplicationMaster排查问题多一层黑匣子不如让一次作业一个 Application 来得直观。提交命令如下我以 Flink 官方示例中的SocketWindowWordCount为例因为它不需要 HDFS 数据文件可以随时开跑。4.2 提交 SocketWindowWordCount 的完整命令与日志观察先在 CDH 任意节点上启动一个 netcat 端口作为数据源nc -l 9999然后提交 Flink 作业到 YARN/opt/flink/flink-1.13.1/bin/flink run \ -m yarn-cluster \ -ynm flink-cdh-wordcount \ -yjm 2048m \ -ytm 4096m \ -ys 2 \ -p 2 \ /opt/flink/flink-1.13.1/examples/streaming/SocketWindowWordCount.jar \ --hostname 你执行nc的节点IP --port 9999参数解释-m yarn-cluster告诉 Flink 客户端使用 YARN 作为资源管理器。-ynm flink-cdh-wordcountYARN Application 的名字方便在 ResourceManager 上识别和 kill。-yjm 2048mYARN 模式下的 JobManager 内存覆盖 flink-conf.yaml 里的对应配置。-ytm 4096mTaskManager 内存。-ys 2每个 TaskManager 的 slot 数要和taskmanager.numberOfTaskSlots一致。-p 2作业并行度。提交后你会看到类似下面的日志这意味着作业注册成功2025-XX-XX 10:00:00,123 INFO org.apache.flink.yarn.YarnClusterDescriptor - YARN application has been deployed successfully. 2025-XX-XX 10:00:01,456 INFO org.apache.flink.runtime.dispatcher.Dispatcher - JobManager is now running此时去 YARN ResourceManager 的 Web UI 页面输入作业名flink-cdh-wordcount应该能看到RUNNING状态点进去能看到 Flink 的 Dashboard 链接。再回到 netcat 终端输入一行hello flinkFlink 的聚合日志里很快会出现(hello,1)和(flink,1)的输出。这表示 HDFS、YARN、网络栈、类加载全部打通。日志查找路径在 YARN 的 container 日志目录常见位置是/opt/cloudera/parcels/CDH/log/userlogs/application_xxx/container_xxx/stdout。如果找不到可以直接在 YARN Web UI 的对应 container 内查看 stdout。4.3 通过 Flink Web UI 和 YARN ResourceManager 查看状态Flink 1.13.1 的 Web UI 默认端口是 8081但提交到 YARN 后Dashboard 端口会随机分配。你可以在 ResourceManager 页面点击 Application 的ApplicationMaster链接跳转到 Flink Dashboard。这个界面有几个地方值得重点看Job Manager标签下的Memory确认堆内和堆外使用率是否和-yjm 2048m对得上。Task Managers标签下的Number of Task Slots确认是不是ys * 并行度这个结果。Checkpoints标签在任务运行一段时间后会显示 Checkpoint 是否成功。如果一直处于PENDING多半是 HDFS 端写 Checkpoint 的目录没有权限我在第 5 章会展开。另一个验证方式是去 HDFS 查看 Flink 在 YARN 模式下的临时目录hadoop fs -ls /tmp/flink-*。正常情况下提交过作业后会看到几个以flink-开头的 Application 专属目录。如果这个目录在 YARN 上访问受限也可以在命令行用hadoop fs -ls /user/yarn-user/.flink检查。5. CDH 上跑 Flink 的 5 个高频踩坑现象、原因与解决5.1 提交时报 NoClassDefFoundError: org/apache/hadoop/conf/Configuration错误现象flink run -m yarn-cluster启动后客户端进程直接抛NoClassDefFoundError: org/apache/hadoop/conf/Configuration然后退出。原因Flink 客户端的 bin 脚本在执行时classpath 里没有加入 lib 下的 CDH jar。虽然你已经把cdh-hadoop-uber-1.0.jar放进了 lib但 Flink 1.13.1 的启动脚本在 org.apache.flink.yarn 包加载前会先加载部分 Hadoop 类而 JVM 的 classpath 扫描顺序是数组顺序lib 目录外没有提前注入。解决在conf/flink-conf.yaml中追加一行env.java.opts: -Xbootclasspath/a:/opt/flink/flink-1.13.1/lib/cdh-hadoop-uber-1.0.jar。这个参数把 CDH jar 放进了 Bootstrap classpath保证最早加载。注意这种方式有些暴力但在 YARN 客户端模式下它是可靠的生产环境如果担心可以用-Dflink.lib结合YARN_CONTAINER_CLASSPATH来配置。5.2 冷启动卡在 ApplicationMaster 注册阶段错误现象提交后 FLink 客户端输出YARN application has been deployed successfully.但 ApplicationMaster 一直处于NEW状态迟迟不进入RUNNINGYARN 日志里反复出现ApplicationMaster is not registered。原因这是我在 CDH6.3.2 上踩过最久的一次。原因是yarn.resourcemanager.am.max-attempts和 Flink 的yarn.application-attempts参数不匹配。CDH 的 YARN 默认限定了 application master 的最大重启次数为 2而 Flink 1.13.1 默认要求 3 次。解决在 flink-conf.yaml 里显式设置yarn.application-attempts: 2另外检查yarn-site.xml里的yarn.nodemanager.pmem-check-enabled如果 CDH 开的是物理内存强校验而你给 TaskManager 的内存大于 NodeManager 的单容器上限也会卡在注册。可以把taskmanager.memory.process.size降低到 NodeManager 允许范围或者修改 CDH 队列配置。5.3 从 Hive 表读数据时报 AbstractMethodError: org.apache.hadoop.hive.ql.Driver错误现象用 Flink SQL 建 HiveCatalog 并读取 CDH 的 Hive 表提交后 TaskManager 报AbstractMethodError: org.apache.hadoop.hive.ql.Driver.getQueryPlan之类的异常。原因CDH6.3.2 的 Hive 是 2.1.1-cdh6.3.2但 Flink 1.13.1 的 Hive 集成包默认针对 Apache Hive 2.1.0 编译。Hive 2.1.1 新增了一个接口方法导致二进制不兼容。另外 Flink lib 里还混进了 CDH 自己的 hive-exec jar和 Flink 的 hive-exec 版本相同但类有细微差异。解决把 CDH 的 hive-exec jar 排除只使用 Flink 发行包里的 hive-connector。在 porm 里做精准排除或者干脆把lib/目录下的flink-sql-connector-hive-2.1.1_2.12-1.13.1.jar里的 Hive 类作为唯一来源。我的经验是在 Flink SQL 初始化时显式指定版本CREATE CATALOG myhive WITH ( type hive, hive-conf-dir /etc/hive/conf, hive-version 2.1.1 );同时删除 lib 下所有 hive-exec-cdh.jar只保留 Flink 自带的。5.4 Kerberos 认证环境下的期限与票据问题错误现象集群开了 KerberosFlink 作业刚启动时正常运行半小时后所有 source/sink 都报javax.security.auth.login.LoginException: No valid credentials provided。原因Flink 的 Kerberos 认证默认只主动登录一次生成的 tgt 票据有时效性。CDH 的 Kafka/HDFS 通常要求票据在 24 小时内有效但 Flink 任务的长时间运行会超过这个时限而 Flink 1.13.1 的security.kerberos.login.use-ticket-cache默认值在某些部署下不会自动续期。解决在 flink-conf.yaml 里加security.kerberos.login.contexts: Client,KafkaClient security.kerberos.login.keytab: /etc/flink/flink.keytab security.kerberos.login.principal: flinkEXAMPLE.COM另外为 Flink 进程单独创建一个定期 renew 的 cron 任务手动刷新 keytab 对应的票据这也是很多 CDH 运维团队在用的兜底方案。5.5 HDFS 小文件导致 Checkpoint 频繁超时错误现象作业运行一段时间后Checkpoint 开始失败失败原因是Exception while performing checkpointTaskManager 日志里出现大量Connection refused或Retries exhausted。原因CDH 集群中你的 HDFS 用户目录下有很多 Spark 作业产生的小文件。Flink 的 Checkpoint 默认写到 HDFS 的/tmp/flink目录该目录下的文件数量过多时NameNode 响应变慢导致一系列 RPC 超时。我在 CDH 上遇到过因为/tmp/flink下有超过 10 万个小文件整个集群 GC 都受影响的案例。解决把 Flink 的 Checkpoint 基础目录改到一块数据盘上专用的路径并配置差分目录state.checkpoints.dir: hdfs://nameservice1/data/flink-checkpoints state.savepoints.dir: hdfs://nameservice1/data/flink-savepoints再配合定期清理 HDFS 回收站的策略让 Checkpoint 目录的文件数量控制在万级以内。如果仍然慢考虑把 state.backend 从rocksdb切换到filesystem两个后端对这种目录的扫描方式不同。6. 用火焰图验证 Flink 作业性能一条命令定位 CPU 热点在 CDH6.3.2 上Flink 作业能跑起来只是第一步真正要投产还需要知道它有没有在某个算子卡死。常规监控看 Flink Dashboard 的 CPU 和网络吞吐但定位到具体方法就得靠火焰图。这里介绍一下我常用的基于 AsyncProfiler 的做法前提是你能在 TaskManager 节点上有 root 或同用户权限。先下载 AsyncProfiler它的 release 包是一个 tar.gz解压后直接有profiler.sh脚本。不要用系统自带的 perf 加 JFR对 Flink 这种长时间运行的服务AsyncProfiler 的挂载方式更安全不需要重启 TaskManager。# 找到你想采样的 TaskManager JVM 进程 PID jps -l | grep TaskManager # 采样 60 秒输出到 /tmp/flamegraph.html ./profiler.sh -d 60 -e cpu -o flamegraph \ -f /tmp/flamegraph_$(date %s).html PID这条命令的含义是-d 60表示采样 60 秒-e cpu表示统计 CPU 周期而非 wall 时间-o flamegraph直接生成 HTML 格式火焰图-f指定输出路径。执行完成后用浏览器打开 HTML 文件你会看到一层层的函数调用栈顶部越宽说明 CPU 在窗口期间消耗越多的函数。实际解读时我每次都会先看三个地方如果宽栈出现在sun.nio.ch或者org.apache.flink.runtime.io.network.netty.NettyServer说明网络收发占了主导主要矛盾在带宽而不是算子逻辑。如果java.util.HashMap或org.apache.flink.runtime.state.heap.HeapKeyedStateBackend很宽说明 state 读写竞争激烈考虑增大 taskmanager 内存或换 RocksDB。如果org.apache.hadoop.hdfs.DFSClient相关调用栈宽说明 HDFS 读写 IO 阻塞先看数据本地性和小文件而不是调并行度。这套方法对 CDH 集群尤其有价值因为很多 Flink 性能问题根源不在 Flink 本身而是 CDH 的 HDFS 或 YARN 资源竞争。通过火焰图你能快速分辨是磁盘、网络还是 GC 导致不靠猜。我自己的习惯是每调整一组参数后重新采样对比火焰图宽度变化如果某个栈的占比从 30% 降到 10%再继续下一步优化否则宁可不动。另外提醒一句AsyncProfiler 不能在 JDK8 的某些小版本上直接使用请在采样前先uname -a和java -version确认系统架构和 JDK 版本匹配否则会报Failed to load shared library。这本质上也是我们做技术落地常说的版本玄学希望这篇笔记能帮你在这个坑上少耗半天顺利把 Flink 1.13.1 跑在 CDH6.3.2 上。本文还有配套的精品资源点击获取
上一篇/下一篇内容由系统自动关联 返回资讯列表 →