ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink+Hudi实时数据湖入湖实践:从环境搭建到性能调优

Flink+Hudi实时数据湖入湖实践:从环境搭建到性能调优 前阵子在做一个实时数据湖的改造Kafka 里的业务变更数据要实时落到 Hudi 表下游再用 Presto 和 Spark 做分析。选型的时候没怎么犹豫就定了 Hudi Flink。用到今天这套组合在线上已经稳定跑了快半年中间踩过不少坑也总结出一些可以直接抄作业的配置。这篇把从环境搭建到性能调优的完整过程写出来给正在搞湖仓一体的朋友做个参考。文章主要面向两类人一类是刚开始接触 Flink 和 Hudi想快速搭一个 demo 看到数据跑通的新手另一类是已经上了 Hudi但被各种参数、异常、性能问题缠身想找排查思路的工程师。1. 项目概述与方案选型1.1 为什么选择 Hudi Flink先说结论如果你的核心诉求是“把变更流实时写入数据湖同时能支持 upsert 和增量查询”Hudi Flink 几乎是当前最顺的组合。Hudi 本质上是一个数据湖存储格式它在 HDFS 或者对象存储上维护表结构、数据文件、索引和元数据。Flink 则负责流式计算和状态管理。二者结合以后Flink 利用 checkpoint 机制保证写入的精确一次语义Hudi 则通过内部的 file group、file slice 和索引机制把 upsert 语义落地到文件层面。跟纯写 Parquet 文件不同Hudi 会把每次写入划分为一个个 commit并且维护 commit 的元数据这样下游既能读取快照数据也能消费增量数据。很多人会问为什么不直接用 Iceberg 或者 Delta Lake。我的理解是Iceberg 在 Flink 集成上做得也不错但在 upsert 场景下的成熟度不如 HudiDelta Lake 官方对 Flink 的支持相对滞后主要还是围绕 Spark 生态。而 Hudi 对 Flink 的原生支持已经比较完整从 Flink SQL 建表、流式写入、流式读取到归档和清理都能通过 SQL 或 API 触发这是生产环境非常看重的。1.2 适用场景与整体架构Hudi Flink 最常见的落地场景是实时数仓的 ODS 层和 DWD 层。比如业务库通过 CDC 同步到 KafkaFlink 消费 Kafka 后直接写入 Hudi ODS 表保留所有历史变更下游再用 Spark 或者 Presto 做批量加工或者用 Flink 读取 Hudi 增量数据做实时关联。这套架构的好处是Kafka 的数据不会被重复消费时丢失因为 checkpoint 会记录消费位点Hudi 写入本身也是事务性的同时 Hudi 表暴露给分析引擎的是一个统一的表视图不需要额外维护一份全量数据。整体数据流大概是这样的业务数据库 - Canal/Debezium - Kafka - Flink 作业 - Hudi(DFS/OSS) - Spark/Presto/Hive在这个流程里Flink 作业承担了最核心的 ETL 逻辑包括数据清洗、字段映射、维表关联以及最后的写入。Hudi 表则同时支撑三种读模式快照读查最新状态、增量读消费某个 commit 之后的变更、读优化模式只读列式文件跳过 log 文件。这些能力在做实时数仓分层时非常有用。2. 环境准备与版本选型2.1 版本兼容矩阵与实测组合Hudi 集成 Flink 最头疼的是版本兼容。不同版本的 Flink 对应不同的 Hudi bundle 包搞错版本会出现各种莫名其妙的序列化异常或者 SQL 校验失败。以我实际的测试和线上使用为例目前比较稳的组合是Flink 版本Hudi 版本Hadoop 版本备注1.13.x0.12.x / 0.13.x3.1.x老项目常用稳定但功能偏旧1.14.x0.13.x / 0.14.x3.1.x / 3.2.x兼容性最好推荐生产1.15.x0.13.x / 0.14.x3.2.x / 3.3.x新项目可以选注意依赖冲突1.16.x0.14.x / 0.15.x3.2.x / 3.3.x需要更新版本的 bundle我在新项目里用的是 Flink 1.14.6 Hudi 0.14.1 Hadoop 3.2.4这个组合踩坑最少。Hudi 0.14 之后对 Flink SQL 的语法更友好automatic table 创建也稳定很多。如果你们团队已经锁定了较低版本的 Flink建议优先考虑 Hudi 0.12/0.13不要强行升到 0.14。注意尽量不要把 Hudi bundle 包和 Flink 自带的 lib 里相同类库混用尤其是 avro、parquet、hadoop-common。版本冲突基本都会导致作业在提交阶段直接失败。2.2 依赖引入Maven 还是 Flink lib如果你用 DataStream API 开发需要在 Maven 项目里引入 hudi-flink 相关依赖。这个时候要注意把 flink 相关的 scope 设置成 provided避免和运行环境的 Flink 冲突。properties flink.version1.14.6/flink.version hudi.version0.14.1/hudi.version /properties dependencies dependency groupIdorg.apache.hudi/groupId artifactIdhudi-flink1.14-bundle/artifactId version${hudi.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_2.12/artifactId version${flink.version}/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge_2.12/artifactId version${flink.version}/version scopeprovided/scope /dependency /dependencies如果你主要在 Flink SQL 客户端里跑就不需要 Maven直接把hudi-flink1.14-bundle.jar放到$FLINK_HOME/lib目录下。这里提醒一下bundle 包已经自带了依赖的 hadoop、avro、parquet 等库放到 lib 之后先用flink sql客户端执行一条简单的 show tables 确认没报类冲突再做正式开发。2.3 快速启动Docker 安装 Flink 并加载 Hudi 插件很多刚接触 Flink 的朋友喜欢用 Docker 快速起一个集群练手这个路子没问题。我自己也经常用来验证 Hudi 的新版本特性比搭物理机快得多。最简单的方式是用 docker-compose 启动一个 Flink standalone 集群version: 2 services: jobmanager: image: flink:1.14.6 ports: - 8081:8081 command: jobmanager environment: - | FLINK_PROPERTIES jobmanager.rpc.address: jobmanager taskmanager: image: flink:1.14.6 depends_on: - jobmanager command: taskmanager environment: - | FLINK_PROPERTIES taskmanager.numberOfTaskSlots: 4 jobmanager.rpc.address: jobmanager启动后进入 jobmanager 容器把 Hudi bundle 包下载到/opt/flink/lib重启 jobmanager 和 taskmanager再打开 Flink SQL 客户端就默认有 Hudi connector 了。Docker 方式适合功能验证生产不建议这样。生产环境建议把 Flink 跑在独立集群或者 K8s 上Hudi 需要的依赖包最好在构建镜像时固定好不要等运行再手动放入。否则后续扩容或者任务重启时环境不一致排查起来非常痛苦。3. 核心实操Flink SQL 实时写入 Hudi3.1 Hudi 建表语法与关键参数Hudi 在 Flink SQL 里的接入方式很简单核心是配置connector为hudi。建表语法兼容 Flink 的标准 DDL但在表属性里有一堆hoodie.开头的参数。新手最容易懵的地方就是参数太多不知道哪些必填哪些可以默认。一个最基础的 Hudi 表定义长这样CREATE TABLE hudi_orders ( order_id BIGINT PRIMARY KEY NOT ENFORCED, user_id BIGINT, product_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3), partition_day STRING ) PARTITIONED BY (partition_day) WITH ( connector hudi, path s3://my-bucket/ods/hudi_orders, table.type COPY_ON_WRITE, hoodie.datasource.write.recordkey.field order_id, hoodie.datasource.write.precombine.field order_time, hoodie.datasource.write.partitionpath.field partition_day, write.tasks 4, hive_sync.enable true, hive_sync.mode hms, hive_sync.metastore.uris thrift://hive-metastore:9083, hive_sync.db ods, hive_sync.table hudi_orders );这里我解释几个关键点。table.type有两个值COPY_ON_WRITE和MERGE_ON_READ。COW 每次写入都会重写受影响的数据文件查询快但写放大明显MOR 用 log file 记录增量写入后台异步合并写路径更轻但查询时需要合并 log查询延迟更高。实时入湖如果用 COW数据量不大时完全没问题如果写入并发很高、更新频繁建议用 MOR配合异步 compaction。recordkey.field是主键字段对应 Hudi 的 record key所有 upsert 都靠它定位记录。precombine.field是关键中的关键当同一条 record 出现多条版本时Hudi 会按照这个字段的值决定哪条数据最新。对于流式环境通常取时间字段比如数据的业务时间或写入时间千万不要用无意义的常量。partitionpath.field是分区字段建议在已经做过分区选择的情况下把它和 Flink 表的分区字段对齐。分区字段选小了比如按天数据倾斜之后小文件会很多选大了比如按小时有可能产生过细的分区。具体还是结合下游查询频率和数据量来定。3.2 完整示例从 Kafka 到 Hudi 的实时入湖下面给一个完整的 Flink SQL 作业示例这是一个我常用的模板。假设 Kafka 里有订单变更流我们把数据清洗后写入 Hudi。先创建 Kafka 源表CREATE TABLE kafka_orders ( order_id BIGINT, user_id BIGINT, product_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3), operation STRING, WATERMARK FOR order_time AS order_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic ods_orders, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, properties.group.id flink-hudi-orders, scan.startup.mode earliest-offset, format debezium-json, debezium-json.schema-include false );然后创建 Hudi 目标表CREATE TABLE hudi_orders ( order_id BIGINT PRIMARY KEY NOT ENFORCED, user_id BIGINT, product_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3), partition_day STRING ) PARTITIONED BY (partition_day) WITH ( connector hudi, path s3://my-bucket/ods/hudi_orders, table.type MERGE_ON_READ, hoodie.datasource.write.recordkey.field order_id, hoodie.datasource.write.precombine.field order_time, write.tasks 4, compaction.async.enable true, compaction.tasks 2, compaction.trigger.strategy num_commits, compaction.delta_commits 5 );最后提交写入任务INSERT INTO hudi_orders SELECT order_id, user_id, product_id, amount, order_time, DATE_FORMAT(order_time, yyyy-MM-dd) AS partition_day FROM kafka_orders WHERE operation DELETE;这段作业跑起来后Kafka 里的每个订单事件都会实时写入 Hudi。commit 完成后你可以在 Spark 或 Presto 里查询这张表看到最新的快照数据。注意我把table.type设置成了MERGE_ON_READ同时开启了异步 compaction。compaction 的触发策略配置成了每经过 5 个 delta commit 就自动合并一次合并任务并发度是 2这样能把 log 文件控制在一个合理的范围内。如果你的数据源是全量 增量混合也可以用 DataStream 的方式直接用 HoodieFlinkStreamer但多数情况 SQL 已经够用而且 SQL 对后续维护更友好。3.3 核心机制解析upsert、索引与合并很多人以为 Hudi 的 upsert 就是把 flink 的输出根据主键删除再插入这么理解虽然大方向没错但实现细节远不止这么简单。Hudi 表由无数个 file group 组成每个 file group 包含若干 file slice。当一条数据进来后先通过索引比如简单的 record key hash 索引或 bucket 索引定位这条数据应该属于哪个 file group然后决定是更新已有 file slice 还是插入新的 file group。如果是更新COW 表会重写整个 file slice 中包含相关记录的所有列式文件MOR 表则先追加到 log 文件等 compaction 再合并。由于 Flink 是流式处理如果直接每条数据都触发文件重写IO 开销会炸掉。所以 Hudi 的 Flink sink 在内部做了缓冲和批量提交数据按检查点或者 buffer 时间攒批达到阈值之后统一写文件并提交 commit。这样既保证了事务性也减少了小文件数量。和 commit 联动的是 Flink checkpoint。Flink 的 checkpoint barrier 传播到 Hudi sink 时sink 会把当前批次的数据刷到存储并生成一个 Hudi commit成功后再确认 checkpoint。如果作业崩溃或者重启Flink 会从最近的 checkpoint 恢复同时 Hudi 也只会保留已经 commit 的数据未提交的临时文件会被清理。所以两边的机制绑定在一起才能实现真正意义上的 exactly-once。索引策略上最简单的SIMPLE_INDEX是拿 record key 和 partition path 直接去遍历相关文件目录定位数据所在位置。数据量小的时候没问题数据量大就会出现明显性能瓶颈。生产上我建议使用BUCKET_INDEX它把 record key 的一致性哈希运算对应到固定数量的 bucket写入和检索都只访问对应 bucket 里的文件。配置方式hoodie.index.type BUCKET, hoodie.bucket.index.hash.field order_id, hoodie.bucket.index.num.buckets 16bucket 数量要结合文件大小和写入并行度来设我一般先按每个 bucket 输出数据量估算确保单个 bucket 落地的数据文件在 128MB 到 512MB 之间太少会小文件多太多则合并压力大。4. 性能调优与问题排查4.1 写入性能与资源调优Hudi 集成 Flink 的性能调优核心围绕三件事并行度、批量大小、小文件治理。先说并行度。Flink 写入 Hudi 的 sink 算子并行度write.tasks不一定要等于 Flink 作业全局并行度。如果上游数据量小设置过大的 write.tasks 会生成很多空文件反而拖垮 commit 性能。我一般让 write.tasks 等于数据源分区的 1 到 2 倍同时观察存储生成的文件大小如果单个文件长期低于 64MB就要降低并行度或者调大 buffer 时间。批量大小由sink.buffer-flush.max-rows和sink.buffer-flush.interval控制。Hudi 0.14 中相关参数可能在不同的版本里有命名变化但思路是一致的攒的行数越多或者间隔时间越长单次写的文件就越大commit 次数也越少。但如果间隔太长数据实时性会下降。生产上我常用sink.buffer-flush.max-rows 100000和sink.buffer-flush.interval 10s作为初始值然后根据实际延迟调整。小文件问题是 Hudi 写入最常见的病。即使设置了合理的 buffer时间长了仍然会因为分区数量多而出现大量小文件。Hudi 提供了 clustering 功能Flink 侧可以通过clustering.schedule.enabled和clustering.async.enabled来开启异步 clustering。简单理解就是对小文件进行合并重写让表保持在健康的文件大小分布。同时clean操作会清理旧版本文件archive会压缩 commit 历史元数据这些都是要打开的。clustering.schedule.enabled true, clustering.async.enabled true, clustering.tasks 4, clustering.delta_commits 10, clean.async.enabled true, archive.async.enabled true注意 clustering 和 compaction 不要同时开得过于频繁两个操作都会占用作业资源。建议 compaction 按 commit 数触发clustering 按 commit 数或者时间周期触发错峰运行。4.2 火焰图分析定位 Hudi 写入瓶颈线上 Flink 作业一旦出现延迟或者吞吐上不去我最先做的事情就是看火焰图。很多性能问题在界面上看不出来比如序列化消耗过高、磁盘 IO 等待、某个算子内部锁竞争只有火焰图能直观看出来。生成火焰图我一般用 async-profiler。先找到 TaskManager 进程 PID在作业运行一段时间后采集./profiler.sh -d 60 -o flamegraph -i 5ms PID hudi_sink_flamegraph.html然后把生成的 HTML 文件下载到本地用浏览器打开直接看哪个函数的自顶向下占比最宽。在 Hudi 写入作业里我见过最多的火焰图热点有两类。第一类是org.apache.avro相关的序列化和 schema 转换。Hudi 内部默认用 Avro 作为内存记录格式如果表字段特别多或者嵌套结构很深序列化成本会高得离谱。这时候可以考虑开启hoodie.table.avro.schema.validate相关调优但本质上要减少不必要的字段把宽表拆成窄表。第二类是文件写入时的parquet编码和压缩。Parquet 的字典编码和压缩级别会影响 CPU如果发现org.apache.parquet相关栈宽可以考虑调整压缩算法比如从snappy改成zstd或者下调压缩级别。还有一个隐藏热点是 Flink 的状态后端。Hudi sink 本身会使用 Flink 状态来缓存数据如果状态后端是 RocksDB频繁序列化状态也会在火焰图上显示为 RocksDB 相关调用。此时可以调大托管内存或者优化状态 TTL避免状态无限增长。4.3 JDBC 连接器异常及常见错误搜热词里经常有人问 Flink 的 JDBC 连接器异常这个在 Hudi 集成场景里也很常见但大多数时候不是 Hudi 本身的问题而是作业里同时用了 JDBC 维表或者 JDBC sink连接管理没做好。最常见的错误是Connection is not available, request timed out、SQLException: No suitable driver found之类。排查步骤一般是检查 flink 运行环境的 lib 目录是否放入了对应的 JDBC driver 包比如 MySQL 的mysql-connector-java.jar。注意 Flink 1.14 之后 driver 需要放到lib而不是用-j参数提交。检查 JDBC 连接池配置。Flink JDBC connector 默认有连接超时和空闲超时建议显式配置sink.buffer-flush.interval和sink.max-retries避免短连接风暴。如果 JDBC 维表关联 Hudi 写入场景最好给维表连接加lookup.cache.max-rows和lookup.cache.ttl降低对数据库的压力否则数据库连接数被打满后就会出现大量超时异常。Hudi 本身常见的异常还有schema 不匹配导致的AvroSchemaException这个通常是因为上游 topic 字段变更而 Hudi 表没有同步还有 checkpoint 超时导致的提交失败这时要看是否因为写 Hudi 的 commit 耗时超过了 checkpoint interval。应对方案是把 checkpoint 超时时间适当调大比如execution.checkpointing.timeout 10min同时调整 Hudi 写入缓冲减少每次 commit 的数据量。4.4 常见问题速查表我整理了一张速查表基本覆盖了从环境搭建到线上运行最常见的几个坑方便大家对照排查。错误现象可能原因解决办法提交作业时 ClassNotFoundExceptionbundle 包未放入 lib或与 Flink 自带依赖冲突确认 hudi-flink bundle 版本与 Flink 版本匹配干净环境验证建表报找不到 connectorSQL 客户端未加载 Hudi connector检查 lib 下的 bundle 包重启客户端写入后查不到数据数据还在 buffer 中或分区路径写错等待 flush检查 partitionpath 字段是否与 sink 表一致主键重复但更新不生效precombine.field 选错了字段改为业务时间或变更时间确保同一主键能比较出新旧MOR 表查询越来越慢compaction 未触发或资源不足增大 compaction.tasks设置合理的 delta_commitscommit 耗时过长导致 Flink checkpoint 超时单批次写入数据量过大或并发太高调整缓冲参数降低写入并行度审计 commit 时间Hive 同步失败metastore 地址或库表名配置错误检查 hive_sync.metastore.uris确认 HMS 版本兼容文件数量爆炸分区过细或 write.tasks 过大合并分区启用 clustering合理设置 write.tasksJDBC 维表查询连接超时连接池太小或缓存未开开启 lookup 缓存增大 max-retries检查 driver 版本这张表不包含所有问题但覆盖了头几个月的使用频率最高的坑。遇到新问题时建议先看 Flink Web UI 上的 TaskManager 日志再看 Hudi 的 commit 元数据基本能定位到是 Flink 侧还是 Hudi 侧的问题。5. 项目总结与经验补充5.1 我踩过的几个坑第一个坑发生在版本升级的时候。之前线上用 Flink 1.13 Hudi 0.12后来想升到 Flink 1.15直接换了高版本 bundle 包结果运行时报了NoSuchMethodError最后查了一圈发现是新的 bundle 里 hudi-common 依赖的某些类来自新版本的 hive 包和集群上的 hive 客户端冲突。解决的办法是排除 bundle 里的旧依赖或者干脆保持 Hudi 版本不变只升其中一个组件。所以每次动版本前我建议先在测试环境跑一遍完整的消费、写入、查询链路不要只看 build 通过。第二个坑是小文件治理没跟上。最开始我把 write.tasks 设置成和 Flink 并行度一样大结果一个小时生成了几千个几十 KB 的小文件。后来开了 clustering 和自动归档同时调大了write.batch.size文件数量才恢复正常。小文件问题一旦积累后续查询性能会断崖式下降甚至元数据服务都会被拖垮所以一定要在写入设计阶段就考虑。第三个坑是 precombine.field 选错了。我当时用了数据生成时间结果上游数据源因为重试导致同一个主键出现了乱序旧数据反而覆盖了新数据。后来改成业务时间并且在写入前对数据按照业务时间做了去重和排序才解决问题。要记住Hudi 不会自己判断哪个版本业务上是最新的它只信任 precombine field。5.2 值得深入的点CDC 入湖与 streaming readHudi Flink 的应用不止简单入湖。我最近在做的一块是把 MySQL 的 binlog 通过 Flink CDC 写入 Hudi然后下游用 Flink streaming read 读取增量数据。这个方案的核心是依赖 Hudi 的增量读取能力Flink 可以定期扫描 Hudi 表最近 commit 产生的数据这样就能基于数据湖实现分钟级的实时数仓链路。具体做法是配置read.start-commit为某个时间点或 commit然后 Flink 消费 Hudi 表时指定hoodie.datasource.query.type incremental。这种方式比直接用 Kafka 重放更灵活因为 Hudi 的增量数据是持久化的即使 Flink 作业重启也可以从指定 commit 开始消费不会丢失。如果你现在已经接了 Kafka 到 Hudi建议下一个阶段试试用streaming read把 Hudi 表当成一个消息队列来复用能减少很多中间环节。5.3 如果面试遇到 Hudi Flink最近身边不少朋友在准备 Flink 面试我发现 Hudi 相关的问题出现频率一下子高了起来。面试官通常不会问太细的参数而是问原理层面的问题Hudi 是怎么支持 upsert 的Flink 写入 Hudi 如何保证 exactly-onceCOW 和 MOR 的区别是什么回答这些问题的关键是把 file group、file slice、index、precombine、commit 和 checkpoint 的关系讲清楚。你可以这么组织思路Hudi 通过索引将记录定位到 file group再按照 record key precombine field 决定是否更新如果不是更新就生成新的 file slice 或者追加到 log 文件Flink 通过 checkpoint 对齐批次在 checkpoint 成功时生成 commit保证数据要么全部写入提交、要么不出现。如果能在回答里举一个实际调优的例子比如通过调整 write.tasks 减小了小文件或者是用 BUCKET 索引解决了更新性能问题面试官会很认可。这个生态发展很快Hudi 官方也在持续优化 Flink 的集成体验。我最后再分享一个小技巧多去看 Hudi Flink 相关 issue 和 release notes有时候一个看似诡异的问题新版本已经静默解决了。升级前尽量把 release notes 里提到的 breaking change 列出来对照自己的作业逐项检查能省下很多排查时间。最好的方式不是停留在会用而是能理解每个参数的底层逻辑这样面对任何版本变化都能快速适应。
返回列表