ARTICLE DETAIL

资讯详情

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

Hive与Doris整合实践:MPP加速离线数仓查询的架构与同步链路详解

Hive与Doris整合实践:MPP加速离线数仓查询的架构与同步链路详解 做离线数仓项目时最常被业务方问的一句话是这张大表能不能跑快点Hive跑一个聚合报表动不动就是五六分钟甚至半小时业务要的是秒级响应。这个矛盾在网约车订单分析、电商日活报表这类场景里特别明显底层的Hive数据一直增长SQL逻辑并不复杂但MapReduce的批处理模型决定了它很难给到交互式体验。于是Doris这类MPP数据库开始越来越多地出现在数仓链路里——Hive继续承担离线ETL和海量存储Doris接手查询分析加速两套引擎配合着用。这篇文章就围绕Hive与Doris整合这条主线把我实际落地过程中的架构选型、部署细节、同步链路、性能调优和踩坑记录完整写出来。适合已经在用Hive做数仓、但被查询延迟困扰的工程师也适合正在选型MPP引擎、想了解Doris到底怎么接入Hive的同学。1. Hive慢在哪里、Doris快在哪里MPP加速依赖的底层差异1.1 Hive查询慢的两个硬伤Hive慢不是没优化好而是它的计算模型天生就不适合交互式查询。Hive默认走MapReduce一个简单的GROUP BY也要经过Map端读取、Shuffle排序、Reduce端聚合中间结果大量落盘。即便很多公司已经把执行引擎换成了Tez或者Spark本质仍然是批处理任务启动有调度开销数据要经过多轮Shuffle磁盘I/O占了大部分时间。另一个硬伤是存储与计算分离带来的网络开销。Hive表的数据放在HDFS上Map任务要跨节点拉数据数据本地性只能尽量保证没法做到极致。当业务方同时跑十几张报表YARN队列一挤查询延迟就会进一步放大。说白了Hive是为跑完就算成功设计的不是为人等结果设计的。1.2 MPP引擎为什么能跑出极速Doris是典型的MPPMassively Parallel Processing架构这个概念的直观理解是一条SQL进来不是由一个任务串行处理而是由几十个BE节点各管各的数据分片Tablet同时干活再汇总。Doris内部有两个核心角色。FEFrontend负责接收MySQL协议请求、解析SQL、生成分布式执行计划、管理元数据和副本调度BEBackend负责真正的数据存储和计算执行。FE把一条SQL拆成多个PlanFragment下发给各BE每个BE只扫描自己本地磁盘上的那一部分数据。配合列式存储、前缀索引、ZoneMap索引和向量化执行引擎数据在内存里一批一批流动全程几乎不落盘。我用一个具体数字说明差距。同样一张5亿行、按天分区的订单明细表Hive里跑统计某城市一周的订单量和GMVSpark引擎大约需要40秒Doris里跑同样的SQL全表走本地索引加并行扫描基本在2秒以内。这就是MPP引擎的价值——你不缺数据缺的是让人等得起。1.3 Hive和Doris的正确分工谁也不能替代谁有一点必须想明白Doris不是用来取代Hive的。Hive的生态完整度、UDF丰富程度、与Spark/Flink的配合能力、以及基于HDFS的超低成本存储都是Doris短期内比不上的。PB级原始数据的清洗加工还是得靠Hive。真正的落地姿势是分工Hive继续做全量数据的离线加工和分层建模产出质量可控的结果表Doris承接结果表的明细查询、固定报表、多维分析和即时探查。数据体量特别大又不要求秒级响应的场景留在Hive需要秒级交互的场景把数据同步进Doris。这也是Hive与Doris整合这句话的核心含义——不是二选一而是把各自的优势拼起来。2. Doris集群落地与Hive Catalog打通部署阶段的取舍2.1 最小生产集群怎么搭Doris部署本身不复杂但有不少细节会被官网文档一笔带过。先讲规模。生产环境我建议至少3个FE节点1个Master 2个Follower通过内置Paxos选主 3个或更多BE节点。测试环境可以1个FE 1个BE甚至单节点跑通但你用测试集群得出的性能结论别直接套到生产BE只有1个时MPP并行度根本体现不出来。机器配置上FE是轻量组件8核16G就够BE是重活主力建议16核64G起步磁盘用多块SSD做多目录。操作系统层面有几个坑要先处理关闭swapDoris官方要求swapoff -a否则内存交换会带来极大的查询抖动。把文件句柄数调到655360以上BE高并发时文件句柄很容易打满。FE和BE节点之间必须做时钟同步chrony或ntp时间偏移超过阈值会出现副本误判、心跳异常。端口要提前规划好FE的8030是Web UI、9030是MySQL协议端口BE的8040用于HTTP、9060用于Thrift、9070用于BE之间的BRPC通信。我第一次部署时就是全用默认端口没规划后来跟公司安全组策略撞了排查半天。下载好官方二进制包后FE目录下执行sh bin/start_fe.sh --daemon启动BE目录下执行sh bin/start_be.sh --daemon启动再登录MySQL客户端执行ALTER SYSTEM ADD BACKEND be_host:9050;BE默认心跳端口是9050注意不是9060我见过有人真在9060上栽过跟头。2.2 创建Hive Catalog的正确姿势集群起来之后最重要的一步是把Hive的元数据接入Doris。Doris从1.2版本开始支持Catalogs可以直连Hive Metastore不需要把数据搬过来就能查询Hive表。创建Catalog的SQL很简单CREATE CATALOG hive_catalog PROPERTIES ( type hms, hive.metastore.uris thrift://hive-metastore-host:9083 );创建完执行SHOW CATALOGS确认存在然后就能用三层命名方式查Hive数据SHOW TABLES FROM hive_catalog.default_db; SELECT COUNT(*) FROM hive_catalog.default_db.ods_order WHERE dt 2025-03-01;这里有个容易搞混的点HiveServer2的地址不是Metastore地址Catalog连接的是Hive Metastore的Thrift端口默认9083不是HS2的10000端口。把这两者搞错是接入失败的头号原因。如果你们的Hive集群启用了KerberosCatalog还需要额外配置hive.kerberos.principal、hive.kerberos.keytab和认证方式。另外Doris对Hive元数据是有缓存的Hive侧新增了分区Doris这边不会立刻看到需要执行REFRESH CATALOG hive_catalog;来刷新。2.3 FE与BE配置里最容易出问题的三处第一个是内存参数。BE的mem_limit默认是物理内存的80%对于混部机器这个比例偏激进我一般调到60%-70%留出给操作系统和监控组件的余量。FE的JVM堆默认8G如果元数据量极大或查询规划复杂建议加到16G。第二个是查询并发限制。FE的max_running_query默认100对于小团队够用但一旦多业务共用集群一个跑飞的query就可能拖垮所有查询。我会配合max_query_mem_limit做双保险宁可让大查询排队也不要因为一个query打光整个集群内存。第三个是BE的 compaction 相关配置。Doris的BE后台会持续做数据合并compaction小文件多的时候compaction压力大会挤占查询资源。建议把cumulative_compaction_check_interval_seconds设为一个可接受的值比如30秒同时给BE配独立的数据目录避免单盘I/O瓶颈。3. 三条数据通路怎么选外部Catalog、批量同步、实时写入3.1 三种整合方式的实际体验打通Hive Catalog之后你其实有不止一条路让Doris用上Hive的数据。我实践下来主流方案是三种。第一种是直接用外部Catalog直查Hive。这种方式零拷贝、不需要同步任务Doris通过Metastore拿到Schema执行时由BE直接扫描HDFS上的ORC或Parquet文件。它的优点显而易见——部署成本最低缺点也很真实查询性能完全取决于Hive表文件本身的质量。如果Hive表是Text格式、小文件一堆、字段类型又不规整Doris扫描起来同样吃力甚至比Spark快不了多少。第二种是批量同步。用SparkSQL或者DataX定期把Hive里加工好的结果表拉到Doris走StreamLoad导入。这是目前生产环境用得最稳的方案数据进入Doris后会按照Doris的存储格式重新组织索引、分区分桶、本地存储全部生效查询速度自然是最快的。代价是多了一套调度任务数据新鲜度只能做到小时级或分钟级。第三种是实时写入。通过Flink Doris Connector把Kafka里的实时数据直接写进Doris同时仍然保留Hive的离线链路。这是有实时数仓需求时的选择数据新鲜度可以到秒级但链路复杂度明显上升数据质量、乱序处理、join策略都要额外考虑。三种方式对比如下通路数据新鲜度查询性能维护成本适用场景Catalog直查Hive取决于Hive表中受文件质量影响低临时探查、文件质量好的表Spark/DataX批量同步小时级/分钟级高索引可生效中数仓结果表、指标表Flink实时写入秒级高高实时大屏、实时风控、实时指标3.2 选型判断标准新鲜度、数据质量、成本很多团队上来就问Doris怎么跟Hive对接其实应该先问我的数据多久更新一次、业务能不能等。如果报表是T1的那实时链路纯属给自己找麻烦如果业务方要求看到分钟级的实时订单变化那Catalog直查也满足不了因为Hive表本身不会分钟级更新。我的选型逻辑是这样的先看Hive表文件质量如果这表是别人随便写的、分区乱、格式杂直接Catalog直查一定会被坑再看数据量级和查询模式如果是业务方高频查询的固定明细表和指标表直接上批量同步如果有实时需求Flink实时和批量同步双通道同时跑Doris用Unique模型配合sequence列做增量merge。3.3 典型双层架构长什么样我现在维护的一个项目就是典型的双通道架构。离线链路是业务日志进KafkaFlink负责把原始数据落到Hive ODS层Hive做DWD到DWS的清洗加工DWS结果表通过Spark批量任务每小时同步到DorisDoris对外提供报表查询和即席分析。实时链路是Flink直接消费Kafka把关键指标实时写入Doris的Unique表供大屏和实时看板使用。这样一套下来Hive仍然是数仓的底座Doris变成了加速层两套引擎职责清晰。最忌讳的是把两边都当成全能选手——让Hive扛实时查询或者让Doris承载整个数仓的ETL都是架构上的错配。4. 从Hive往Doris搬数据的实操细节小文件治理、分区对齐与表模型4.1 Hive小文件问题不解决同步和查询两头吃亏Hive优化小文件这个词大家都不陌生但放到Hive与Doris整合的语境里小文件问题会被放大。同步任务读Hive表时一个小文件对应一个Map任务文件越多任务越多调度开销和HDFS NameNode压力一起涨数据同步进Doris之后小文件还会变成Doris底层的多个Tablet增加compaction负担。所以我在做同步前一定会先治理Hive侧的源表。做法比较朴素用INSERT OVERWRITE重写一遍目标分区同时用DISTRIBUTE BY指定分布列让数据按固定粒度落到指定数量的文件。比如INSERT OVERWRITE TABLE dws_order_daily PARTITION (dt 2025-03-01) SELECT /* REPARTITION(20) */ user_id, city_id, order_cnt, amount FROM dwd_order_wide WHERE dt 2025-03-01 DISTRIBUTE BY user_id;REPARTITION(20)或DISTRIBUTE BY可以控制最终文件数量避免出现一个分区几百个小文件的情况。同时可以打开Hive的自动合并参数hive.merge.mapred.filestrue、hive.merge.size.per.task128000000让MapReduce在输出阶段自动合并小于阈值的小文件。4.2 分区对齐与Doris建表模型选择同步之前Hive和Doris两边的分区定义必须对齐。Hive按dt分区Doris也按dt分区同步SQL里写WHERE dt ${date}只处理增量分区避免每次全量重导。Doris的分区名我习惯用p20250301这种纯数字加前缀格式不带连字符省得某些SQL里引号问题。Doris建表模型这一块是整合方案里最关键的决定之一。Doris有三种表模型Duplicate模型明细存储不去重不聚合适合原样保留Hive明细数据。Aggregate模型预聚合存储适合指标同步SUM、MAX、MIN、REPLACE等聚合方式在建表时定死。Unique模型主键唯一适合实时增量场景靠主键做update。从Hive同步结果表时如果只是报表查询不需要更新用Aggregate模型最省事查询时聚合结果直接读连GROUP BY都省了。比如我要同步一张用户每日订单汇总表CREATE TABLE dws_user_order_daily ( user_id LARGEINT NOT NULL, dt DATEV2 NOT NULL, order_cnt BIGINT SUM DEFAULT 0, order_amount DECIMAL(20, 6) SUM DEFAULT 0 ) AGGREGATE KEY(user_id, dt) DISTRIBUTED BY HASH(user_id) BUCKETS 16 PROPERTIES ( replication_num 2 );如果只是同步Hive的明细流水后续可能还有更新或删除就选Unique模型用sequence列解决乱序覆盖问题CREATE TABLE dwd_order_detail ( order_id BIGINT NOT NULL, dt DATEV2 NOT NULL, user_id LARGEINT, order_status INT, update_time DATETIME ) UNIQUE KEY(order_id) DISTRIBUTED BY HASH(order_id) BUCKETS 32 PROPERTIES ( replication_num 2, function_column.sequence_type DATETIME );4.3 StreamLoad写入label、两阶段提交与并发控制数据从Hive同步到Doris底层走的是StreamLoad导入。StreamLoad是Doris提供的批量导入接口支持HTTP方式提交。它的工作方式是这样的客户端把数据流式推给BEBE边接收边写入最后返回导入结果。小数据量可以直接用curl测试curl --location-trusted -u admin:your_password \ -H label:sync_dws_user_order_daily_20250301 \ -H column_separator:| \ -T /data/sync/dws_user_order_daily_20250301.csv \ http://fe_host:8030/api/dws/user_order_daily/_stream_load生产环境量大时我一般用Spark配合Doris的Spark Connector或者直接写StreamLoad客户端。有几个经验值得记下来每一批导入都要设置labellabel是幂等标识。同一批数据如果因为网络问题重试label不变Doris会返回AlreadyExist不会重复导入。这是防止数据翻倍的第一道防线。大批量数据用两阶段提交enable_two_phase_committrue先预提交等所有数据都成功后再COMMIT。如果中途失败就ABORT避免看到半个分区的脏数据。并发数要控制。同步任务开的并发太高BE的写入压力会很大反而拖慢整体速度甚至触发compaction拥堵。我通常把总并发控制在BE数量的2-4倍写入速度用max_filter_ratio0来兜底——严格模式下宁可导入失败也不能静默丢掉脏数据。5. 同步之外的另一半工程Hive侧表结构配合与聚合下推5.1 字段类型和文件格式决定外部数据源能下推多少把数据同步进Doris之后Hive侧的表结构看起来就不那么重要了但如果还用Catalog直查或者同步任务本身要解析Hive文件字段类型的小问题会变成大麻烦。先说日期。Hive里最常见的反模式是用string存日期比如dt2025-03-01虽然看起来没错但Doris读外部数据时需要通过字符串解析出日期转换成本很高还妨碍分区裁剪。好的做法是Hive侧就用date类型或者至少在Doris建表时显式用DATEV2并保证同步时能正确映射。然后是decimal和char。Hive的decimal精度如果定义得比较随意比如decimal(20,10)同步进Doris时如果精度不匹配会出现loss precision的报错或者数据被截断。我的做法是在Hive ETL层就把金额、比率这类字段统一规范成固定精度Doris侧DECIMAL(20,6)或DECIMAL(27,9)两边对齐再同步。文件格式建议统一用ORC或Parquet加Snappy/ZSTD压缩不要用TextFile。TextFile在Doris Catalog直查时扫描效率很低ZoneMap下推也发挥不出来ORC/Parquet自带统计信息Doris扫描时能跳过大量无关stripes查询快很多。5.2 分区裁剪、统计信息与CBO的关系无论走Catalog直查还是批量同步后查询写SQL时都得有分区裁剪意识。Hive表是分区表你没写分区条件Doris只能全表扫描写对分区条件扫描量可能只剩几十分之一。这个收益比任何引擎优化都来得直接。Doris 2.x的优化器已经比较成熟CBO会根据统计信息决定表连接的执行顺序和方式。但它对Hive Catalog外部表的统计信息掌握是有限的做复杂JOIN时规划不一定最优。我的处理是核心报表数据一定要同步成Doris内表让优化器拿到准确统计信息对外部表的即席查询尽量控制表连接的数量和过滤条件。另外记得给Doris内表定期执行ANALYZE TABLE让统计信息保持新鲜。我见过一个案例Doris内表数据量翻了几倍但统计信息没更新CBO选择了错误的Hash Join策略查询从2秒退化到30秒跑一次ANALYZE就恢复了。5.3 Hive UDAF的复杂逻辑留在Hive算还是算完再入Doris查热搜词时看到有人在问Hive自定义UDAF函数这个跟我们的整合方案有直接关系。Doris的内置函数很多比如approx_count_distinct、percentile、窗口函数都支持但它在自定义UDAF上的生态远不如Hive丰富。你很难把Hive里沉淀多年的业务口径UDAF原封不动搬到Doris。我的经验是复杂业务口径的UDAF留在Hive侧计算算完的结果同步进DorisDoris只做简单聚合。换句话说Doris处理的是已经算好口径的指标而不是从原始明细重新推导指标。这样做的好处是口径只在Hive一处定义Doris不参与业务逻辑两边不会因为口径不一致吵架。反过来如果Doris侧确实需要高效去重统计可以用Bitmap类型配合BITMAP_UNION做预聚合把原始ID去重后物化在Doris表里查询时直接读取预计算结果比每次跑COUNT(DISTINCT)快得多。6. 上线后我踩过的几个坑从Flink写Hive到Doris查询报错6.1 Flink写Hive数据查不到不是没写入是没commit做整合方案时很多人会顺手用Flink把数据写到Hive作为实时转离线的一条链路。然后就会遇到热搜词里那个经典问题Flink sink Hive表数据不入表。我第一次遇到时也懵了。Flink任务明明显示成功Hive表里却查不到任何数据。后来查了Flink Hive Streaming Sink的机制才明白Flink写Hive默认是事务性的数据写完后要以_COPYING_后缀暂存在HDFS目录里只有Checkpoint完成才会触发事务提交把_COPYING_文件rename成正式文件。如果你没开Checkpoint或者Checkpoint总是失败那数据就一直处于半提交状态Hive自然查不到。解决办法很直接StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 开启checkpointHive Streaming Sink依赖checkpoint触发提交 env.enableCheckpointing(Time.milliseconds(60000));同时记得配置分区提交策略比如sink.partition-commit.policy.kindsuccess-file确保分区数据完整后才可见。排查时也可以直接看HDFS目录找.part开头和_COPYING_后缀的文件就能判断数据到底写没写进去。这个坑背后其实引出一个判断Flink直接写Hive本身适合做离线数仓的ODS层但不适合做需要秒级查询的实时服务。所以整合方案里实时数据我都是Flink同时写Hive留底和Doris服务查询而不是试图让Hive承担实时读取。6.2 Catalog查询报missing的定位思路热搜词里有个presto doris错误的missing跟我遇到的Doris外部数据源报错很类似。现象是用Doris查询Hive Catalog表时报出类似missing xxx column或missing partition的信息。这类问题八成出在元数据不一致上我按下面几步排查Hive侧的表结构是不是刚改过Doris对Catalog元数据是有缓存的加列、改类型后没有刷新查询就会按旧Schema去定位数据出现missing。执行REFRESH CATALOG hive_catalog;是最快的验证手段。列名大小写是否一致Doris默认对表名、列名做了小写归一化如果Hive侧的列名带大写两侧匹配不上也会报missing。我的经验是Hive建表时就把所有列名统一成小写。Hive表的SerDe类型Doris是否支持某些自定义JSON SerDe的表Doris读不了报错信息还特别隐晦。这种表我一般先在Hive侧落地成ORC格式的中间表再被Doris读取。排查时先开FE的审计日志把查询的实际SQL和涉及的表名、分区信息打出来基本能定位到是Schema问题还是权限问题。权限问题的话留意Ranger / Kerberos的映射配置。6.3 StreamLoad重复执行数据翻倍我见过最疼的一个坑就是同步任务重跑导致Doris里数据翻倍。表象是Doris里的订单金额比Hive大了一截查下来发现同步脚本跑了两遍。原因基本就是label没用好。StreamLoad的幂等靠label实现如果每跑一次都生成一个新label那么同一批数据就能被导入两次。我的规范是所有同步脚本的label都带上业务名和日期比如sync_dws_user_order_daily_20250301同一分区同一逻辑批次的重复执行label保持一致Doris会自动返回AlreadyExist直接跳过重复导入。另外用两阶段提交时事务没有COMMIT之前数据是不可见的如果脚本在COMMIT之前退出需要执行ABORT清理事务否则会一直占用导入资源影响后续同步。这两点配合好了同步任务怎么重跑都不会污染数据。6.4 BE内存与查询并发的平衡问题MPP引擎的快本质是拿内存换时间。Doris集群用久了容易遇到一个现象某个业务方跑了个超大JOIN或者全表无过滤查询BE内存瞬间飙升同一时间其他所有查询全部变慢甚至失败。控制手段我在前面部署部分提过这里再补充几个实战经验。BE的mem_limit别给满我常用的是物理内存的60%-70%剩下给操作系统做Page Cache反而对扫描性能有帮助FE的max_query_mem_limit要设超限查询直接拒绝或排队对于实在要跑的复杂大查询开BE的Spill数据溢写到磁盘能力不要让一个查询把集群内存打穿。扩容的时候也要有预期新加BE节点后旧节点上的Tablet不会立刻均衡过去Doris后台按批次迁移整个过程可能持续几小时甚至更久视数据量而定。所以扩容尽量安排在业务低峰期扩容后再关注BE之间的数据均衡度。BUCKETS数量设得不合理后续调整成本很高建表时按总数据量除以单个Tablet 2-5GB的规模来定别拍脑袋。回到最开始的问题Hive与Doris整合这件事技术选型并不难难的是把链路里的每个细节都想清楚。我个人的体会是先别急着建一堆同步任务先用Catalog直查把Hive表摸清楚哪些表是业务高频查询的、哪些文件质量差再针对性地设计同步方案。同步链路一定要做好label幂等和监控告警宁可任务跑慢一点也不要重复跑或者静默丢数据。Doris的查询确实快但它不是银弹它把Hive侧的文件质量问题、Schema规范问题都提前暴露了出来逼着你把数仓基础打扎实。配合数据分层和两套引擎的合理分工这套架构跑起来之后你会发现业务问的为什么这么慢慢慢变成了能不能再加几个报表。
返回列表