简介Apache Flink 1.13.1强化了实时流处理能力Cloudera CDH 6.3.2则是企业级Hadoop生态集成平台这份资源面向大数据平台工程师、运维人员与实时计算开发者解决Flink与CDH之间的Parcel安装、YARN资源调度和作业运行管理等集成问题。压缩包共5个文件整体约299.52MB文件间职责清晰且互为补充其中Parcel格式的二进制分发包适配Scala 2.11与CentOS 7环境JAR文件涵盖Flink核心库和YARN客户端JSON提供元数据信息SHA校验文件保证完整性导入CDH后即可在集群内统一分发与版本校验。已有3495人学习下载适合希望通过Parcel机制快速完成部署、避免自行编译整合的CDH用户尤其适合已有YARN基础、首次接触Parcel部署的团队。借助包内文件可依次完成Parcel激活、flink-conf.yaml资源参数调整、内存与TaskManager槽位规划并通过YARN客户端提交实时作业后续结合Flink Web UI或CDH管理界面进行状态监控与故障排查明显缩短部署周期提升大数据实时任务的交付效率。1. 给 CDH 6.3.2 单独部署一套 Flink 1.13.1装上去只是开始如果你所在的团队还在用 CDH 6.3.2又想把实时计算落到 Flink 上大概率会卡在一个尴尬位置CDH 的组件体系里没有一套能直接跑的社区版 Flink控制台自带的版本要么太老要么跟 Hive、Hadoop 的私有依赖互相打架。我在生产环境里折腾过两轮之后固定下来的组合就是 Flink 1.13.1 搭配 CDH 6.3.2。这个组合不算新但胜在稳定Flink 1.13.1 的 Hive 集成已经成熟YARN 会话模式在 CDH 的调度器下跑得很顺而且网上踩坑记录够多出了问题搜得到方向。这篇文章不讲空泛的架构直接说清楚从安装配置到部署要动哪些文件、改哪些参数、提交作业后怎么验证以及那些让你半夜爬起来看日志的坑。2. 版本选型与依赖准备为什么是 1.13.1以及 lib 目录里该放什么2.1 CDH 6.3.2 的 Hadoop 和 Hive 版本决定了你没法直接抄官方文档CDH 6.3.2 这一版的核心组件版本是 Hadoop 3.0.0-cdh6.3.2、Hive 2.1.1-cdh6.3.2Java 环境是 JDK 8。这套版本组合跟 Apache 原生发行版的差异在于Cloudera 在 Hadoop 和 Hive 的源码里打了不少自己的补丁包名、类路径、依赖版本都动过。官方 Flink 1.13.1 文档里默认假设你连的是 Apache Hadoop 2.x 或 3.x 的标准发行版直接按文档配第一步就可能挂在 classpath 上。我在第一次部署时犯过一个典型错误直接下载了 flink-1.13.1-bin-scala_2.12.tgz解压后把flink-conf.yaml改了改就提交到 YARN结果作业一直报ClassNotFoundException: org.apache.hadoop.fs.FileSystem。原因很简单Flink 发行包从 1.11 开始就不再打包 Hadoop 相关依赖需要自己往lib目录里放一个兼容的 shaded Hadoop jar。这个 jar 放不放、放哪个版本直接决定了你能不能跑起来。还有一个容易被忽略的点CDH 6.3.2 的 YARN 用的是 Fair Scheduler 还是 Capacity Scheduler会影响你提交 Flink 作业时的队列参数写法。我这边集群用的是 Capacity Scheduler提交时必须显式指定队列名否则作业会被丢到 default 队列而 default 队列往往没有分配资源作业就一直卡在 ACCEPTED 状态。这不是 Flink 的问题是 CDH 的调度策略但你得知道它会影响yarn-session.sh的-qu参数。2.2 两个关键 jarflink-shaded-hadoop 和 hive-exec放错目录就是黑匣子依赖准备的第一个动作是确认$FLINK_HOME/lib目录下有什么。干净解压出来的 Flink 1.13.1lib目录里只有核心引擎和几个自带连接器绝对没有 Hadoop shaded jar。我一般会做这样一件事FLINK_HOME/opt/flink-1.13.1 cd $FLINK_HOME/lib # 先看缺什么 ls -lh *.jar | awk {print $5, $9} # 从 Maven 仓库拉 Flink 官方打包的 Hadoop 3.2.0 shaded jar wget https://repo1.maven.org/maven2/org/apache/flink/flink-shaded-hadoop-3-uber/3.2.0-10.0/flink-shaded-hadoop-3-uber-3.2.0-10.0.jar这个 jar 是 Flink 社区专门用来桥接 Hadoop 生态的 uber 包里面把 Hadoop 的hadoop-common、hadoop-hdfs、hadoop-mapreduce-client-core等核心依赖全部重定位过了不会跟 CDH 自带的其他组件冲突。为什么选 Hadoop 3.2.0 这个 shade 版本而不是 2.10.1因为 CDH 6.3.2 本身是 Hadoop 3.0.0虽然主版本一致但 RPC 协议细节有差异我用 3.2.0 的 shaded jar 跑了读写 HDFS、提交 YARN、读 Hive 表这些操作都没有出现协议不兼容的问题。2.10.1 那个是给 Hadoop 2.x 集群准备的在 CDH 6.3.2 上用的话偶尔会在提交 YARN 时出现 token 相关的异常。第二个关键 jar 是 Hive 的执行包。Flink 1.13.1 对 Hive 的版本支持是 1.2.2、2.1.1、2.3.6、3.1.0 这四类CDH 6.3.2 的 Hive 版本是 2.1.1-cdh6.3.2跟社区版 2.1.1 并不是同一个 jar。正确做法是把 CDH 自己打包的hive-exec拷贝到 Flink 的 lib 目录# 确认 CDH 的 Hive 版本 ls /opt/cloudera/parcels/CDH/lib/hive/lib/hive-exec-*.jar # 拷贝到 Flink 的 lib 目录避免权限问题 cp /opt/cloudera/parcels/CDH/lib/hive/lib/hive-exec-2.1.1-cdh6.3.2.jar /opt/flink-1.13.1/lib/这里有个很隐蔽的坑hive-exec包里面带着javax.jdo、org.datanucleus等一堆依赖这些类跟 Flink 自带的某些 jar 里的类会在运行时打架。我的经验是只拷贝hive-exec不要图省事把整个hive/lib目录全部丢进 Flink 的 lib不然作业启动时出现NoSuchMethodError的概率极高而且错误信息指向的类往往跟实际冲突的 jar 没有直观关联排查起来非常痛苦。做完这两步之后再检查一遍lib目录里有没有flink-sql-connector-hive-2.1.1_2.12-1.13.1.jar这个文件。如果后面要用 Flink SQL 写 Hive 表这个 jar 是必需的。我没有放在这个阶段安装而是把hive-site.xml的路径提前记下来后面配 Hive Catalog 的时候会用到。依赖这层准备完毕接下来就是配置文件。3. 从安装配置到部署YARN 会话模式跑通第一个作业3.1 flink-conf.yaml 里必须改的参数Flink 1.13.1 的配置文件集中在$FLINK_HOME/conf/flink-conf.yaml。默认配置能跑本地模式但要部署到 CDH 的 YARN 上至少要动下面这些参数。我直接贴一份经过生产验证的配置片段jobmanager.memory.process.size: 4096m taskmanager.memory.process.size: 8192m taskmanager.memory.managed.fraction: 0.4 taskmanager.numberOfTaskSlots: 4 parallelism.default: 4 high-availability: zookeeper high-availability.zookeeper.quorum: cdh01:2181,cdh02:2181,cdh03:2181 high-availability.storageDir: hdfs:///flink/ha rest.bind-address: 0.0.0.0 rest.bind-port: 8081 yarn.application-attempts: 4 yarn.properties.xxx: # 占位实际不用配逐项说明一下。jobmanager.memory.process.size和taskmanager.memory.process.size是 Flink 1.10 之后引入的新版内存模型直接写进程总内存比旧版的jobmanager.heap.size更直观。按照你手头每台 NodeManager 的物理内存来定我这边的机器是 64GBTM 给 8GB一个节点最多起 7 个 TM。taskmanager.memory.managed.fraction默认是 0.4这个值是给排序、Join 用的托管内存如果作业里大量使用 RocksDB 状态后端我会把比例降到 0.3给 RocksDB 腾空间。高可用配置指向 Zookeeper这是因为 CDH 6.3.2 自带 ZK 服务直接用现成的。high-availability.storageDir必须放在 HDFS 上JobManager 的元数据会写到这个目录。yarn.application-attempts建议设成 4如果设太大JobManager 反复异常重启会导致整个 YARN 应用被集群标记为失败。还有一类参数是 web 端口rest.bind-address设为0.0.0.0是为了方便本机之外的机器打开 Web UI如果你有安全组管理可以收紧。这里我要单独提醒这句话不要直接照抄我给的端口和内存值CDH 集群上每个队列的资源配额不同你的容量规划需要按队列实际配置来调。我见过有人把taskmanager.memory.process.size配成 32GB结果 Container 请求超过队列上限作业永远起不来。3.2 启动 YARN 会话与提交作业的完整命令配置改完后启动会话模式。我习惯用 detached 模式启动这样终端关掉不影响会话进程cd /opt/flink-1.13.1 ./bin/yarn-session.sh \ -nm flink-yarn-session \ -qu root.flink \ -jm 4096m \ -tm 8192m \ -s 4 \ -d参数含义-nm是 YARN 上显示的应用名便于在 ResourceManager 的页面上快速找到-qu指定队列名这个队列名必须跟 CDH 集群里实际存在的队列对上否则提交直接报错-jm和-tm分别指定 JobManager 和 TaskManager 的进程内存会覆盖flink-conf.yaml里的对应配置-s指定每个 TaskManager 的 slot 数我一般设为 4这样 8GB 内存分配 4 个 slot每个 slot 差不多 2GB兼顾内存隔离和并行度弹性-d是 detach 模式。启动成功后屏幕上会打印类似JobManager Web Interface: http://cdh01:8081的信息同时 YARN ResourceManager 页面上也会出现这个应用。注意如果启动后没有打印 Web 地址而是卡在Deploying cluster大概率是队列资源不够或者 HDFS 写权限有问题去 ResourceManager 日志里搜Application master的报错。接下来提交一个测试作业。Flink 发行包自带 examples我习惯先用它验证链路./bin/flink run \ -m yarn-session \ -d \ -p 4 \ ./examples/batch/WordCount.jar \ --input hdfs:///tmp/wordcount_input.txt \ --output hdfs:///tmp/wordcount_output-m yarn-session表示附加到已有的会话-p 4覆盖默认并行度。如果作业正常跑完HDFS 上能读到输出文件说明 Flink 已经能在 CDH 的 YARN 上正常调度容器。我在生产里用这个方式跑的第一个作业就是做一个简单的 Kafka 到 HDFS 的渠道验证网络和 HDFS 权限。检查作业状态也有固定套路./bin/flink list -m yarn-session # 输出里能看到 RUNNING 状态的作业和对应的 JobIDflink list是最直接的验证方式。作业运行中如果觉得 Web UI 访问不了检查rest.bind-address和防火墙端口跟 Flink 本身没关系。3.3 资源配置与队列隔离的取舍YARN 会话模式适合场景你有多个小作业轮流跑不想每次提交都重新拉起一个 Flink 集群。但会话模式有一个隐患多个作业共享同一个 JobManager如果某一个作业把 TM 内存吃满其他作业会被拖死。我在 CDH 上经历过的翻车现场两个 Flink SQL 作业跑在同一个会话里其中一个作业的 RocksDB 状态快速增长TM 的 GC 时间飙升另一个作业的延迟从秒级变成分钟级。排查半天发现是两个作业在抢同一个 TM 的堆内存。从那以后凡是状态量大、吞吐高的作业我都会用 Application 模式单独提交。Flink 1.13.1 支持flink run -t yarn-application每个作业一个独立的 YARN 应用资源隔离更彻底代价是作业启动时间从秒级变成分钟级因为要重新拉起 JobManager 和 TaskManager。./bin/flink run \ -t yarn-application \ -Djobmanager.memory.process.size2048m \ -Dtaskmanager.memory.process.size4096m \ -Dtaskmanager.numberOfTaskSlots2 \ -Dparallelism.default2 \ -Dyarn.application.queueroot.flink \ ./my-job.jar-D参数可以逐个覆盖配置不用改 flink-conf.yaml。如果你对作业的资源需求比较确定直接走yarn-application模式少踩会话模式互相干扰的坑。4. Flink SQL 集成 HiveCatalog 配置与 sink 不落表的排查4.1 Hive Catalog 的 jar 与 hive-site.xml 对齐Flink 1.13.1 的 Hive 集成是HiveCatalog它能让你在 Flink SQL 里直接读写已有的 Hive 表。前提有三个hive-execjar 在 lib 目录、Hive 版本匹配、能读到hive-site.xml。第四章开头我已经把hive-exec放到 lib 了这里主要说 Catalog 配置。启动 SQL 客户端cd /opt/flink-1.13.1 ./bin/sql-client.sh embedded在客户端里执行CREATE CATALOG cdh_hive WITH ( type hive, default-database default, hive-conf-dir /etc/hive/conf ); USE CATALOG cdh_hive; SHOW TABLES;hive-conf-dir必须指向 CDH 上 Hive 的配置目录那里有hive-site.xml里面写着 Hive Metastore 的地址。如果SHOW TABLES报错先看是不是metastore uris连不上CDH 里 Hive Metastore 服务一般在端口 9083。如果报错是ClassNotFound: org.apache.hadoop.hive.conf.HiveConf说明hive-exec没放对地方或者 SQL 客户端启动的不是你放 jar 的那个 Flink 目录。还有一类常见报错是Unsupported class file major version这说明 SQL 客户端用了 JDK 11 以上启动而 CDH 的 Hive 相关 jar 是按 JDK 8 编译的。我的建议是统一用 JDK 8 跑 Flink设好JAVA_HOME再启动客户端别在这种基础问题上浪费精力。4.2 写入 Hive 表checkpoint 与提交策略Flink SQL 写 Hive 表本质是流式写入但 Hive 表本身不是流式存储。Flink 的处理方式是先把数据写到 HDFS 上的临时文件等 checkpoint 完成后再做提交把临时文件变成正式分区文件。如果你没开 checkpoint数据就会一直卡在临时文件里看起来像“没写进去”。一个标准的写入流程是这样SET execution.checkpointing.interval 60s; SET execution.checkpointing.mode EXACTLY_ONCE; SET parallelism.default 4; USE CATALOG cdh_hive; INSERT INTO dwd_orders SELECT order_id, user_id, amount, order_time FROM kafka_orders;execution.checkpointing.interval设成 60 秒意思是 Flink 每 60 秒做一次 checkpoint每次 checkpoint 完成时Hive sink 才会提交一批文件。如果你把这个间隔改小比如 10 秒提交频率会升高HDFS 上会产生大量小文件。改大也不行比如 10 分钟那数据延迟就有 10 分钟。生产上我一般根据业务容忍度取 30 到 60 秒。EXACTLY_ONCE模式对 Hive sink 来说意味着提交阶段会做原子性重命名能避免下游读到半截文件。这个参数在数据量小的时候感知不强但一旦作业失败重启你就会发现它有多重要——没有它重启后会出现重复数据和缺失数据并存的情况。4.3 排查步骤一个 sink 不落表现场“Flink sink hive 表数据不入表”是我被问得最多的问题几乎每个 CDH 生产环境都遇到过。我总结出一个固定排查顺序按顺序查基本能在十分钟内定位。第一步看 HDFS 上目标分区的临时文件是否存在。查法如下hdfs dfs -ls /user/hive/warehouse/dwd_orders/ # 发现大量 part-xxx-xxx.inprogress 文件inprogress文件是尚未提交的临时文件。如果文件存在但一直是.inprogress说明 checkpoint 没有完成。去 Flink Web UI 看 Checkpoints 页面如果显示FAILED通常原因是状态后端写入 HDFS 超时或者作业重启太频繁导致 checkpoint 无法对齐。第二步看作业日志里有没有提交阶段的报错。我遇到过的错误长这样Caused by: java.io.IOException: Failed to commit file ... .part-xxx这个错误在 CDH 上最常见的根因是权限运行 Flink 的 YARN 用户对 Hive 的 warehouse 目录没有写权限。CDH 的 Hive 表存储目录通常是/user/hive/warehouse属主是 hive 用户而你提交 Flink 作业的用户可能是 flink 或其他 Linux 账号两者权限不一致提交阶段就会炸。解决方法是给对应的 HDFS 目录加 ACL或者让 Flink 作业跑在跟 Hive 同一个用户组下hdfs dfs -chmod -R 775 /user/hive/warehouse/dwd_orders第三步检查是不是分区表但没写分区字段。Hive 分区表要求 INSERT 语句里带上所有分区列漏掉一个Flink 会把数据写到错误路径看起来也是不入表。第四步确认 SQL 客户端会话里USE CATALOG之后操作的表确实属于当前库别在默认库和cdh_hive库之间搞混了。这套顺序我基本每次都有效。再有想不起来时就先开一个 60 秒的 checkpoint再随便写一个查询往 Hive 表丢数据看文件从.inprogress变成正式文件链路就通了。5. 常见问题与避坑记录JDBC 连接器、CDC 与跨集群迁移5.1 JDBC 连接器异常ClassNotFound 与驱动覆盖Flink 1.13.1 的 JDBC 连接器是生产环境最常见的连接器之一但它在 CDH 上有一个高频翻车点驱动类找不到。现象Caused by: java.lang.ClassNotFoundException: com.mysql.cj.jdbc.Driver原因很直白Flink 发行包里不会带 MySQL 驱动的第三方 jar你需要自己放。但这里有个陷阱很多同事只放了mysql-connector-java-5.1.x.jar这个老版本驱动对应的是com.mysql.jdbc.Driver类而你的 JDBC URL 里写了com.mysql.cj.jdbc.Driver自然找不到。解决方式wget https://repo1.maven.org/maven2/mysql/mysql-connector-java/8.0.28/mysql-connector-java-8.0.28.jar mv mysql-connector-java-8.0.28.jar /opt/flink-1.13.1/lib/放完 jar 重启 Flink 会话再提交作业前先确认驱动类名和 URL 一致。Flink SQL JDBC 表声明可以这样写CREATE TABLE mysql_dim ( id INT, name STRING ) WITH ( connector jdbc, url jdbc:mysql://10.0.0.5:3306/dim?serverTimezoneAsia/Shanghai, table-name user_dim, driver com.mysql.cj.jdbc.Driver, username flink, password flink123 );serverTimezoneAsia/Shanghai是 8.0 驱动的强制要求不写会报时间时区错误。驱动类名选 8.0 的com.mysql.cj.jdbc.Driver不要图兼容写旧类名否则会出现executeQuery() can not be used for update这类误导性报错。5.2 JDBC 连接器异常超时与连接池参数另一个高频问题是作业运行十几分钟后开始抛Communications link failure任务自动重启重启之后又正常过一阵子再断。现象Caused by: com.mysql.jdbc.exceptions.jdbc4.CommunicationsException: Communications link failure我遇到过两个根因。第一个是 MySQL 的wait_timeout默认是 8 小时但 Flink JDBC sink 的连接池里有空闲连接超过这个时间没被使用MySQL 服务端把连接断掉了。Flink JDBC 连接器提供了两个参数WITH ( sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 2s )第二个根因是max-rows太小导致频繁 flush连接反复创建销毁。如果你的数据量不大但吞吐有峰值可以适当调大buffer-flush.max-rows到 5000interval调到 5 秒减少连接建立次数。记住一点这两个参数是防御性配置目的是控制批量提交节奏不是越大越好。5.3 CDC pipeline 部署server-id 与 binlog 权限Flink CDC 在 CDH 生产环境上的部署坑比 JDBC 连接器更隐蔽。先声明一个依赖问题Flink 1.13.1 要配合 Flink CDC 2.1.x 使用需要单独下载flink-sql-connector-mysql-cdc-2.1.1.jar放到 lib 目录否则 SQL 客户端不认识mysql-cdc连接器。现象ERROR: binlog connection is not available, please check binlog config原因基本是 MySQL 的 binlog 没开或者权限不对。CDC 需要 MySQL 用户具备RELOAD、SELECT、REPLICATION SLAVE、REPLICATION CLIENT四个权限。另外一个高频坑是server-id冲突我先贴一个能跑的声明CREATE TABLE cdc_orders ( id INT PRIMARY KEY, order_no STRING, amount DECIMAL(10,2) ) WITH ( connector mysql-cdc, hostname 10.0.0.10, port 3306, username flink_cdc, password xxx, database-name trade, table-name orders, server-id 5400-5404, scan.incremental.snapshot.chunk.key-column id );server-id是 CDC 连接器伪装成 MySQL 从节点时用的 ID如果跟线上已有的从节点重合MySQL 主库会直接断开连接现象就是作业启动后几秒内失败报错里带着slave with the same server_uuid这类关键字。解决方式是把server-id设成一个范围比如5400-5404连接器会在范围内取一个没被占用的值。chunk.key-column是增量快照阶段分片用的指定主键列即可。还有一个部署注意点CDC pipeline 如果目标端是 Hive那 sink 阶段同样依赖 checkpoint你需要把 CDC 的作业和 Hive sink 的 checkpoint 间隔配合好我这里一般设 30 秒。5.4 跨 CDH 集群迁移先测带宽再定并行度把数据从一个 CDH 集群迁移到另一个 CDH 集群用 Flink 读源集群的 Kafka 或 Hive写到目标集群的 Hive这个场景现在很常见。但很多人上来就调并行度忽略了两个集群之间的带宽限制。现象作业并行度从 4 提到 8吞吐几乎没有变化但背压指标一路飘红。原因大概率是跨集群网络带宽成了瓶颈。我在做迁移之前会强制走一个流程先用iperf3测两个集群网关之间的实际带宽# 在目标集群的一台机器上启动服务端 iperf3 -s -p 5201 # 在源集群的一台机器上启动客户端测 60 秒 iperf3 -c 10.0.0.7 -p 5201 -t 60如果测出来的带宽只有 200Mbps而你要迁移的数据每天有几十 GB那 Flink 并行度再高也没用。这时候正确做法是控制作业的写入速率或者引入 Kafka 做中间缓冲。还有一种比测带宽更直接的验证方式在目标集群上用hdfs dfs -put写一个同样大小的测试文件记录耗时换算实际吞吐再决定 Flink 作业的并行度这样不会一上来就被反压打崩溃。6. 投产前最后一步火焰图、savepoint 与恢复演练6.1 打开火焰图接口看热点Flink 1.13.1 内置了实验性的火焰图功能在conf/flink-conf.yaml里加一行rest.flamegraph.enabled: true打开后Web UI 进入某个作业的页面点某个 subtask就能看到 FlameGraph 标签页。火焰图对定位 CPU 热点非常直观我曾经用它发现一个 SQL 作业里 UDF 做了大量字符串split占 CPU 超过 60%改成自定义解析之后吞吐直接翻倍。6.2 用 savepoint 做一次容灾演练投产前必须验证 savepoint 机制不然你都不知道作业到底能不能恢复。操作分两步先给运行中的作业做 savepoint# 从 flink list 拿到作业 ID ./bin/flink savepoint jobid hdfs:///flink/savepoints -yid application_xxx提交成功后从返回信息里抄下 savepoint 路径。然后杀作业再从这个路径恢复./bin/flink run \ -t yarn-application \ -s hdfs:///flink/savepoints/savepoint-xxxx \ ./my-job.jar恢复成功的标志是新作业的状态能接上旧作业的进度比如 Kafka 消费位点从 savepoint 时刻继续而不是从头消费。这个演练我每次上线前都强制做特别是改完状态后端或者并行度之后不验证就直接上线纯属赌运气。从那以后我每次在 CDH 6.3.2 上部署或升级 Flink 作业都强制走一遍这个流程检查 lib 目录、确认 flink-conf.yaml 的资源参数、提交后立刻验证 checkpoint 是否成功、最后做一次 savepoint 恢复。这套流程帮我挡住了至少三次上线事故。希望帮到你。本文还有配套的精品资源点击获取