简介:面向CDH 6.3.2平台的大数据工程师,这份资源提供了Apache Flink 1.13.1的完整离线部署组件包,解决了Flink在Cloudera企业级Hadoop生态中安装、分发与YARN调度集成的痛点。包内共5个文件,涵盖Flink核心库JAR、YARN客户端JAR、Parcel安装包、SHA校验文件及manifest元数据JSON,压缩包整体约299.52MB,文件类型清晰,适配Scala 2.11与CentOS 7环境。在CDH中可通过Parcel机制导入并激活,结合conf/flink-conf.yaml完成资源、日志与交互参数配置,再利用YARN会话或Standalone模式启动服务,实现作业提交与Web UI监控。整个部署路径兼顾企业集群的规范性与运维便捷性,尤其适合需要将实时计算能力无缝接入既有Hadoop平台的技术团队。目前已有3495人学习下载,是实践Flink on CDH部署时颇具参考价值的高质量资源。
1. flink-1.13.1 与 CDH6.3.2:为什么这对组合值得你花时间
在 CDH6.3.2 这种基于 Apache Hadoop 3.0.0 的发行版上跑 Flink 1.13.1,核心诉求通常是:不想为了用 Flink 而单独维护一套 YARN 集群,也不想升级 CDH 去赌生态兼容。CDH6.3.2 的 YARN 默认支持的是 Spark 2.4.0 和 Hive 2.1.1,它们和 Flink 1.13.1 的 Flink on YARN 模式、Hive 方言、流批一体能力在实际生产里有不少已知的兼容性缺口,但通过正确的配置和依赖裁剪,这些坑绝大部分是可以绕过去的。这篇文章会从为什么选型、怎么部署、怎么调优、到踩坑记录,完整走一遍这个组合的落地路径,适合正在做 CDH 平台之上实时计算选型,或者已经在用老版本 Flink 但想迁移到 1.13 的工程师。
2. 为什么是 Flink 1.13.1:版本定位与 CDH6.3.2 的兼容性分析
2.1 Flink 1.13.1 解决了哪些前代痛点
Flink 1.13 是一个承上启下的版本。在此之前,1.11 和 1.12 引入了 Hive 集成和 PyFlink 的显著增强,但 DataStream 和 Table API 的语义统一问题一直存在。1.13.1 在这个基础上把TableEnvironment的配置项和DataStream的ExecutionConfig做了大量对齐,特别是table.dynamic-table-options.enabled这个开关让建表时不再依赖全局配置,可以在 DDL 里直接写'scan.startup.mode' = 'earliest-offset'这类选项,DBA 和业务方各改各的,互不污染。
另外,1.13.1 的 Checkpoint 机制在处理背压时有明显改进。execution.checkpointing.checkpoints-after-tasks-finished这个参数允许在部分算子完成之后再做 checkpoint,而非必须等全图完成,这对 CDC 任务(比如从 MySQL 同步到 Kafka)非常关键,binlog 源不需要等 sink 全部写完才开始记录状态。
2.2 CDH6.3.2 的 YARN 与 Zookeeper 约束
CDH6.3.2 自带的 YARN 版本是 3.0.0,它和 Flink 1.13.1 的flink-yarn模块在ApplicationMaster的通信协议上基本兼容,因为 Flink 用的是 YARN 的 Client API 而不是内部 RPC。但有一个细节容易踩坑:CDH 的 YARN 默认开启了yarn.resourcemanager.ha.enabled(如果部署了多个 RM),而 Flink 1.13.1 的 yarn client 在读yarn-site.xml时,对yarn.resourcemanager.ha.rm-ids的解析在部分 CDH 定制配置下会失效,导致找不到 Active RM。
Zookeeper 方面,CDH6.3.2 自带的是 3.4.5 版本,Flink 1.13.1 的 HA 和 Kafka 客户端对 zk 客户端版本的兼容性良好,但注意不要在 Flink 的 lib 目录里放一个高版本的 zookeeper.jar,它会和 CDH 的 YARN 代理产生 NoSuchMethodError。
2.3 组件版本对照表
| 组件 | CDH6.3.2 自带版本 | Flink 1.13.1 要求/建议 | 兼容性说明 |
|---|---|---|---|
| Hadoop | 3.0.0 | 2.10.0+ 或 3.x | 编译时用 2.10.0 的 flink-shaded-hadoop-2-uber 即可,CDH 的 hadoop-client 在运行时优先 |
| Hive | 2.1.1 | 1.2.1 或 2.x | 用 flink-sql-hive-connector 2.6.0 版本,内部兼容 Hive 2.1 |
| Kafka | 2.2.0 (CDH 附带) | 2.4.1+ | Flink 1.13.1 的 kafka connector 需要 flink-connector-kafka_2.12-1.13.1.jar,依赖客户端 2.4.1 |
| Zookeeper | 3.4.5 | 3.4.x | 无需额外引入 zk.jar,CDH 自带的即可 |
| Scala | 2.11 (CDH Spark) | 2.12 或 2.11 | 1.13.1 默认发布 scala_2.12 版本,CDH 的 Spark 不影响 Flink 独立运行 |
2.4 选型结论:什么场景下值得用这个组合
如果你的生产环境已经是 CDH6.3.2,迁移成本最低的实时计算方案就是 Flink 1.13.1。先把 Flink 部署在 YARN 上,后续需要切换 Hive 方言或读取 Iceberg 表时,1.13.1 的 Hive 集成是稳定的。不建议直接上 Flink 1.14 或 1.15,因为它们的 Hive 方言版本要求 Hive 3.1,CDH6.3.2 需要额外维护 Hive 3 依赖,复杂度陡增。也没有必要为了用 1.13 去升级 CDH,升级 CDH 的代价远大于 Flink 升级。
3. 从零到一:在 CDH6.3.2 上部署 Flink 1.13.1 的完整步骤
3.1 准备阶段的三个前置条件
第一个前置条件:确认 CPU 架构和 JDK 版本。Flink 1.13.1 官方编译包要求 JDK 8,CDH6.3.2 的 YARN 节点默认是 JDK 8,但如果某台机器装了 JDK 11,直接把JAVA_HOME指到 JDK8 再启动。第二个前置条件:确保所有 YARN 节点(NodeManager)都能访问 Flink 的 dist 包路径,Flink on YARN 模式下 AM 会从 HDFS 拉取 dist 包到每个 NM 的本地目录,路径写错或者权限不够会一直卡在 SUBMITTED 状态。第三个前置条件:把/etc/hadoop/conf下的yarn-site.xml、core-site.xml、hdfs-site.xml软链接到 Flink 的conf目录里,别用复制,CDH 的配置是动态生成的,软链接才能保证 RM 切换后配置自动更新。
3.2 下载并组织 Flink 目录
在 CDH 集群的一台管理机上(比如部署了 YARN Client 的机器)执行:
# 下载 1.13.1 scala 2.12 版本 wget https://archive.apache.org/dist/flink/flink-1.13.1/flink-1.13.1-bin-scala_2.12.tgz tar zxvf flink-1.13.1-bin-scala_2.12.tgz mv flink-1.13.1 /opt/flink-1.13.1 # 创建 HDFS 上的 Flink 目录,并上传 dist 包(flink-yarn 模式需要) hdfs dfs -mkdir -p /tmp/flink-dist hdfs dfs -put /opt/flink-1.13.1/flink-1.13.1-bin-scala_2.12.tgz /tmp/flink-dist/ # 软链 CDH 的 Hadoop 配置 ln -s /etc/hadoop/conf/core-site.xml /opt/flink-1.13.1/conf/core-site.xml ln -s /etc/hadoop/conf/hdfs-site.xml /opt/flink-1.13.1/conf/hdfs-site.xml ln -s /etc/hadoop/conf/yarn-site.xml /opt/flink-1.13.1/conf/yarn-site.xml # 确认 scala 版本一致性 ls /opt/flink-1.13.1/lib/ | grep scala这里的scala_2.12是为了兼顾 Flink 1.13.1 里已经用 scala 2.12 编译的 Table API 和 CEP 库。如果用户的业务代码里有 Scala 2.11 编译的依赖,就需要下载flink-1.13.1-bin-scala_2.11.tgz,这个选择必须提前定,不能中途换。
HDFS 上传完成后,检查一下 dist 包的权限,确保 YARN 的yarn用户能读。CDH 默认的 HDFS 权限控制比较严格,经常出现Permission denied导致 AM 启动失败。
3.3 修改 Flink 配置:内存、并行度与资源队列
Flink 的conf/flink-conf.yaml是核心配置文件,建议按以下参数初始化:
# 内存配置,按每台 YARN 节点的物理内存减去系统预留 jobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 8192m # 并行度默认值,建议先设为 yarn 核数的 60% parallelism.default: 4 # Checkpoint 配置,务必开启 state.backend: rocksdb state.checkpoints.dir: hdfs:///tmp/flink-checkpoints execution.checkpointing.interval: 60s execution.checkpointing.timeout: 30s # 关闭全局动态表选项,避免 DDL 冲突 table.dynamic-table-options.enabled: true # Flink YARN 提交时的队列名,必须存在于 CDH YARN 中 yarn.application.queue: default其中state.backend: rocksdb的选择原因:CDH6.3.2 的 YARN 容器默认内存有限,如果状态量超过 1GB,用 HeapStateBackend 很容易 OOM,RocksDB 可以将状态溢出到磁盘,但代价是序列化和反序列化开销。生产环境建议一开始就配置 RocksDB,因为后续增量 Checkpoint 会稳定很多。
队列名default建议改成 CDH 里实际存在的队列(比如realtime),如果队列不存在,Flink 提交时会直接报Queue does not exist,这个错误信息很迷惑,经常被误判为 RM 问题。
3.4 用 YARN 模式提交一个最小 Flink 任务
在 Flink 目录下执行:
cd /opt/flink-1.13.1 ./bin/flink run \ -m yarn-cluster \ -ynm flink_job_test \ -yjm 2048 \ -ytm 8192 \ -ys 4 \ -p 8 \ ./examples/streaming/WordCount.jar参数说明:
-m yarn-cluster:提交到 YARN 模式,1.13.1 的yarn-cluster会被解析为yarn的 application 模式。-ynm:指定 YARN application 的名称,方便在 RM 界面定位任务。-yjm和-ytm:分别指定 JobManager 和 TaskManager 的内存。-ys:每个 TaskManager 的 slot 数,结合-p并行度控制容器数量。
执行完成后检查 RM 的 UI,观察 application 状态变为 RUNNING。如果一直停留在 ACCEPTED,打开 YARN 的日志看是否在拉取 dist 包,或者 AM 是否因为镜像问题启动失败。
3.5 初始化 SQL CLI 并验证 Hive 方言
Flink 1.13.1 的 SQL CLI 在 CDH 上需要额外放置 Hive 连接器,否则加载 DDL 时报 ClassNotFound:
# 在 lib 目录放置 Hive 连接器(根据你的 Hive 版本选 2.6.0 对应项) cp flink-sql-hive-connector_2.12-1.13.1.jar /opt/flink-1.13.1/lib/ # 初始化 SQL 客户端,加载 Hive 方言配置 ./bin/sql-client.sh embedded \ -init ./conf/sql-init.sqlsql-init.sql内容为:
SET execution.result-mode=tableau; SET parallelism.default=2; SET table.sql-dialect=hive;该配置让 SQL CLI 默认走 Hive 语法,比如支持CREATE TABLE ... STORED AS PARQUET。不需要 Hive 方言时,把最后一行改成default即可切换回 Flink 原生语法。
4. 连接 Kafka 与 Hive:常见连接器配置与 SQL 实战
4.1 Kafka Source 与 Sink 的配置模板
生产中最常见的链路是 Kafka -> Flink -> Hive(或 Kafka -> Flink -> Kafka)。Flink 1.13.1 的 Kafka connector 需要额外放置flink-connector-kafka_2.12-1.13.1.jar到 lib 目录,该 jar 依赖 Kafka 客户端的版本是 2.4.1,而 CDH6.3.2 的 Kafka 是 2.2.0,运行时实测兼容。
一个流式 ETL 的 SQL 建表如下:
CREATE TABLE kafka_source ( id BIGINT, name STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'ods_user_log', 'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092', 'properties.group.id' = 'flink_etl_group', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ); CREATE TABLE kafka_sink ( id BIGINT, name STRING, cnt BIGINT ) WITH ( 'connector' = 'kafka', 'topic' = 'dwd_user_cnt', 'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092', 'format' = 'json' );参数注意点:
scan.startup.mode建议先设为earliest-offset用于验证,生产可改为latest-offset。properties.group.id建议为每个任务独立设置,避免多个 Flink job 共用 group 导致 offset 混乱。WATERMARK必须和事件时间的ts保持一致,否则窗口计算会一直等待迟到数据。
4.2 Sink 到 Hive 表:核心配置与数据延迟问题
Hive Sink 在 Flink 1.13.1 中的稳定场景是 Partition 写入,和 Hive 2.1.1 配合时需要特别注意 FileSystem 的格式配置:
CREATE TABLE hive_sink_table ( id BIGINT, name STRING, dt STRING, hr STRING ) PARTITIONED BY (dt, hr) WITH ( 'connector' = 'hive', 'path' = 'hdfs:///user/hive/warehouse/dwd_user', 'sink.partition-commit.trigger' = 'partition-time', 'sink.partition-commit.delay' = '1 h', 'sink.partition-commit.policy.kind' = 'metastore,success-file', 'format' = 'parquet' );partition-time是生产推荐的提交策略,它会根据分区字段的值和当前时间比较,延迟 1 小时提交,防止迟到数据写入后又被覆盖。如果使用process-time,则每次 checkpoint 都会检查分区时间,可能出现小文件碎片化。
4.3 SQL 提交前的检查清单
在运行上述 SQL 之前,务必检查:
- Flink lib 目录下有没有
flink-sql-hive-connector,如果没有,在 SQL CLI 里执行SHOW TABLES都会报错。 - Hive Metastore 的地址是否正确,CDH 的 hive-site.xml 里
hive.metastore.uris默认是 Thrift 地址,如果不通,SQL 建表会卡在初始化阶段。 hive-site.xml需要被 Flink 访问到。将 Hive 的配置文件软链到 Flink conf 下是最简单的方式。
5. CDH6.3.2 上 Flink 运行常见问题:踩坑记录与排查指南
5.1 提交任务后一直处于 ACCEPTED 状态
现象:执行flink run -m yarn-cluster后,YARN 的 RM UI 上 application 一直显示 ACCEPTED,既不调度也不失败。
原因:Flink 的 dist 包没有上传到 HDFS,或者上传到的目录不在 YARN Client 的搜索路径中。CDH 的 YARN 对本地目录有白名单限制,Flink 默认把分发包放到临时目录,但在 CDH 上临时目录经常被清理或不可写。
解决:将 dist 包手工放到/tmp/flink-dist目录并修改flink-conf.yaml中的yarn.application.connector.dist.path指向该路径。另外,确认 YARN 节点的yarn.nodemanager.local-dirs有足够的磁盘空间,Flink 会在 NM 本地目录缓存 dist 包,空间不足时也会卡住。
5.2 RocksDB 状态后端导致 JVM 频繁 Full GC
现象:任务运行数小时后,TaskManager 的 GC 时间占比超过 30%,甚至出现OutOfMemoryError: Direct buffer memory。
原因:RocksDB 的默认配置会调用系统的内存分配,且不受 Flink 的 TaskManager 堆内存管理约束。在 CDH6.3.2 上,NM 容器的内存限额(CGroup)会强杀超过限制的进程。
解决:在flink-conf.yaml中设置 RocksDB 的 managed memory:
state.backend.rocksdb.memory.managed: true state.backend.rocksdb.block.cache-size: 128mb state.backend.rocksdb.writebuffer.size: 64mb同时调大 JVM OverHead 比例:taskmanager.memory.jvm-overhead.fraction: 0.2。这两个参数配合后,RocksDB 的内存会被 Flink 纳入统一管理,不会突破容器限额。
5.3 写入 Hive 表数据不落盘
现象:SQL 任务正常执行,Kafka 源有数据流入,但 Hive 表查询始终为空。
原因:这是 1.13.1 的 Hive connector 的一个已知行为,如果 Flink 任务不是以BATCH模式运行,Hive Sink 默认不会主动提交分区,导致数据停留在 staging 目录。
解决:设置流式写入 Hive 的提交策略,前面示例中的sink.partition-commit.trigger必须设置为partition-time,另外检查sink.partition-commit.delay是否比批处理窗口长,如果写入频率低于延迟值,分区可能永远不提交。一个快速验证方式是临时把 delay 设为1 s,任务运行一分钟后查询 Hive 表,能查到数据再调整回生产值。
5.4 JDBC 连接器报ClassNotFoundException
现象:使用 JDBC Sink(例如写入 MySQL)时,报错找不到org.apache.flink.connector.jdbc.JdbcSinkFunction。
原因:JDBC 连接器在 1.13.1 中不是 lib 的默认组件,需要单独下载flink-connector-jdbc_2.12-1.13.1.jar放入 lib 目录。CDH 环境通常没有外网权限,需要提前下载后分发到所有 TaskManager 节点。
解决:离线部署时,把整个 Flink 的 lib 目录打包,在flink-conf.yaml里通过yarn.application.connector.dist.path指定打包好的 dist 包地址,避免每台 NM 手工放 jar。注意如果有多个 Flink 工程,依赖的 Jar 尽量统一放在 lib 下,避免冲突。
5.5 并行度过高导致 HDFS 小文件爆炸
现象:写入 Hive 的分区文件数量等于并行度,一个分区下出现几十个几十 KB 的小文件,Hive 查询性能急剧下降。
原因:Flink 的 Hive sink 默认按并行度写入,没有做文件合并。1.13.1 中auto-compaction功能不完整。
解决:先将 Hive Sink 并行度控制在 2 或 3,或者写入后额外跑一个 Hive SQL 做INSERT OVERWRITE ... SELECT合并文件。如果在 1.13.1 里使用sink.partition-commit.policy.kind = metastore,success-file,并配合主键或时间字段做分桶,也可以缓解。
6. 进阶:用 SQL Client 做流批一体任务与验证方法
6.1 用同一个 SQL 跑批和流
Flink 1.13.1 的 Table API 实现了 SQL 的批流统一,在 CDH 上可以通过SET execution.runtime-mode = batch或streaming切换同一个 DDL 的执行模式。对于 Hive 表的读取,批模式下走 MapReduce 输入,流模式下走文件监听。
一个实用的验证方式:用同一个 Kafka -> AGG -> Hive 的 SQL,在streaming模式下观察窗口聚合结果,确认无误后,将execution.runtime-mode切到batch,对同一天的历史数据重跑,对比聚合结果是否一致。如果不一致,通常是 Watermark 或状态 TTL 配置导致的问题,在流模式下检查table.exec.state.ttl是否太短。
6.2 自定义 Data Source 与 Sink 的注册方式
如果业务需要对接非标准数据源(比如自研 MQTT),在 1.13.1 中可以用TableSourceFactory和TableSinkFactory实现,但更快的验证方式是直接写一个 DataStream 的 UserFunction,再通过StreamTableEnvironment转成 Table 注册。
Java 代码关键部分:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tEnv = StreamTableEnvironment.create(env); // 自定义 Source,实现 SourceFunction DataStream<Row> stream = env.addSource(new MySourceFunction(), TypeInformation.of(Row.class)); // 将 DataStream 注册为 Table,供 SQL 查询 tEnv.createTemporaryView("user_log", stream, $("id"), $("name"), $("ts").rowtime()); // 之后即可用 SQL 查询 Table result = tEnv.sqlQuery("SELECT id, COUNT(*) FROM user_log GROUP BY id");注册自定义 Sink 时,推荐继承RichSinkFunction,在open()方法里创建连接,在invoke()里做反序列化和写入。务必在close()里释放连接资源,否则 TaskManager 重启后会出现连接泄漏。
6.3 通过 Flame Graph 定位性能瓶颈
生产环境出现反压时,不要只盯 YARN 的 CPU 监控。Flink 自带火焰图工具,在 Web UI 的 Job 详情页任务节点上右键执行火焰图采样。能明确是序列化热点还是算子内逻辑热点。
比如在 CDH 上常见的一个性能问题:Kafka Source 反序列化慢,导致整个作业吞吐上不去。开启火焰图后,如果JSONDeserializationSchema的 CPU 占比超过 60%,把format换成avro或者protobuf,吞吐一般能提升一倍以上。
6.4 一套自检清单,上线前逐条验证
我在每次上线新 Flink 任务前,会按固定顺序做以下验证,这套流程能挡掉 90% 的生产事故:
- 检查 YARN 队列权限,确认任务不会跑到别的组队列里,CDH 的队列配额错误是静默的,任务能跑但资源被限。
- 用
./bin/flink list -m yarn-cluster查看当前所有任务的状态,确认没有同组任务互相争抢 slot。 - 开启 checkpoint 后重启任务,确认从 checkpoint 恢复的时间在可接受范围内,如果 RocksDB 的恢复时间过长,适当调大
state.backend.rocksdb.thread.num。 - 模拟一次 Kafka 集群抖动(比如停掉一台 broker),确认 Kafka Source 的
properties.enable.auto.commit设为 false 且手动 offset 提交正常,否则重启任务可能重复消费。 - 最后一条也是血泪教训:不要把 Flink 的 lib 目录下不需要的 jar 删掉,尤其是
flink-table-planner和flink-table-runtime,1.13.1 默认两个都需要,删了 SQL Client 会直接启动失败。
CDH6.3.2 上部署 Flink 1.13.1 这条路,我自己走下来的体感是:「版本匹配」这件事实在太关键。网上很多教程默认 Flink 的 dist 包自带 Hadoop 客户端,但 CDH 的 YARN 客户端不认那一套。只要你把 Hadoop 配置软链做好、Hive connector 选对版本、RocksDB 内存交给 Flink 托管,大部分任务都能稳定跑起来。希望这些经验能帮你少走几趟弯路,祝顺利。
本文还有配套的精品资源,点击获取