尧图精选

Flink 1.13.1 与 CDH6.3.2 集成部署实战与踩坑指南

🕒 发布时间:2026/10/1 5:22:53 📁 来源:尧图网络
简介面向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 要求/建议兼容性说明Hadoop3.0.02.10.0 或 3.x编译时用 2.10.0 的 flink-shaded-hadoop-2-uber 即可CDH 的 hadoop-client 在运行时优先Hive2.1.11.2.1 或 2.x用 flink-sql-hive-connector 2.6.0 版本内部兼容 Hive 2.1Kafka2.2.0 (CDH 附带)2.4.1Flink 1.13.1 的 kafka connector 需要 flink-connector-kafka_2.12-1.13.1.jar依赖客户端 2.4.1Zookeeper3.4.53.4.x无需额外引入 zk.jarCDH 自带的即可Scala2.11 (CDH Spark)2.12 或 2.111.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.1CDH6.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 8CDH6.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 很容易 OOMRocksDB 可以将状态溢出到磁盘但代价是序列化和反序列化开销。生产环境建议一开始就配置 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-modetableau; SET parallelism.default2; SET table.sql-dialecthive;该配置让 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 memorystate.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 DataStreamRow 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-runtime1.13.1 默认两个都需要删了 SQL Client 会直接启动失败。CDH6.3.2 上部署 Flink 1.13.1 这条路我自己走下来的体感是「版本匹配」这件事实在太关键。网上很多教程默认 Flink 的 dist 包自带 Hadoop 客户端但 CDH 的 YARN 客户端不认那一套。只要你把 Hadoop 配置软链做好、Hive connector 选对版本、RocksDB 内存交给 Flink 托管大部分任务都能稳定跑起来。希望这些经验能帮你少走几趟弯路祝顺利。本文还有配套的精品资源点击获取
上一篇/下一篇内容由系统自动关联 返回资讯列表 →