ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Hive与Hudi整合实战:解决增量数据读写难题的全流程指南

Hive与Hudi整合实战:解决增量数据读写难题的全流程指南 Hive做增量数据这件事但凡真正在数仓或者数据平台上干过的人心里都有数。最典型一个场景一张网约车订单表一天几千个分区凌晨的调度任务把昨天的增量数据全量刷一遍跑两三个小时算快的碰上数据倾斜直接跑穿告警线。更头疼的是业务说我要更新昨天某几条订单的状态Hive这张表压根不支持行级更新只能把整分区重写。我自己在好几个项目里被这件事折磨过直到把Hudi引进来和Hive做整合增量数据的读写才算有了个正经解法。这篇就围绕Hive与Hudi整合这条主线把我从环境搭建、建表、同步、增量查询到生产调优的完整过程掰开揉碎讲清楚。适合正在搞数据湖、或者准备给Hive集群引入增量处理能力的工程师参考基本原理和踩坑细节我都会说到。1. Hive做增量处理的老大难以及Hudi为什么能补位先说清楚一个事实Hive本身不是为增量而生的。它的设计哲学是一次写入、多次读取底层文件都是不可变的Parquet/ORC。所以你在Hive里做增量通常只有两条路要么靠分区裁剪只读最近几个分区要么靠一个时间戳字段全表扫一遍再过滤。两条路都是歪路前者要求业务严格按分区组织数据后者随着表体积膨胀耗时会变得完全不可接受。1.1 痛点一没有行级更新能力Hive 3.x虽然有了ACID支持支持INSERT/UPDATE/DELETE但实际用起来限制很多表必须做成ORC格式、必须开启事务、每次操作会产生一堆delta文件而且得定时做compaction合并。在多任务并发写同一个表的时候性能问题更加明显。我试过在Hive 3.1.3上开启ACID跑一个千万级别的更新任务小文件瞬间把表的文件数拉高了几倍后续所有查询都被拖慢。而且Hive ACID的并发控制很粗糙多Session同时更新同一张表很容易出现死锁或者版本冲突。生产上我对它的定位就是能用但别指望靠它做大规模增量更新。1.2 痛点二增量拉取要付出全量扫描的代价假设你有张10亿行的订单表每天新增两百万行你想取昨天变化的订单最朴素的做法就是SELECT * FROM orders WHERE update_time 2025-01-01 00:00:00 AND update_time 2025-01-02 00:00:00;这个SQL虽然语义对但物理执行上MRScan依然会读取所有parquet文件的row group只是最后把不符合的行过滤掉。在10亿行级别的表上过滤性能再优化磁盘IO和CPU开销都省不下来等于你用增量查询的语法付出了全量扫描的成本。1.3 Hudi的补位逻辑在HDFS上放一层可变更的数据集HudiHadoop Upserts Deletes and Incrementals做的事情简单说就是给HDFS上的数据文件增加了一套事务日志 索引 时间线机制。它把一张表的数据拆成一个个Base文件Parquet每次写入commit都记录在时间线Timeline上写入时通过索引定位记录所在文件实现行级更新读取时通过时间线知道哪个文件是新的实现快照读或增量读。关键点是Hudi生成的文件依然以Parquet为主Hive的InputFormat完全能识别所以Hudi可以和Hive共用一套元数据Hive Metastore。这就是整合的基础Hudi不替代Hive它只是把底层的文件组织方式变得更聪明而Hive还傻乎乎地以为自己在读一张普通表。从使用层面看整合完之后你的查询入口还是Hive但底层得由Hudi的StorageHandler来解析文件路径和增量提交信息。这个机制后面详细讲。2. 环境准备版本兼容矩阵与bundle包的选择这个环节最容易劝退人。Hudi的版本和Spark、Hive的兼容关系非常微妙网上教程版本混乱很多坑都是版本不匹配导致的。我最后稳定跑通的生产组合是Hudi 0.13.1 Spark 3.3.2 Hive 3.1.3 Hadoop 3.1.3我用的是CDH 6.3.x的影子但自带的Hadoop版本兼容验证过没问题。2.1 版本兼容矩阵速查我整理了一张实际踩坑后校准过的对应关系供参考Hudi版本Spark版本Hive版本备注0.10.x3.0.x / 3.1.x2.3.x / 3.x架构验证可用但缺一些新特性0.12.13.2.x / 3.3.x2.3.x / 3.x较稳定常见的生产版本0.13.13.2.x / 3.3.x2.3.x / 3.x我最终选用的版本对Hive sync支持最完整0.14.x3.3.x / 3.4.x2.3.x / 3.x新版建议新项目直接用选0.13.1的另外一个原因是它内置的Hive同步工具HoodieHiveSyncClient对HMS的兼容性最好同步表结构时不太容易报ClassNotFound。0.14.x我也试过但当时Spark 3.4和已有集群上的一些UDF冲突降回0.13.1反而省心。2.2 怎么拿bundle包Hudi打出了多个profile的bundle jar分别对Spark 2、Spark 3.0、Spark 3.1、Spark 3.2/3.3。如果你直接用Spark 3.3那应该下hudi-spark3.3-bundle_2.12-0.13.1.jar不建议下载不带spark版本的纯hudi-bundle那种往往不包含Spark的适配器启动spark-sql时会报方法签名不匹配。如果是离线环境可以提前下载好jar包放到每个节点启动时用--jars参数指定即可。2.3 编译Hudi源码时的两个注意点如果团队想自己改源码或者公司安全要求必须内网编译那需要用Maven编译对应profilemvn clean package -DskipTests -Dspark3.3 -Pflink-1.15编译时有几个坑值得说一下编译前确认JDK版本Hudi 0.13.1在JDK8和JDK11下都能编但某些模块比如flink整合在JDK11下需要额外参数。如果只要Spark模块别加-Pflink否则编译时间翻倍还会下载一堆你根本用不到的依赖。编译完成后找到packaging/hudi-spark-bundle/target/目录下的bundle包那才是你真正要放到节点上的东西。2.4 启动spark-sql的关键参数我每次连Hudi都会用一个固定的启动脚本核心参数如下spark-sql \ --jars /opt/hudi/hudi-spark3.3-bundle_2.12-0.13.1.jar \ --conf spark.sql.hive.metastore.version3.1.3 \ --conf spark.sql.hive.metastore.jarsmaven \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.sql.catalogImplementationhive \ --conf spark.driver.extraClassPath/etc/hive/conf \ --conf spark.executor.extraClassPath/etc/hive/confspark.sql.hive.metastore.version必须和你的HMS版本严格一致否则同步表结构时容易抛MetaException。spark.serializer用KryoHudi内部对象序列化要求。spark.sql.catalogImplementationhive让Spark把HMS作为catalog这样spark-sql里建的表才会同步到Hive。另一点如果HMS走的是Kerberos认证记得先kinit否则后面所有Hudi写操作都会挂在认证上报错信息还特别隐晦。3. 建表与Hive看得见Hudi的同步机制环境准备好之后真正的整合操作其实就浓缩成一句话用Spark SQL建一张Hudi表Hudi通过HMS的同步钩子自动在Hive里注册一张外表指向同样的数据路径。这张外表就是Hive查询Hudi数据的入口。3.1 用spark-sql创建一张COW表我把网约车订单场景简化了一下建表语句如下CREATE TABLE ods_mock_order ( ts BIGINT, order_id STRING, driver_id STRING, passenger_id STRING, fare DOUBLE, city STRING ) USING HUDI TBLPROPERTIES ( type cow, primaryKey order_id, preCombineField ts, hoodie.datasource.hive_sync.enable true, hoodie.datasource.hive_sync.mode hms, hoodie.datasource.hive_sync.metastore.uris thrift://hadoop01:9083, hoodie.datasource.hive_sync.db default, hoodie.datasource.hive_sync.table ods_mock_order, hoodie.table.name ods_mock_order ) PARTITIONED BY (city);对着参数说几个重点typecowCopy On Write写时复制。后面会对比MOR这里先记住COW适合更新频繁但查询要求实时一致的场景。primaryKey主键是Hudi做upsert的定位依据。Hudi没有主键约束的概念它更像一个逻辑主键用来判断一条记录是新增还是更新。preCombineField这个字段很关键。如果同一批次写入的数据里同一个主键出现多条记录Hudi就用preCombineField来选最新的一条。我一般直接用业务时间戳不用Hudi内部时间。3.2 同步机制HoodieHiveSyncTool到底做了什么表创建成功后Hudi的后台同步工具HoodieHiveSyncTool会做三件事在HMS里注册一张Hive表Location指向Hudi表在HDFS上的路径。给这张Hive表的SerDe设置成org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDeInputFormat设置成Hudi提供的HoodieParquetInputFormat。在Hive表属性里注入hoodie.table.name、hoodie.table.type等信息让Hive的查询引擎能识别出这是一张Hudi表。你在Hive CLI里执行SHOW CREATE TABLE ods_mock_order会看到类似CREATE TABLE ods_mock_order( ts bigint, order_id string, driver_id string, passenger_id string, fare double) PARTITIONED BY ( city string) ROW FORMAT SERDE org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe STORED AS INPUTFORMAT org.apache.hudi.hadoop.hive.HoodieParquetInputFormat OUTPUTFORMAT org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat LOCATION hdfs://hadoop01:8020/user/hive/warehouse/ods_mock_order TBLPROPERTIES ( hoodie.table.typeCOPY_ON_WRITE, ... )这就是Hive和Hudi整合的核心桥梁。Hive读这张表时InputFormat是Hudi提供的实现它会读取Hudi的时间线决定这个分区下哪些Base文件是当前可见的。比较老的、已经被更新过的Base文件就直接过滤掉了。3.3 插入数据并验证Hive端可见建完表先插入一批数据INSERT INTO ods_mock_order VALUES (1700000000000, 1001, d001, p001, 35.5, shanghai), (1700000000000, 1002, d002, p002, 42.0, beijing);这条INSERT的底层实际是Spark执行了一次bulk_insert生成两个parquet文件按city分区各一个。这时候去Hive CLI执行SELECT * FROM ods_mock_order;如果返回了两行说明整合链路已经通了。有些刚上手的朋友在这里会遇到一个问题Hive能查出表结构但查不出数据或者报文件找不到。原因多半是Hive用户没有HDFS目录的读权限或者LOCATION路径配置不一致用dfs -ls /user/hive/warehouse/ods_mock_order确认一下路径就行。注意如果你在Hive CLI里建了一张同名表但指向不同路径那和Hudi表是两个表Hive侧查询会看不见Hudi新写入的文件。所以务必让Hudi自动同步不要去手动建Hive表。4. 增量查询、快照查询、时间旅行三种姿势的实操与取舍整合之后Hudi表在Hive里能查到但这只是基础。真正体现增量处理价值的是接下来的三种查询模式。4.1 快照查询快照查询是默认方式查的是当前最新可见数据等价于SELECT * FROM ods_mock_order。它读的是每个文件Slice里最新的Base文件和普通Hive查询的执行计划几乎一样唯一区别是InputFormat相关代码会在并行读取时做一次文件列表裁剪。快照查询适合日常报表、即席分析也是Hive侧最自然的查询方式。性能上如果你的Hudi表文件排布合理基本和直接查Parquet表差不多。4.2 增量查询增量查询是Hudi的核心能力。它的语义是给我一次commit之后发生变化的所有行。在Spark SQL里我们通常借助Hudi表自动带出的_hoodie_commit_time隐藏字段来做过滤SELECT * FROM ods_mock_order WHERE _hoodie_commit_time 20250101000000;这个字段是Hudi写入时自动打上的commit时间戳精确到毫秒。执行计划上Hudi的InputFormat会利用时间线信息只读取发生变更的文件Slice而不是整表扫描。这里有个容易踩坑的点增量查询时如果对非隐藏字段加了过滤条件可能会触发Spark把全部文件重新扫描因为Hudi的文件裁剪逻辑必须基于_hoodie_commit_time才能生效。所以增量查询的推荐写法就是把commit_time放在WHERE最前面其他过滤条件放到子查询里。举个例子要算昨天每个城市新增订单数和GMVSELECT city, COUNT(*) AS cnt, SUM(fare) AS gmv FROM ( SELECT order_id, city, fare FROM ods_mock_order WHERE _hoodie_commit_time 20250101000000 AND _hoodie_commit_time 20250102000000 ) t GROUP BY city;批跑任务把前一天凌晨的commit_time边界算好跑这个查询Spark扫描的文件量只和当天写入的文件有关系。在增量数据占比很小的情况下性能比全量扫描能提升一个数量级。4.3 时间旅行查询时间旅行查询是指查询某个历史时间点的快照常用于数据回刷、审计或者我当时看的数为什么和现在不一样的排查。Hudi在0.13.1里通过Spark DataFrame的asOfInstant来支持val df spark.read.format(hudi) .option(as.of.instant, 20250101120000000) .load(hdfs://hadoop01:8020/user/hive/warehouse/ods_mock_order) df.createOrReplaceTempView(ods_mock_order_snap)SQL层面直接用spark-sql命令行不好指定option我一般是写一个小的Scala脚本封装。这个功能在排查数据质量问题的时候非常有用。比如业务质疑今天某个订单金额不对你可以直接定位到它写入commit时间对应的快照看那份数据当时长什么样。4.4 MOR表的读优化查询如果建表时用了typemorMerge On Read那Hudi会拆成Base文件和增量Log文件读的时候要把两者合并会造成一定的读放大。为了缓解Hudi提供了一种读优化查询只读Base文件不管Log里的增量SELECT * FROM ods_mock_order_mor_ro;注意Hudi在HMS里注册MOR表的时候会自动生成两张表一张表名带log语义一张表名_ro。_ro就是读优化视图适合对实时性要求不高的场景。我个人的选型经验查询压力大、更新量不是特别极端的场景优先用COW如果写入方是流水型大量追加、偶尔更新用MOR能省下大量写放大成本。具体对比我放到最后一部分讲。5. 上线之后躲不开的三个生产问题小文件、同步失效与写放大环境通了增量查询也跑通了真正到了生产环境你会发现一堆坑在等着。这里说三个我实际遇到并解决过的每个都有完整排查过程。5.1 小文件问题Hudi表越写越碎Hudi默认并发写的时候如果并行度设置不当每个Shuffle分片都会输出一个小Parquet文件。你设了1000并行度就生成几千个小文件。几轮commit下来文件数爆炸NameNode压力增大Hive查询的计划步骤变多性能骤降。这个问题的根源在于并行度和文件大小没有联动。Hudi针对插入时的文件大小有参数控制但很多人建表之后根本不去调默认配置偏保守。我的做法是两条线同时改ALTER TABLE ods_mock_order SET TBLPROPERTIES ( hoodie.parquet.small.file.limit 134217728, hoodie.insert.shuffle.parallelism 4 );hoodie.parquet.small.file.limit设成128MB意味着小于128MB的文件Hudi会尝试把新记录追加到旧文件里。insert.shuffle.parallelism控制插入时的shuffle分区数这个值不需要设太大我通常按executor核数的两倍来定。如果文件已经碎到没法靠参数挽救就开inline clustering内联聚类让Hudi在commit次数达到阈值后自动合并文件SET TBLPROPERTIES ( hoodie.clustering.inline true, hoodie.clustering.inline.max.commits 4, hoodie.clustering.plan.strategy.sort.columns order_id );这里有个细节clustering和compact在COW表里是同一个概念但MOR表里clustering针对的是Base文件compaction针对的是Log文件别搞混。5.2 Hive同步失效增量数据查不到遇到过一次诡异现象Spark SQL写入Hudi表后Hive CLI查询时数据量对不上。排查过程是这样一步步缩小的第一步确认写入成功。在HDFS上看对应分区文件确实存在commit时间线也有新commit。第二步查HMS里Hudi表的INPUTFORMAT。发现还是普通的MapredParquetInputFormat而不是Hudi提供的HoodieParquetInputFormat。这就意味着Hive绕过了Hudi的文件裁剪逻辑直接读路径下所有文件把旧的Base文件和新文件一起读了导致结果重复。原因是我当时手动在Hive里执行了ALTER TABLE ... SET LOCATION这个操作把Hudi注入的StorageHandler信息覆盖掉了。解决办法是重建同步信息/opt/hudi/hoodie-hive-sync.jar \ --database default \ --table ods_mock_order \ --base-path hdfs://hadoop01:8020/user/hive/warehouse/ods_mock_order \ --partitioned-by city \ --sync-mode hms \ --hive-metastore-uris thrift://hadoop01:9083跑完之后再查HMSINPUTFORMAT回来了。这件事也提醒我Hudi表在Hive里的注册信息非常脆弱千万不要手贱去改它的表结构或者Location。如果确实要改正确姿势是先用Hudi的工具同步一次再去动HMS。5.3 写放大问题COW的代价COW表每次upsert都会重写受影响的数据文件如果一个文件里只有一行需要更新整个文件都要重写。这就是写放大。在更新分布很分散、每个文件都要碰一下的场景下写放大指数增长。我们有个订单状态更新任务每天要更新上一小时导入的约五万行但主键比较分散。用COW表跑一轮写了近200GB的数据而真正变更的数据不到1GB。后来改用MOR表CREATE TABLE ods_mock_order_mor ( ts BIGINT, order_id STRING, driver_id STRING, passenger_id STRING, fare DOUBLE, city STRING ) USING HUDI TBLPROPERTIES ( type mor, primaryKey order_id, preCombineField ts, ... ) PARTITIONED BY (city);同样的任务写放大从200GB降到了不到20GB。代价是查询发生更新的文件时需要同时读Base文件和Log文件做合并读性能有一定下降。但我们在查询层做了缓存对生产影响不大。这里给一个通用决策表帮助选型场景推荐表类型原因订单明细、交易流水大量追加少更新COW读性能好文件规整状态类数据频繁更新小范围记录MOR写放大低实时性可接受用Hive做常规报表查询为主COW查询稳定不需要合并log增量数据实时性要求高且写入量大MOR写入快读时通过compaction补偿6. 一个真实增量指标计算实验与最终选型建议前面讲的都是机制和参数最后用一个实际跑过的例子收尾。这个例子是我在一个网约车综合数据项目里做过的订单增量指标小时级报表简化后非常适合拿来复现。6.1 实验场景设计数据源模拟订单流每小时写入一次每次约20万条其中大约5%是对之前订单的更新比如司机取消、费用调整。目标Hive端每小时跑一次增量SQL统计最近两小时新增订单数、取消率、每小时GMV。对比对象同一张全量Hive表直接读全部文件过滤时间 vs Hudi增量查询。6.2 数据写入写入用Spark Structured Streaming接Kafkabatch模式每小时触发一次import org.apache.hudi.DataSourceWriteOptions._ import org.apache.hudi.config.HoodieWriteConfig import org.apache.hudi.hive.HiveSyncConfig inputDF.write.format(hudi) .option(OPERATION.key(), UPSERT_OPERATION_OPT_VAL) .option(RECORDKEY_FIELD.key(), order_id) .option(PRECOMBINE_FIELD.key(), ts) .option(hoodie.datasource.hive_sync.enable, true) .option(hoodie.datasource.hive_sync.use_jdbc, false) .option(hoodie.datasource.hive_sync.mode, hms) .option(hoodie.datasource.hive_sync.metastore.uris, thrift://hadoop01:9083) .option(hoodie.datasource.write.hive_style_partitioning, true) .option(HoodieWriteConfig.TABLE_NAME.key(), ods_mock_order) .mode(Append) .save(hdfs://hadoop01:8020/user/hive/warehouse/ods_mock_order)这里UPSERT_OPERATION_OPT_VAL会把新数据和历史数据做比对有相同主键就更新没有就插入正好模拟订单状态变化。6.3 Hive侧增量SQL任务调度脚本里先算好两个时间边界BEGIN_TIME$(date -d -2 hours %Y%m%d%H%M%S) END_TIME$(date -d -1 hour %Y%m%d%H%M%S)然后跑SQLINSERT OVERWRITE TABLE dm_order_hourly_stats SELECT city, COUNT(DISTINCT order_id) AS order_cnt, SUM(CASE WHEN status cancelled THEN 1 ELSE 0 END) / COUNT(*) AS cancel_rate, SUM(fare) AS gmv FROM ( SELECT order_id, city, fare, status FROM ods_mock_order WHERE _hoodie_commit_time ${BEGIN_TIME} AND _hoodie_commit_time ${END_TIME} ) t GROUP BY city;因为这个SQL的表名是Hudi自动同步到HMS里的Hive侧的执行引擎还是走HoodieParquetInputFormat所以只读取两个小时commit里涉及的文件而不是整个大表。6.4 实测结果这张表跑到5亿行量级、3000个分区时全量扫描方案大概要4分20秒才能出结果。Hudi增量查询方案稳定在22秒左右文件读取量只有全量方案的1/10到1/20。代价是建表和写入时的配置多了一些但换来的是下游分析任务的稳定性和时效性。另外在这轮实验里我把小文件治理加进去了。每小时写入20万条靠着hoodie.insert.shuffle.parallelism4和hoodie.parquet.small.file.limit128MB的配合整个表的分区文件数保持在一个很健康的水平Hive端查询执行计划步骤少了很多。6.5 我对HiveHudi组合的最终建议如果你的场景是数据分析的活儿还得靠Hive生态但数据源里的更新和新增想要做到小时级甚至分钟级可见那HiveHudi这个组合值得认真考虑。选型落地时这几条经验应该能帮你少走弯路表类型优先按读写比决定别盲从COW就是好或MOR就是好的说法。写多读少、更新密集用MOR读多写少用COW。建表时就把小文件参数设好不要等到表大了再来调。上生产前用一个模拟数据把文件数验证一遍比事后加Clustering省力得多。Hudi自动同步HMS表结构以后把Hive那边的手动DDL权限收掉防止有人为了加个注释把StorageHandler覆盖了。增量查询的WHERE条件里_hoodie_commit_time一定要放在最外层不要放进子查询再包一层否则Spark可能把裁剪逻辑优化掉。批任务里用时间边界跑增量查询比依赖Hudi自带的Streaming增量模式更容易和大数据平台现有的调度体系融合。最后说一个个人体会接触Hudi之前我一直觉得数据湖是个营销词直到用增量查询把小时级报表的扫描量降了一个数量级才真正认可它的价值。Hive和Hudi的整合不是把Hive替换掉而是给Hive插上了一双能处理增量数据的翅膀至少在未来的两三年里这张组合拳依然是离线数仓里性价比很高的方案。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进