ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

CDH 6.3.2 部署 Flink 1.13.1 实战:Classpath 与 YARN 配置全解析

CDH 6.3.2 部署 Flink 1.13.1 实战:Classpath 与 YARN 配置全解析 简介面向需要在CDH 6.3.2集群中集成Apache Flink的运维与开发人员这份资源包提供了Flink 1.13.1的完整分发与部署组件解决了在Cloudera企业大数据平台上安装、配置和运行Flink的兼容性问题。压缩包共5个文件包含约299.52MB的核心内容主要文件类型有Parcel二进制包、YARN客户端Jar包、Flink核心库Jar包以及Parcel元数据JSON和SHA校验文件覆盖从分发校验到YARN提交作业的完整链路。目前已有3492人学习下载说明该部署方案具有较高的参考价值。借助其中的Parcel文件可快速导入并激活到CDH集群通过YARN客户端Jar包实现作业提交与资源调度同时SHA文件可确保包体完整性适合需要高效落地Flink on YARN的生产环境参考。1. 为什么 Flink 1.13.1 放进 CDH 6.3.2 不是解压就能用我第一次把 flink-1.13.1 部署到 CDH 6.3.2 集群时提交命令按下去不到十秒就迎面撞上一个标准的NoClassDefFoundError。群里问了一圈十个人有八个让我“换个版本试试”但没人说清楚问题本质CDH 的 Hadoop 不是社区 HadoopFlink 的发行包和 YARN 之间需要一座桥。这座桥没搭好换多少个 Flink 版本都是一样的结果。很多从纯 Apache Hadoop 社区版转过来的人会默认“Flink 解压就能连接 HDFS/YARN”。这个预期在 CDH 上不成立原因不是 Flink 有毛病而是 CDH 6.3.2 的组件体系比你想象的要“封闭”。1.1 CDH 的 Hadoop 版本和社区版不是一回事CDH 6.3.2 自带的 Hadoop 版本是3.0.0-cdh6.3.2。这个版本在社区 Hadoop 3.0 的基础上打了很多 Cloudera 自己的补丁很多客户端路径、配置项、依赖传递关系也和社区版有差异。Flink 1.13.1 的官方二进制发布包集成验证时主要对标的是 Apache Hadoop 2.x/3.x 的开放接口不会专门拿 Cloudera 的补丁版做回归测试。这意味着什么意味着你在 CDH 上运行 Flink必须让 Flink 的 JVM 进程在启动时看到 CDH 的 Hadoop 客户端 jar而不是随便拿一个社区版 hadoop-client 顶上去。否则即使能启动 JobManager一旦读写 HDFS 或向 YARN 申请资源就会出现各种各样的 ClassNotFound、NoSuchMethodError。还有一个更隐蔽的点CDH 6.3.2 这种基于 parcel 的管理方式会把 Hadoop 客户端拆到/opt/cloudera/parcels/CDH/lib/hadoop/等路径下而 Flink 默认是不会自动去这些目录找 jar 的。所以这本质上不是一个“安装 Flink”的问题而是一个“JVM classpath 供应链”的问题。1.2 Flink 运行期需要的是“客户端 classpath”不是一个安装包Flink 与 Hadoop 的耦合点其实只有三个HDFS 文件系统、YARN 客户端、Kerberos/DelegationToken 认证。这三块在 Flink 编译期都是 optional 的运行期却必须有类可加载。因此 Flink 的每个进程——客户端、JobManager、TaskManager——都需要在 classpath 里找到 Hadoop 相关类。很多人的做法是在flink/lib里塞一个flink-shaded-hadoop-3-uberjar 解决所有问题这种方式确实简单粗暴。但也有人只设置HADOOP_CLASSPATH环境变量也能跑得稳稳当当。两条路线分别对应不同的坑后面我会展开讲。这里先确立一个概念Flink on CDH 的部署工作七成是在解决 classpath三成才是改配置。2. 动手前先清点环境JDK、Hadoop Gateway 和 YARN 队列与其等报错再回头补环境不如在解压 Flink 之前先把三个前置条件检查一遍。这三件事通常不会写在 Flink 官方文档里但我在 CDH 6.3.2 上踩过坑之后每次都先做。2.1 JDK 版本不要无脑追新Flink 1.13.1 官方支持 JDK 8 和 JDK 11但 CDH 6.3.2 自身的 HDFS、YARN、Hive 等服务基本都是在 JDK 8 上编译和验证的。生产环境最稳的组合是给 Flink 客户端和所有 YARN 容器都统一用 JDK 8比如/usr/java/jdk1.8.0_181-cloudera。有的节点装过高版本 JDK客户端在本地提交时没毛病但 YARN 容器启动时如果解析到了另一套 JDK就可能出现莫名其妙的UnsupportedClassVersionError或者 GC 参数不兼容。我的建议是在flink-env.sh里显式设置JAVA_HOME不要依赖节点的默认环境变量。env.java.opts里也可以把-Djava.library.path指到/opt/cloudera/parcels/CDH/lib/hadoop/lib/native避免后续访问 HDFS 时报本地库相关的警告。2.2 三条命令确认集群底子在部署前先找一台需要安装 Flink 的边缘节点执行三条命令hadoop version echo $HADOOP_CONF_DIR hadoop classpath | tr : \n | grep cdh | head -n 20第一条确认 Hadoop 客户端确实是 CDH 6.3.2 的版本第二条确认核心配置目录存在第三条验证 CDH 的 jar 路径能通过hadoop classpath正常输出。如果第三条返回空后面的 Flink 大概率连 HDFS 都连不上。还要确认/etc/hadoop/conf目录下至少存在core-site.xml、hdfs-site.xml、yarn-site.xml三个文件。CDH 的/etc/hadoop/conf经常是指向 Cloudera 动态生成的配置目录的软链缺文件时 Flink 在提交 YARN 任务时拿到的是一个残缺集群配置报错会很晚定位成本非常高。2.3 边缘节点要有 Hadoop Gateway 角色这一步很容易漏。只装了 CDH 的 Agent不一定有 Hadoop Gateway 角色。没有 Gateway节点上连hadoop命令都可能不存在更别谈hadoop classpath能否输出 CDH 的 jar 路径。在 Cloudera Manager 里找到集群对应的主机添加角色实例时选择 Gateway部署到 Flink 所在节点即可。这个角色本身不启动常驻服务只是把 Hadoop 客户端配置和库装好。装完之后再跑一次上面的三条命令确认输出正常再开工。整个过程五分钟能省掉后面至少半天的排查时间。3. CDH 6.3.2 上安装 Flink 1.13.1 的四步实操环境准备好之后真正的安装其实不复杂。以下是我在 CDH 6.3.2 上部署 flink-1.13.1 时总结的四步流程每步都有取舍逻辑。3.1 下载 flink-1.13.1-bin-scala_2.12.tgz 并做本地离线分发下载时选flink-1.13.1-bin-scala_2.12.tgz不要选 Scala 2.11 版本。Flink 1.13 是 Scala 2.12 的分水岭后续很多 connector 和用户代码都默认按 2.12 编译选 2.11 后续会遇到版本匹配问题。如果集群环境是离线的你需要先在有外网的机器上把 tar 包下载好传到/opt/flink或自定义目录。CDH 6.3.2 本身做离线安装时需要提前准备 parcel 和本地 HTTP 源但 Flink 没有官方 parcel直接 tar 分发即可。建议保留一份原始 tar 包作为离线仓库分发时不要直接解压到运行目录而是先解压到临时目录再调整配置最后复制到各节点。这里有个认识误区要纠正解压后 Flink 的lib目录里没有 Hadoop jar不是因为你下载了“不含 Hadoop”的特殊版本而是从 Flink 1.13 开始官方发布包默认就不捆绑 Hadoop。所以下一步要做的不是找“完整版 Flink”而是主动补齐 Hadoop classpath。3.2 第一道选择题HADOOP_CLASSPATH 还是 flink-shaded-hadoop-3-uber这是整个部署中最关键的一个选择两条路各有适用场景。方式 A使用HADOOP_CLASSPATH环境变量。在flink-env.sh里加export HADOOP_CONF_DIR/etc/hadoop/conf export HADOOP_CLASSPATH$(hadoop classpath)这种方式的好处是完全复用 CDH 的 Hadoop 客户端不需要额外下载任何 Hadoop 相关 jarCDH 升级后只要重新执行一次命令即可。缺点是如果 Flink 提交到 YARN 后JobManager 和 TaskManager 容器没有继承这个环境变量就会在运行时缺类。方式 B下载flink-shaded-hadoop-3-uberjar 放进flink/lib。这种方式能让每个 Flink 进程在启动时自动加载 Hadoop 类不依赖环境变量传递YARN 模式下最省心。但要注意这个 uber jar 里的 Hadoop 版本和 CDH 补丁版本可能有差异一旦用到 CDH 特有补丁的接口会出现NoSuchMethodError这类问题。我的经验是先用方式 A 跑通一个简单任务如果 TaskManager 容器里反复报 Hadoop 类缺失再考虑方式 B 作为兜底。两种方式不要同时上否则 lib 和 classpath 各有一份 Hadoop 类加载顺序不可控。3.3 修改 flink-conf.yaml 前先算 YARN 容器内存flink-conf.yaml里最常改的是内存配置但这里不是随便填。先看 CDH 的 YARN 资源配置尤其注意yarn.scheduler.maximum-allocation-mb和yarn.nodemanager.resource.memory-mb。如果 TaskManager 申请的内存超过上限Flink 会反复报Could not allocate container而且日志在 Flink 侧显示得不够直观容易误判为网络问题或 YARN 故障。一个常用的保守算法容器最大 8GB任务并行度 4 个 slot可以设置jobmanager.memory.process.size: 1536m taskmanager.memory.process.size: 6144m taskmanager.numberOfTaskSlots: 4不要把taskmanager.memory.process.size直接顶到 8GB系统本身还要留一些给 JVM overhead 和 YARN 容器本身。另外内存单位必须写m或g只写数字在部分配置解析中会被忽略最终使用默认值导致后续行为完全不符合预期。如果状态量较大建议开启 RocksDB 作为 state.backendstate.backend: rocksdb这行配置在 YARN 模式下不会直接带来内存问题但能避免默认的 HashMapStateBackend 把大量状态塞进堆内存。3.4 是否让 Cloudera Manager 接管 Flink很多从 CM 管理 Spark 过来的人会习惯性地想在 CM 里也看到 Flink 服务。但 CDH 6.3.2 默认不管理 Flink也没有官方 parcel。强行用 CSD 脚本接入 CM维护成本高于收益而且 Flink 如果是按 YARN 模式运行它的 JobManager 和 TaskManager 生命周期由 YARN 调度并不需要像 HDFS 那样常驻服务。我通常把 Flink 当作一个“外部计算客户端”来维护在边缘节点放一份 Flink 目录用脚本管理版本和环境变量提交任务时通过 YARN 获取资源。这样 CDH 集群的稳定性、升级节奏都不受影响Flink 版本升级也只是替换目录的事。4. YARN 上提交作业per-job 和 application 怎么选Flink 1.13.1 在 CDH 6.3.2 上最常用的两种 YARN 提交模式是 per-job 和 application。它们的区别不只是命令不同对集群运维方式也有影响。4.1 提交命令和最容易漏掉的 HADOOP_CONF_DIRper-job 模式提交命令./bin/flink run -t yarn-per-job \ -Dyarn.application.nameMyFlinkJob \ -p 2 \ -c com.example.MyJob \ /path/to/your-job.jarapplication 模式提交命令./bin/flink run-application -t yarn-application \ -Dyarn.application.nameMyFlinkJob \ -c com.example.MyJob \ /path/to/your-job.jar两者的核心差异在 main() 方法执行的位置。per-job 模式下客户端在自己的 JVM 里执行 main()然后为这个任务单独启动一个 JobManagerapplication 模式下main() 在 YARN 里的 ApplicationMaster 中执行客户端提交完就可以退出。对于一次性跑批任务application 模式更符合 YARN 的作业语义。无论哪种模式都必须确保HADOOP_CONF_DIR和HADOOP_CLASSPATH在提交客户端可用。尤其是HADOOP_CONF_DIRFlink 找 yarn-site.xml 依赖它。有人只设了HADOOP_CLASSPATH结果 Flink 一直连接不到 ResourceManager就是因为缺少核心配置文件目录。对比项per-jobapplicationmain() 执行位置客户端YARN ApplicationMaster客户端退出后任务不受影响不受影响适用场景交互式调试、临时任务生产定时任务资源申请时机客户端触发AM 触发4.2 一个真实失败案例NoClassDefFoundError 的完整排查链路有一次我在 CDH 6.3.2 上提交一个稍复杂的作业本地 IDE 运行正常但一换成-t yarn-per-job就报java.lang.NoClassDefFoundError: org/apache/hadoop/yarn/api/records/ApplicationId。我的排查链路是这样的先确认客户端环境变量是否生效。执行echo $HADOOP_CLASSPATH如果为空说明flink-env.sh里的 export 没被加载。再确认 Flink 的启动脚本是否拿到了 Hadoop classpath。可以用./bin/flink info会加载客户端环境如果这个命令都报 Hadoop 类缺失就是提交客户端的问题。如果客户端正常但作业跑到一半 TaskManager 报同样的 ClassNotFound说明容器环境没有继承变量。此时需要检查flink-conf.yaml中是否配置了yarn.taskmanager.env.HADOOP_CLASSPATH或对应前缀的环境变量透传不同 Flink 小版本前缀略有差异以 1.13.1 版本实际支持的 YARN 配置项为准。用yarn logs -applicationId appId查看 YARN 容器日志重点看 TaskManager 启动命令里的 classpath 是否包含 CDH 的 hadoop 路径。最后确认是否误删了 Flink 自带的flink-shaded-hadoop相关 jar导致 lib 目录里缺基础依赖。最终根因是容器环境没有拿到 HADOOP_CLASSPATH客户端环境正常给了我一种“配置已生效”的错觉。这也是 Flink on CDH 最常见的隐蔽问题本地没问题客户端没问题容器一启动就缺类。5. 跑通后要盯住的三个高频问题第一课跑通之后真正磨人的是这些边边角角的问题。我把高频的三种现象整理出来每个都对应一个明确的处理方向。5.1 lib 目录里出现双份 Hadoop 类时的冲突规律如果你按方式 B 在lib里放了flink-shaded-hadoop-3-uber同时又习惯性地在flink-env.sh里设置了export HADOOP_CLASSPATH那么 JobManager/TaskManager 的 classpath 里会存在两套 Hadoop 类。表面症状通常是某个接口出现NoSuchMethodError并且只在特定方法调用时触发。这类问题最麻烦的地方在于日志不直接经常被误认为用户代码问题。处理原则很简单二选一。用方式 B 就把HADOOP_CLASSPATH注释掉用方式 A 就把 lib 里的 Hadoop uber jar 移出。如果在 CDH 上同时要读 HDFS、连 Hive、提交 YARN我更倾向用方式 A因为它能跟随 CDH 的版本变化升级 CDH 时不需要同步换 Flink lib 里的 jar。5.2 YARN 容器内存申请失败的常见算式错误有次同事说 Flink 起不来日志里一直有Could not allocate container。我检查发现他把taskmanager.memory.process.size写成了4096没有带mFlink 解析成了另一种单位实际申请的内存远超 YARN 配额。还有一次是对多 slot 的算法理解偏了。要记住taskmanager.memory.process.size是单个 TaskManager 进程的总内存而 TaskManager 进程后续会被切分给多个 slot。如果并行度是 8你只给一个 TM 进程 2GB每个 slot 只有 256MB 左右作业一进来就频繁 OOM。如果真的遇到 YARN 容器上限不足优先调小 Flink 侧的内存申请而不是去改yarn.scheduler.maximum-allocation-mb。这个参数修改后需要重启 ResourceManager影响整个集群代价很大。为避免这种局面部署前就把 Flink 的内存上限定在 YARN 容器最大值的 75% 左右留出余量。5.3 想连 Hive Catalog先处理 CDH 的 Hive 版本差异CDH 6.3.2 自带的 Hive 版本是2.1.1-cdh6.3.2。Flink 1.13.1 官方提供的 Hive 集成包一般是按 Apache Hive 2.3.6 或 3.1.2 编译的直接把这些 connector 包丢到 Flink lib 里连接 CDH 的 Hive Metastore 时经常出现 protobuf、guava 或 Hive 内部接口的版本冲突。比较稳的方法是把 CDH 自己提供的hive-exec、hive-metastore、hive-common以及它们依赖的libfb303、libthrift等 jar从/opt/cloudera/parcels/CDH/lib/hive/lib/拷贝到 Flink 的 lib 目录而不是使用官方 Hive connector 自带的那套 jar。同时要确保hive-site.xml能被 Flink 读到。最简单的做法是把 hive-site.xml 放到HADOOP_CONF_DIR指定的目录下或者在flink-conf.yaml里配置 Hive Metastore 的地址。CDH 的 Hive Metastore 如果走的是 MySQL 元数据库网络连通性、账号权限也要提前验证否则 Flink 建表时会报出很底层的 JDBC 错误。6. 热词延伸CDH 6.3.2 旁边装 Doris版本适配看什么很多人搜“cdh6.3.2 使用 doris 什么版本适配”其实把问题问偏了。Doris 是一个独立的 MPP 分析型数据库它不依赖 HDFS、不依赖 Hive也不是 CDH 生态里的组件。CDH 和 Doris 之间不存在“版本适配”的问题真正要关心的是数据通路和连接器版本。6.1 Doris 并不需要“CDH 版”但需要想清楚谁做查询入口CDH 6.3.2 适合做离线数仓的数据加工和调度Doris 适合做高频查询和实时 OLAP。两者放在一起常见架构是Kafka - Flink 1.13.1 - DorisCDH 继续承担离线 ETL再通过 Doris 的 Catalog 或外部表把部分维度数据同步到 Doris 或反过来。Doris 的部署本身只看 Linux 内核版本、JDK 版本和 CPU 指令集和 CDH 是 6.3.2 还是 6.3.0 没有关系。如果你是在 CentOS 7 上离线安装提前把 Doris 的 FE/BE tar 包下载好规划好端口即可。不要被“CDH 适配版本”这类说法带偏。6.2 用 Flink 连 Doris 时Connector 版本怎么查Flink 1.13.1 连接 Doris需要使用flink-doris-connector。版本选择的唯一依据是 Doris 官方仓库中提供的 Flink 版本兼容矩阵不要凭感觉用最新版。官方为不同 Flink 主版本维护了对应的 connector 分支Flink 1.13 时期常见的是 1.2.x 这一组但请务必以你下载时官方 README 标记的对应关系为准。查证步骤很简单打开 Doris 官方 GitHub 的flink-doris-connector项目。查看 README 中的“版本兼容”或“Support Flink Version”表格。找到 Flink 1.13 对应的 connector 版本和 Scala 版本。在有网环境提前下载 connector jar 和它依赖的 doris sdk jar随 Flink lib 目录一起离线分发。如果你已经跑通了 Flink on CDH 6.3.2Doris 的接入不会增加什么复杂部署量。它只是把另一个 connector jar 放入 Flink lib然后提交一个包含DorisStreamLoad或 SQL 方言的任务而已。这套“CDH 管数据加工、Flink 做实时管道、Doris 承担查询”的组合我实际维护了大半年最深的体会是 classpath 一定要单一来源。无论是 Hadoop 也好Doris connector 也罢同一个类别只保留一套 jar宁缺毋滥。只要你把客户端的HADOOP_CLASSPATH管好把 YARN 容器内存算准flink-1.13.1 在 cdh6.3.2 上完全可以稳定运行很久。本文还有配套的精品资源点击获取
返回列表