ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

CDH与ES统一迁移:基于HDFS+Iceberg+StarRocks的大数据底座实践

2026/10/7 3:41:11 拓冰建站 浏览量
CDH与ES统一迁移:基于HDFS+Iceberg+StarRocks的大数据底座实践 2024年帮一家城商行做了件“减法”把运行多年的CDH集群连同散落在各业务部门手里的Elasticsearch集群陆续请下了生产环境的C位换成了一套以HDFS、Iceberg、StarRocks为核心的企业级大数据底座。说是减法是因为原来技术栈太杂——CDH管批处理ES管检索和明细查询两套体系各存各的数据口径经常打架做完之后数据只存一份批、流、查都在同一条链路里跑。这篇不是理论宣讲想把整个案例里最关键的判断、参数、迁移步骤和踩坑细节整理出来。如果你也在纠结处理CDH改造、收敛多套ES或者正琢磨怎么把大数据平台拉成一套统一底座这篇应该正好对路。1. 为什么下决心动这套组合CDH和ES的“三座大山”1.1 看似最稳的CDH维护成本到底有多高先说背景。这家城商行的数据规模没有互联网大厂那么夸张原始加汇总也就200TB上下但业务链路一点都不简单每天凌晨核心系统、信贷、手机银行的数据要落到贴源层接着跑清洗、汇总、指标加工早上七点前必须出前一日的报表。原来这套全靠CDH撑着Hive是绝对主力偶尔用SparkSQL救火。前几年还算凑合后来越来越吃力。凌晨开始的批处理经常跑到七点半应用运维的电话直接打到我这里。加节点、加并发都试过瓶颈出在Hive自己身上小文件多、执行计划优化弱单靠加机器根本解决不了。加上Cloudera的商业模式调整之后license费用一年比一年高每次升级都像在走独木桥。更要命的是整个平台维护就靠三个人既要看HDFS、又要盯YARN、还要管Hive元数据出一次任务失败恢复就能耗掉大半天实在耗不起。这里补充一个细节CDH不是不好而是对一个IT团队只有十几个人的城商行来说太重了。它的定位是“全家桶式的企业级发行版”安装、升级、授权、排障都依赖原厂支持。一旦不想花这笔钱团队就得把整个生态的运维知识补齐代价非常高。我们当时的判断很直接既然Hive在海量小文件场景下的性能瓶颈短期内靠调优解决不了那不如换一套更适合“中等数据量、复杂跑批”的引擎底座把维护复杂度打下来。1.2 Elasticsearch集群越建越多数据孤岛越来越重ES这边的问题更让人头疼。风控、运营、渠道几个部门各自搭过Elasticsearch集群有的是为了查交易明细有的是为了看应用日志还有的干脆拿ES当宽表引擎来跑报表。结果就是同一份数据源被同步了好几遍同样一笔流水在三个集群里各存一份存储膨胀比数仓还严重——本地磁盘从4T加到12T都还是不够。ES对内存又特别敏感。JVM堆开小了容易频繁FullGC开大了留给操作系统的页缓存不够索引分片、脑裂、字段映射冲突这些事几乎每个月都要处理一次。其中最难忘的是有一次索引损坏折腾了一整个通宵去做数据恢复最后还是丢了一部分近端数据从上游重新同步才补上。这种隐性维护成本在选型的时候根本不会写在供应商的PPT里。还有一个更致命的问题ES和Hive数仓对不上账。数仓里算出的交易量和ES检索出来的数量经常差几个百分点业务部门来问的时候谁也说不清楚哪套是准的。后来核查才发现两边同步链路用的字段过滤条件不一样消费Kafka的位点也不一致。数据链路分裂久了连“正确答案”都找不到这可比跑批慢一两个小时严重多了。1.3 统一底座到底要解决什么问题所以领导的诉求其实很简单能不能把平台收敛成一套这个“一套”不是指把所有组件物理装进一个集群而是满足三个约束数据只存一份、链路统一查询和分析在一个引擎里跑权限、监控、审计都在同一套体系里完成。这就等于明确了两件事CDH里的Hive要被替换ES里的检索和明细查询场景也要想办法迁移到新的统一引擎上。范围听起来挺大但拆开之后是有次序的先定存储底座再定计算引擎最后处理存量迁移和双跑验证。这里最忌讳的就是“重构式思维”老想着把组件一锅端换掉实际做法是按场景收敛一步步把流量导到新平台上。2. 技术选型不追求时髦但要替得动2.1 先拆场景再定技术栈技术选型不是拍脑袋而是先把现有场景摆到桌面上。我们列了一张场景矩阵把原本CDH和ES承载的东西逐条对号入座场景原技术栈核心痛点新方案离线批处理/数仓ETLCDH Hive/Spark跑批慢、小文件多、License重HDFS Iceberg Spark / StarRocksOLAP多维分析报表Hive Elasticsearch口径不一致、存储膨胀StarRocks明细查询/分页检索Elasticsearch深分页、内存瓶颈、多集群重复StarRocks日志关键字检索Elasticsearch各团队自建、数据重复、恢复难StarRocks 日志结构化 冷热分层实时/准实时入湖手工同步/DB同步链路长、延迟高Flink CDC / Stream Load 入 Iceberg这张表一出来选型方向就很清楚了不是找一个组件去等价替换ES而是把ES的职责拆掉——一部分交给OLAP引擎一部分交给数据湖表格式一部分直接消灭重复存储。2.2 为什么选“HDFS Iceberg StarRocks”这套组合当时也对比了不少方案。HDFS作为分布式存储继续保留是因为它成熟稳定换掉它等于重新趟一遍文件系统级别的坑没必要。真正需要动的是“表”这一层原来Hive把元数据和数据文件绑得太紧ACID能力弱小文件治理全靠人工所以我们把表格式换成了Iceberg让数据湖层具备快照、时间旅行和自研化的文件组织形式。查询引擎最终选了StarRocks理由就三条。第一它支持完整的SQL语法业务部门原来写SQL的习惯不用改从ES迁过来的查询也能基本平移。第二StarRocks的索引能力比一般MPP引擎强既有ngram_bf索引也有倒排索引日志检索里一部分关键词搜索场景可以无缝接住。第三它的External Catalog可以直接查Hive和Iceberg上的表这意味着新旧切换不用“先搬迁再查询”新平台可以先上线旧数据逐步归档。对比下来Doris和StarRocks同源很多能力接近但StarRocks在索引、外部目录、资源组这几个维度上更贴合我们的诉求。ClickHouse则更适合单表极速分析多表关联和SQL完整性在复杂报表场景下还不够顺手。这里没有谁一定比谁强关键是匹配自己的场景结构。另一个选择是继续用OpenSearch这类“类ES”产品但我们最终没有考虑因为那只是在检索域换了个名字数据孤岛、多套集群、口径分裂的问题一点没减反而又多了一根支柱要维护。2.3 架构形态和边界怎么划新底座的逻辑分层很清晰最下层是统一的HDFS存储中间是Iceberg做表格式和元数据管理再往上是StarRocks统一承接查询、报表、检索和一部分ETL加工。数据接入层用Flink CDC和Stream Load把源端关系库、日志、消息队列的数据汇到Iceberg或直写StarRocks。边界也提前划好StarRocks不是用来存全量冷数据的它主要负责“高并发查询和近热数据”历史归档数据放在Iceberg上需要的时候再通过外部目录或导入方式取用。这样做的好处是既保住了查询性能又不让OLAP集群被数据量拖垮。城商行这种数据规模PB以下这样的分层既简单又够用不搞湖仓两套大集群运维负担可控。3. 从规划到落地底座迁移实操全流程3.1 容量评估先算账再采购这一步特别重要建议所有准备替换CDH或整合ES的团队都把账先算清楚不要凭感觉定节点数。我们当时按日志保留90天、数仓数据累积一年来估算日志日新增约0.5TB压缩后按0.4折算每天净增约0.2TB数仓和明细数据日新增约0.3TB压缩后约0.12TB。加一起每天新增约0.32TB一年新增约116TB。现有存量数据约150TB整体数据规模约266TB按HDFS三副本计算裸存储需求约800TB。如果按12个数据节点、每节点配12块4T盘来估算单台可用容量约44TB12台合计约528TB覆盖中等增长没问题但为了留出缩容余量最终第一批上了14个数据节点后续按需横向扩容。硬件选型上管理节点用了3台承担NameNode、JournalNode和StarRocks FE数据节点就是计算和存储一体CPU选了32核、内存128G磁盘用HDD加一块SSD做热数据缓存。城商行预算有限没必要追求全闪存StarRocks本身有本地缓存机制SSD缓存可以显著提升热点查询的命中率。3.2 新底座搭建和参数调优搭建过程本身不复杂顺序是先把HDFS拉起来再部署StarRocks最后初始化Iceberg Catalog。我们直接用已有的Hive Metastore作为Iceberg Catalog的元数据后端省掉再引入一套独立元数据服务的成本。真正花时间的是StarRocks的配置。这里给一个我们最终稳定运行的参数参考基于128G内存的BE节点mem_limit设为80Gstorage_page_cache_limit设为50G。这个比例很关键——页缓存开太大查询执行内存不够高并发时容易OOM开太小缓存命中率上不去查询RT波动明显。我踩过一版把mem_limit拉到100G的坑双跑验证期间每天晚高峰必现BE内存溢出后来把页缓存降到50G、给查询执行和compaction留足空间问题才消停。建表也有一套约定。明细大表用Duplicate Key模型加日期分区、客户号分桶汇总指标表用Aggregate模型需要频繁更新的维度表用Primary Key模型。举个明细表的例子CREATE TABLE dwd_trade_detail ( trade_no BIGINT, customer_no BIGINT, channel_id INT, amount DECIMAL(15,2), biz_time DATETIME ) ENGINEOLAP DUPLICATE KEY(trade_no) PARTITION BY RANGE(biz_time) () DISTRIBUTED BY HASH(customer_no) BUCKETS 48 PROPERTIES ( replication_num 3, storage_medium SSD, storage_cooldown_time 2025-01-01 00:00:00 );分桶数不是越大越好我们按单桶数据量在200MB到500MB之间来定客户号分桶48跑批和明细查询的并行度都够用。3.3 数仓迁移从Hive到Iceberg怎么改才不翻车Hive表迁到Iceberg有两条路。一种是用Spark SQL的CREATE TABLE ... AS SELECT直接重建数据另一种是通过设置table_type属性把原表原地升级。我们多数表用了第一方案因为可以顺手重排分区、清掉小文件、统一列类型。迁移语句很直接关键在迁移前的表结构梳理CREATE TABLE iceberg_db.dwd_trade_detail USING iceberg PARTITIONED BY (dt) AS SELECT trade_no, customer_no, channel_id, amount, biz_time, dt FROM hive_db.dwd_trade_detail_old;迁移过程中的第一个大坑是类型映射。Hive里的varchar、char在Iceberg里会统一归成string但如果源表里有些字段是decimal(38,18)、有些是decimal(10,2)在CTAS时不提前统一下游SparkSQL查询会经常报类型转换错误。我们提前写了一个脚本扫描所有表的字段类型把跨层不一致的地方全部改成统一标准再启动迁移。第二个坑是分区字段策略。原来Hive表习惯用dt string存“2024-06-01”这种字符串Iceberg支持转换后的date类型但考虑到现有调度代码大量依赖dt‘${date}’这种传参方式我们最终没有改成date保留了string分区字段。技术上不是最优但迁移代价最小事实证明这个妥协很值改一处调度核心逻辑的成本比改几百个任务高得多。第三个坑是小文件。Iceberg本身有compaction机制但不会自动跑到完美状态。生产上我们每分钟可能有几十个流式任务产小文件一开始快照数量膨胀得非常快元数据查询都要好几秒。后来设置了定时compaction同时给快照保留设了7天有效期定期执行expire_snapshots把过期快照清掉元数据才算稳定下来。3.4 检索和查询场景迁移从ES到StarRocksES迁移到StarRocks最怕的就是“拿ES的思维用StarRocks”。ES里的text类型全文检索、嵌套文档、from/size深分页这些如果硬搬过去性能会很差。我们的做法是先把场景拆成三类第一类是纯结构化查询比如“某客户某天的交易明细”、“某渠道近一月交易金额”。这类直接转成SQL过滤用分桶和分区裁剪就能跑得很好比ES的query_string方案快且省内存。第二类是关键字搜索比如查日志里某条报错关键字。StarRocks支持ngram_bf索引我们在需要检索的字段上建了ngram索引把原来的match查询改写成LIKE加索引条件。怎么说呢体验上跟ES的match有点差异但日常日志关键字搜索基本够用。真正做复杂语义检索的场景城商行也不多没必要硬上搜索引擎那一套。第三类是深分页。ES的from/size深分页翻到后面性能断崖式下降StarRocks同样不推荐跳过深分页。我们把所有明细查询都改成了游标式翻页按排序键做增量偏移。样例SQL是这样的-- 旧ESfrom10000, size100 -- 新SQL基于排序键游标翻页 SELECT trade_no, amount, biz_time FROM dwd_trade_detail WHERE biz_time 2024-06-01 00:00:00 AND (biz_time, trade_no) (2024-06-01 12:30:00, 1024001) ORDER BY biz_time, trade_no LIMIT 100;字段映射上ES的keyword对应StarRocks的STRINGES的long对应BIGINTdate对应DATETIME。ES里如果是数组类型StarRocks没有原生数组索引支持我们统一做了拆行或者JSON解析避免到了新平台才发现字段类型根本对不上。ES原有的数据同步链路也要改。之前是Logstash一把梭现在改成日志先入Kafka再通过Flink做清洗落地Iceberg做冷归档近30天热数据直写StarRocks。对账层面每天定时跑count和sum核对确保两边指标一致。我们当时定了一个硬规则数据不同步超过万分之一就要阻断上线不放过任何一条链路。3.5 双跑验证与灰度切流新旧平台并行跑了大概一个月这个周期不能短至少要覆盖一个完整的结息周期和月度报表周期。双跑期间的验证分三层做。第一层是数据量核对。每天定时用SQL做源端和新端行数、金额sum、最大最小值、count distinct的比对例如SELECT COUNT(*) AS cnt_new, SUM(amount) AS amt_new, COUNT(DISTINCT customer_no) AS cust_new FROM starrocks_db.dwd_trade_detail WHERE dt 2024-06-01;第二层是抽样明细比对。按客户号和日期分层抽样5万条逐字段比较任何字段不一致都自动生成差异报告。第三层是报表结果比对拿核心指标交易量、活跃客户数、余额汇总做日级核对允许0.01%以内的容差因为两边数据加工时点不同存在跨秒差异是正常的但要能解释差异来源。切流顺序我们分了四批先切纯查询类场景出了事影响面小回滚也容易再切准实时指标和报表查询然后切日志检索因为终端用户的使用习惯改变大留足培训时间最后才把离线ETL和批处理彻底从Hive切走。整个切流期间旧集群保持数据同步运行了30天直到新平台稳定度超过99.9%才真正按下旧的关机键。4. 踩坑复盘那些文档里查不到的细节4.1 常见问题速查表迁移过程中有一堆问题整理成一个速查表算是给后来人避雷现象排查方向解决办法StarRocks高峰期BE OOMmem_limit/页缓存配置过高调低mem_limit限制查询并发给compaction留内存查询偶发超时、RT跳动大scan rows过高、缓存命中率低查看慢查询日志调整分桶键增加SSD缓存比例Iceberg快照数量膨胀未定时执行expire_snapshots配置定时compaction和快照过期策略保留最近7天Hive迁移后查询报类型错误varchar/decimal映射不一致迁移前统一全表字段类型规范写入自动扫描脚本ES数组字段迁移失败StarRocks不支持原生数组索引拆行或JSON解析迁移前做字段类型清洗新旧平台数据对不上ES的doc_count包含嵌套文档改为按_source级或明细行级核对不能信doc_countHive原任务依赖大量自定义UDF新引擎不支持旧UDF优先用SQL改写改写不了再注册等价UDF到新引擎调度任务并发高峰CPU跑满StarRocks资源组未配置按业务线划分资源组给核心报表保证配额4.2 内存配置的教训StarRocks的BE内存配置是上线初期最大的坑。我印象很深第一次双跑时查询并发一起来BE直接OOM所有该节点的查询全部失败。查了半天发现不是SQL问题就是mem_limit设了100G把内存全塞给了页缓存查询执行时只剩20G不到并发一高必然爆。后来参考官方推荐做了调整总内存的60%左右给page cache剩下40%给查询执行、compaction、导入等动态内存。128G的机器最终稳定在mem_limit 80G、storage_page_cache_limit 50G。打个比方内存就像办公室页缓存是茶水间查询任务是工位茶水间占太多工位就不够坐整个办公室都转不动。内存以外磁盘也值得说。HDFS数据节点和StarRocks BE不建议共用数据目录因为Iceberg的compaction和StarRocks的compaction同时打到同一块盘上IO争抢会让查询RT迅速恶化。我们后来把SSD单独留给了StarRocks热数据缓存机械盘只放HDFS冷数据两类compaction互不干扰。4.3 权限管控和行级安全怎么做城商行的数据安全要求是顶格的不是连上数据库就能看全量。ES时代各集群各管各的账号要做一次跨库权限审计简直噩梦。新底座把权限统一到一套体系里行级和列级都实现了。行级权限主要应对“客户经理只能看自己名下客户”这类需求。StarRocks 3.x版本自带行列级权限能力低版本就用视图包一层CREATE VIEW v_customer AS SELECT * FROM customer WHERE manager_id current_user_id()把视图授权给对应角色底层表只开放给管理员。列级权限就更直接敏感字段手机号、身份证号只允许指定角色查询其他角色默认不可见。这样一来新增一个人力报表账号只需要建角色、挂权限两个动作不用再像以前那样挨个ES索引配映射。审计需求也要提前设计。所有查询会记录到审计日志表定期归档到Iceberg保留180天。这个设计在当时看多花了点时间后来审计抽查时派上了大用场几秒钟就能定位谁在什么时间查过哪些敏感字段。4.4 监控体系怎么收敛原来CDH一套监控ES一套监控之间没有任何联动。新底座直接统一到Prometheus加Grafana所有组件暴露同一套指标。重点盯这几类HDFS的NameNode RPC延迟、坏块数、DataNode磁盘水位StarRocks的查询RT分位数、BE compaction score、Stream Load失败率Iceberg的快照数量、小文件数量、commit延迟。告警阈值我们磨了几轮才稳定。比如BE compaction score超过一定水位持续15分钟才告警避免短时波动骚扰值班查询P95超过3秒持续10分钟才触发低于这个阈值不打扰。监控收敛完之后现在值班看板就是一块屏从底层存储到查询引擎的状态全在里面再也不用像以前那样切四五个系统反复比对。5. 一些个人体会和留给后来者的彩蛋整个项目做下来我最大的体会是替换CDH和ES这类组合关键不在“新组件比旧组件性能好多少”而在“一套底座到底能不能覆盖80%以上的常见场景”。城商行的数据规模和技术团队配置决定了它承受不起多套技术栈长期并存每一步选型都应该让自己离“统一”更近而不是越换越杂。最后分享一个我觉得特别实用的小技巧StarRocks的External Catalog可以直接挂载旧的Hive数据源。也就是说迁移期根本不需要“先搬数、再验证”新平台上线后查询引擎能直接跨catalog同时访问新表和旧表。这个能力让我们的双跑验证从“数据同步式双跑”变成了“同引擎跨源对比”效率高了很多也帮我们大胆地把切流周期从三个月压缩到一个月。如果你也在做类似的数据平台收敛建议先把这个能力用起来迁移压力会小非常多。