关于我们

质量为本、客户为根、勇于拼搏、务实创新

< 返回新闻公共列表

2026 Hive 跑批越跑越慢:先别急着加 CPU,问题多半在小文件上

发布时间:2026-09-28

先别急着加 CPU

跑批慢了,第一反应通常是"机器不够了"。但把作业时间摊开来看,真正花在读数据和算数上的那一段,往往连一半都不到,剩下的全交给了分区 listing、文件 open、块信息拉取、任务申请和 JVM 启动。加 CPU 只能加速那不到一半的部分,剩下的一大半,你加多少核它都不会变快。

来看一个假设场景(以下为典型部署思路,并非特指某一真实客户)。某公司的离线数仓有一张按小时分区的行为明细表,每天凌晨跑一次汇总,产出几张报表宽表。上线第一年,作业跑 40 分钟。第二年同样的代码、同样的数据产出类型,跑 1 小时 40 分。第三年变成 4 小时,已经压到早上 8 点还没出数,业务一上班就开始催。

中间他们做过两件事:一次是给计算队列加了资源,把可用核数提上去;一次是把单节点从 16 核换成了 32 核。两次都是加完之后头两天快了一点,一周之后又回到原样。原因不复杂——作业变慢的那部分时间,跟 CPU 核数基本没关系。

这篇文章要讲的就一句话:跑批变慢的根因多数不在计算资源,而在「文件数量 + 元数据规模」——每个小文件都要一次独立的 map 任务与一次元数据访问;先治理小文件和分区粒度,再谈加机器,否则加多少核都是把时间花在调度和打开关闭文件上。

在往下看之前,先把这几条当成全文的骨架:

一次跑批的时间可以拆成六段,其中"提交编译""元数据 listing""切片与任务划分""任务调度启动"这四段的时间开销主要由文件数和分区数决定,跟 CPU 无关。

文件数比数据量致命。一个 128MB 的文件和一个 1MB 的文件,调度开销几乎一样;1000 个 1MB 文件的总数据量不到 1GB,但要启 1000 个任务。

"明明没在算,时间却在涨",涨的通常就是元数据那一层:分区数增长让 metastore 的 listing 变慢,文件数增长让元数据节点的堆内存和 RPC 压力上升。

小文件不是某一次写错造成的,是五种常见写入习惯日积月累攒出来的:分区粒度太细、写入并发过高、微批太频繁、补数重复写、覆盖写不做清理。

等到"真正计算"那一段时间占比已经很高、CPU 和内存持续打满、横向扩容后时间近似线性下降,这时候加机器才有意义。

把一次跑批的时间拆开:算数的时间占多少

要判断慢在哪,先把一次跑批从提交到结束完整拆成六段。这六段的成本驱动力完全不同,混在一起看就只能得出"慢了"这个结论。

第一段:提交与编译

客户端把 SQL 提交上去,引擎做词法语法解析、语义校验(表存不存在、字段对不对、权限够不够)、生成逻辑计划、过优化器、生成物理计划。这一段跟数据量完全无关,跟 SQL 复杂度强相关。一条几十行、带七八个 JOIN 和一堆子查询的 SQL,编译几十秒到几分钟都不稀奇。表分区特别多的时候,这一段也可能被拖住,因为优化器做分区裁剪要去元数据里取分区键的取值范围。

这一段的特点是:它每次都要重跑一次。作业跑得再快,这段时间是净支出。对于频繁提交的小作业,把常用逻辑固化成预编译视图或者合并多个小 SQL 成一个大 SQL,收益比优化数据处理本身更直接。

第二段:元数据 listing

这一段是本文的主角。引擎要拿到"这个表的哪些分区符合 WHERE 条件",然后对每个命中的分区,列出它下面的所有文件,拿到文件大小、块信息、副本位置。注意这是两轮:先取分区列表,再逐个分区取文件列表。按天分区、查一年数据,第一轮就是 365 个分区;每个分区下如果有 200 个文件,第二轮就要拉 7 万多条文件记录。

这段时间是纯延迟,没有任何计算在里面。它在作业监控界面上往往显示不出来——很多界面把它归到"提交"或者干脆不计,看到的就是作业"卡住了"几十秒甚至几分钟。

第三段:切片与任务划分

拿到文件列表之后,引擎按文件大小切成 InputSplit。规则很直接:单个文件如果小于块大小,通常独占一个 split(可合并的输入格式除外);大于块大小的文件按块切开。所以 split 数量 ≈ 有效文件数量,而不是总数据量除以块大小。这一步确定了后面要起多少个任务,也是"文件数 → 任务数"这条链路的关键环节。

第四段:任务调度与启动

每个 split 要申请一个 container:资源管理器排队、分配到节点、下载依赖、启动 JVM(MapReduce 模式下)、加载类。这一套流程的固定成本在几百毫秒到几秒之间浮动,取决于集群繁忙程度和依赖包大小。集群空闲时快,夜间几十个作业一起跑时慢,还可能出现"资源等待"把墙钟时间再拉长一截。

这四段加起来,在很多长期没做治理的老集群上能占到总时长的六成以上。剩下的两段时间才真正跟 CPU 有关。

第五段:数据读取与计算

这一段时间跟数据量成正比,也确实跟 CPU、内存、磁盘吞吐有关:解压列式文件、反序列化、过滤、JOIN、聚合、排序。这是我们"感觉上"应该在花时间的地方,也是加 CPU 唯一能加速的一段。注意排序和 JOIN 阶段还要 spill 到磁盘,这部分又跟磁盘 IO 有关,跟 CPU 关系不大。

第六段:落盘与提交

结果先写到临时目录,再由 driver 侧做 rename 提交,同时更新元数据里的分区信息和统计信息。动态分区写入时,这一步要为每个新分区写元数据记录;落盘的分区数越多,提交越慢,而且这一步是单点的(driver 串行),加再多计算节点也没用。还有一个容易被忽略的坑:driver 要为落盘的每个文件走一次 rename,文件数一多,这里又是几十秒到几分钟。

小文件是怎么攒出来的:五个常见来源

没有哪个团队是故意写小文件的。它们全都是"当时看很合理"的设计,在数据量涨了两个数量级之后才变成问题。按常见程度排:动态分区写入并发过高第一,分区粒度太细第二,微批写入第三,重跑补数第四,覆盖写不做清理第五。

来源一:分区粒度太细,按小时甚至按分钟

按小时分区听起来很自然:能按小时回溯、能快速删掉某小时的数据、还能天然支持小时级增量。问题在于它是乘数因子。一张按分钟分区的表,一天 1440 个分区,一年 52 万多个分区;每个分区哪怕只有 3 个文件,一年也是 150 万个文件。

更隐蔽的问题是"每个分区多少数据"。假设这张表每天总量 500GB,按小时分,每小时平均 20GB,这个还算健康;但如果每天总量只有 5GB,按小时分,每小时就 200MB,这时候哪怕每个小时只有 10 个文件,单文件也只有 20MB——离目标文件大小差了一个数量级。分区粒度必须跟单分区数据量匹配,不能跟时间戳的精细度匹配。判据很简单:单个分区的数据量应该显著大于"目标单文件大小 × 该分区期望的文件数",比如目标是每个文件 256MB、每个分区不超过 20 个文件,那单个分区至少要有 5GB 才撑得住。

来源二:动态分区写入时 reducer / task 数量过多

这是最常见也最容易被漏掉的一条。动态分区写入时,每个 reducer(或每个写 task)在它负责的每个分区键下都会生成一个或多个文件。表面上看你只是写一张表,实际落地的文件数是"分区数 × 写并发数"。

举个例子:某天的数据分布在 24 个小时分区里,写入作业有 200 个 reducer(很多引擎的 shuffle 分区默认值就是 200),理论上最多会产生 24 × 200 = 4800 个文件。实际会少一些,因为不是每个 reducer 都持有每个分区键的数据,但量级在这。如果一个 reducer 每个分区只有 5MB 数据,那落下来的就是 4800 个 5MB 的文件。

而且这个数还会随数据量变化:数据涨了,引擎自动把 reducer 数调大,文件数跟着涨,单文件大小却没怎么变。这就是为什么"数据量明明只涨了三倍,文件数涨了十倍"。

来源三:流式 / 微批写入

每分钟落一批的增量写入,一年下来是 52 万批。哪怕每批只写一个分区下的一个文件,一年也是 52 万个文件。如果每批要写多个分区(按业务 key 分桶或多个指标表),文件数还要再乘。

这类写入天然和"大文件"冲突:单个批次的时间窗口决定了数据量上限,一分钟的数据量再大也大不到哪去。要破这个局只有两条路:一是拉开批次间隔(一分钟改成五分钟或十分钟),二是在批次之外加独立的合并任务,定期滚动做一次分区内的文件合并。

来源四:上游重跑、补数时重复写入

上游表重算、补历史数据、修 bug 之后回填,这些操作本身没问题,问题在于写入方式。如果用 append 追加写而不先清掉目标分区,同一份数据会在同一个分区里留下两套文件:一套是旧版本,一套是新版本。查询时如果逻辑是按主键取最新,逻辑上没错,物理上文件数翻倍了。

还有一种更常见的:回填作业本身跑了三次(第一次失败、第二次数据不对、第三次才成),前两次留下的中间文件没有通过 staging 机制清理干净,静静躺在分区目录里。这种情况下 SELECT COUNT(*) 的结果可能是对的(如果有去重逻辑),但 listing 和 split 的数量已经翻倍。

来源五:insert overwrite 生成新文件但旧文件未清理

覆盖写在语义上是"替换整个表或分区",但这个"替换"依赖于实现是不是干净。几种常见情况会让旧文件留下来:用了不支持 overwrite 语义的存储格式或外部表配置;写到了 S3 这类对象存储上但清理逻辑因为请求失败而静默跳过;写入过程中 task 失败重试产生的残留文件没有回归一清理作业清掉。

还有一个大家心照不宣的情况:用 INSERT OVERWRITE 覆盖分区时,如果开了动态分区且某些分区键在这次写入里没有数据,那些分区不会被"清成空",而是保持原样。这本是合理设计,但如果下游以为被覆盖了,就会出现新旧数据混在一个分区里的诡异结果。

为什么 1000 个 1MB 比 1 个 1GB 慢得多

这一节把"文件数比数据量致命"拆成三笔账:调度账、打开账、尾巴账。

第一笔:调度账

前面说过,每个 split 对应一个任务。1 个 1GB 的文件按 128MB 块切成 8 块,起 8 个任务;1000 个 1MB 的文件,每个文件独占一个 split,起 1000 个任务。数据总量几乎一样(1000MB ≈ 0.98GB),任务数差了 125 倍。

每个任务有固定开销:申请资源排队、分配到节点、拉依赖包、启动容器、注册、执行、上报、清理。把这部分记为 T_task(典型量级在数百毫秒到数秒,视集群繁忙程度和引擎而定,非实测值)。1000 个任务就是 1000 × T_task。哪怕集群能并行 200 个容器,也要 5 轮,每轮都带着各自的尾巴。

更麻烦的是,这 1000 个任务的实际执行时间极短——读一个 1MB 的文件、解压缩、过滤、输出,可能几百毫秒就结束了。任务的实际工作时间远小于它的调度启动开销,这时候整个作业本质上在"为了打开文件而创建容器"。把 1000 个 1MB 合并成 8 个 128MB,调度开销从 1000 份降到 8 份,这部分时间几乎可以全额省下来。

第二笔:打开账

每个文件的读取不是"直接读"这么简单。客户端要先向元数据节点发起请求拿这个文件的块位置列表,然后跟对应的数据节点建立连接、协商、读取、关闭,读完还要处理校验。这套流程对于一个 128MB 的文件来说完全可以忽略——读几十秒的时间里有几十毫秒花在建立连接上。但对于 1MB 的文件,建立连接的开销可能跟读取本身一个量级。

在列式存储上还有一笔额外成本:ORC/Parquet 这类格式在读之前要先读 footer(文件尾部的元数据区),拿到 row index、stripe 分布、schema 和统计信息,才能决定读哪些 row group。哪怕最终一个 row group 都不读,这个 footer 也必须先拉。1000 个文件就是 1000 次 footer 读取,每次都要一次完整的 open→seek→read→close。

压缩也受影响。列式压缩依赖字典和数据局部性,一个文件内的数据越多样、越连续,压缩比越好。把本来可以放在一起的数据切碎到 1000 个文件里,每个文件单独建字典,压缩比会明显下降,反过来又让磁盘要读更多字节。

第三笔:尾巴账

任务数越多,慢节点(speculative/straggler)出现的概率越高,而作业的结束时间取决于最慢的那个任务。1000 个任务里有一台机器磁盘抖了一下、网络闪了一下,整作业就要多等几十秒。8 个任务的情况下遇到同样问题的概率低得多。

用输入合并格式能救一部分。CombineHiveInputFormat、CombineTextInputFormat、以及 Spark 侧的自适应查询执行里的分区合并(spark.sql.adaptive.coalescePartitions.enabled、spark.sql.adaptive.advisoryPartitionSizeInBytes)都是把多个小 split 并成一个大 split,直接砍掉任务数。但要注意:它们只能在已经拿到文件列表之后做合并,救不了第三段之前的开销。也就是说 listing 慢、metastore 压力大、driver 串行 rename 慢,这些它一个都管不了。它是止痛药,不是根治。

元数据那一层:没在算,时间也在涨

这一节解释那句让人最困惑的抱怨:"我什么都没改,数据量也没涨多少,为什么一天比一天慢。"答案通常在两个元数据服务上。

Hive Metastore:分区数增长带来的是雪崩式开销

Metastore 的后端是一张关系库,分区信息存在 PARTITIONS 表里,每个分区一行。一个上千分区的表,一次查询要从这张表里拉出上千行,再做过滤、排序、分页返回给客户端。这张表的行数增长会让本来几十毫秒的查询变成几秒。

更糟的是后面那一轮:对每个命中的分区,客户端(或 metastore 内部)要去拿这个分区下的文件列表。分区数 × 每分区文件数,这才是真实的 listing 规模。分区数从 365 涨到 8760(按天改按小时),如果每个分区的文件数不变,listing 粒度直接翻了 24 倍。

还有一个独立性很强的问题:统计信息。很多集群配了自动收集列级统计信息,分区数一多,收集任务本身就要扫成千上万个分区,成了一个常年吃资源的后台负担,还会锁表。hive.stats.autogather 在大表上要非常谨慎。

元数据节点:堆内存和 RPC 都被按文件数收费

以 HDFS 体系为例,命名空间是常驻内存的:每个文件的 inode、每个块的 block 对象、以及块到数据节点的映射都在 heap 里。业界常引用的经验值是"每个块约占 150 字节堆内存"(该数值为广泛引用的经验估算,非本文实测)。按这个口径估:1000 万个文件、每文件 1 个块,光 namespace 就是 1.5GB 级别的堆;如果这些文件平均 3 副本,副本映射的账还要另算。堆一大,GC 停顿就明显,而元数据服务的 GC 停顿是全集群性的——所有人一起卡。

RPC 压力同理。客户端的一次 listing 不是一次 RPC 拿到全部,而是分页串行拉取。文件越多,页数越多,往返次数越多;几十个作业同时进来做 listing,RPC 队列就开始排队。表现出来就是"这条 SQL 明明只读一小部分数据,提交之后却要等很久才开始有任务跑"。

这给选型带来一个非常明确的结论:元数据节点的需求是"大内存 + 稳定的元数据盘 + 独占",而不是"多核"。它不计算,它需要的是把命名空间完整放进内存、把元数据写盘的延迟压到最低、并且不被计算任务的资源占用干扰。

如果底层用了对象存储(S3/OSS/COS 一类),情况类似但更极端:这类存储的 listing 是按前缀分页的,通常每页 1000 条,几万个文件的目录下要几十次往返 HTTP 请求,单次延迟远高于本地 RPC。这就是为什么对象存储上的表对分区粒度格外敏感,而且各家的 SDK 都建议尽量避免深目录 + 大文件数的组合。

治理动作清单:合并、改分区粒度、控并发、清过期

下面是可以直接拿去排期的六项动作,按投入产出比排序。

动作一:调整分区粒度

优先级最高,因为它一次性解决"分区数"和"文件数"两个乘数。默认建议:离线跑批的原始层按天分区,除非单个分区的数据量已经大到 > 数百 GB 且查询经常只要其中一小段时间。按小时分区只在两种情况下合理:单日总量达到 TB 级且下游经常只扫一两个小时;或者数据生命周期极短(只保留几天)且就是为小时级查询服务的。

如果确实需要小时级的查询能力,更好的做法是按天分区 + 在分区内保留小时字段,让下游扫描时用谓词下推(ORC/Parquet 的 row group 级统计能帮你跳过大部分数据),而不是靠分区目录来切。

改分区粒度是一次性重构成本。做法通常是新建一个按天分区的目标表,把历史数据按天重写过去,切完之后删掉旧表。这一步要提前算好:读写两倍的存储临时开销、重写期间的文件数、以及对下游作业的影响窗口。

动作二:控制写入并发与 reducer 数

目标是让"每个 task 在每个分区里写出的文件接近目标大小"。三个抓手:

调 shuffle 分区数。Spark 侧用 spark.sql.shuffle.partitions,Hive on Tez 侧用 hive.exec.reducers.bytes.per.reducer 和 hive.tez.auto.reducer.parallelism。很多集群还挂着默认的 200,在数据涨了十倍之后早就该调了(调大还是调小取决于你的单分区数据量,不是一味加大)。

写之前显式重分区。动态分区写入前加一句 DISTRIBUTE BY 分区键(或 repartition(分区键)),让同一个分区键的数据汇集到同一批 task,直接把"分区数 × 并发数"降成"分区数 × 每分区目标文件数"。这是单次改动收益最明显的技巧。

末尾强制聚合成批。如果只是单纯文件太多,可以在写出前 coalesce(N),N 按"总数据量 ÷ 目标文件大小"算。注意 coalesce 会把压力集中到少量 task 上,数据倾斜严重的场景要先解决倾斜。

动作三:合并小文件(compact / 重写分区)

已经攒下的存量要靠合并来清。几条路:

重写分区。INSERT OVERWRITE TABLE t PARTITION(dt='...') SELECT * FROM t WHERE dt='...' DISTRIBUTE BY 业务键,配合合适的 reducer 数,把这个分区重写成十几个大文件。这是最通用、可控性最好的方式,代价是要占一份临时存储和一份计算。

用表格式自带的合并能力。如果用的是 Iceberg,可以直接走 rewrite_data_files 过程并按目标文件大小配置策略;Hudi 有 clustering;Hive 事务表(ACID)有 major compaction(ALTER TABLE ... COMPACT 'major')。这些机制的优势是元数据一致性由框架保证,不需要自己写"重写 + 切换"的双阶段逻辑。

存储格式自带的合并。ORC 表如果是普通表(非事务表),部分版本支持 ALTER TABLE ... CONCATENATE 做 strip 级别合并,但它只能合并 stripe,且要求文件没有其他中间改动,适用面比前两种窄。

排期建议:按天分区的表,每天作业跑完之后对当天分区做一次合并。把它做成作业流里的一层,而不是攒到季度清理。存量历史上百万文件的分区,一次性全并在生产时段跑容易把元数据服务打挂,建议按天分批、错峰执行,并给 listing 请求限速。

动作四:用批量写入替代频繁微批

一分钟一批改成五分钟一批,文件数直接降到五分之一,数据新鲜度从"最多迟一分钟"变成"最多迟五分钟"。这个取舍要跟业务确认,多数 T+1 报表场景根本不需要分钟级新鲜度。

如果必须保留分钟级延迟,就在写入侧加缓冲层:先落到消息队列或本地持久化队列,再由独立的批量落盘任务按"攒够 N 条或超过 T 秒"触发一次写入。这条路跟工业数采里的"批量入库"的核心是同样的道理:把频繁的小动作改成低频的大动作,接收端和存储层的固定开销从每份数据摊一次变成每批摊一次。

动作五:定期清理过期分区与临时表

很多集群里躺着几个月前调试留下的临时表、跑了半小时失败留下的 staging 目录、以及早就没人用但还在每天增长的中间表。它们不产生价值,但持续占用元数据条目和 namespace。

要做的几件小事:ALTER TABLE ... DROP PARTITION 按保留周期清理(建议写成定时作业而不是靠人记得);给临时表 / 测试库设置自动过期和清理作业;给保留周期固定的核心表配 TBLPROPERTIES 里的保留属性并在文档里写明;失败作业产生的 staging 目录要有固定的清理任务(很多引擎会在作业成功时自动清,但失败时需要外部兜底)。

还有一个容易被忽视的:删除数据在对象存储 / 部分配置下可能有延迟回收策略(回收站 / 版本),要看实际是不是真的放掉了。

动作六:启用合适的输入合并策略

这是收尾的一道防线,也是成本最低的一项。开启 hive.input.format=org.apache.hadoop.hive.ql.io.CombineHiveInputFormat(多数较新版本已经是默认),配好 split 的上下限参数(mapreduce.input.fileinputformat.split.minsize.per.node / per.rack / split.maxsize);Spark 侧打开 AQE 的分区合并并设好目标值,通常设成块大小或两倍块大小。

记住前面那句提醒:它只治已经被切碎的任务数,不治元数据。做完前五项之后再开它,效果最好;只开它而前面的都不做,等于把账从 CPU 挪回 listing。

把前面三笔账和分区粒度放在一起看,趋势就很清楚了:

分区粒度 单分区文件数(典型) 总文件数量级 listing 与调度开销 计算时间占比 建议
按分钟 1–5 个,单文件常在 1–20MB 百万至千万级(一年维度) 极高,listing 与任务启动常常数倍于实际计算时间 通常不足三成 不适合离线跑批原始层,仅用于生命周期极短且明确只查最近几分钟的场景
按小时 5–50 个,取决于写入并发 十万至百万级(一年维度) 高,"提交到第一个任务跑起来"这一步就明显变慢 约三成到五成 仅用于单日总量 TB 级以上且下游确实按小时回溯的表,其余情况改成按天
按天(未合并) 20–200 个,受动态分区写入并发影响 万至十万级(一年维度) 中等,可控但随文件数缓慢上升 约五成半到七成 推荐的默认粒度;写完必须补一步按天分区的合并,并控制写入并发
按天 + 合并后 5–20 个,单文件接近或大于块大小 千至万级(一年维度) 低,任务数接近"数据量 ÷ 块大小"的自然值 七成半以上 目标形态;把合并成本显式排进作业流,每天必做一次,而不是攒到季度清理

以上为典型量级的示意,用于说明趋势,非实测数据。

判断信号:什么时候才真的该加机器

治理做到什么程度可以停手、可以开始采购?看下面几个信号,前面的满足了才轮到后面的。

信号一:作业时间构成变了。从资源调度界面或者引擎的执行计划里,取"作业提交到第一个任务开始的时间"和"所有任务执行完毕的时间"。如果前者(代表提交编译 + listing + 切片 + 排队)占总时长的比例已经降到两成以内,说明非计算开销被压下去了;如果它还占一半以上,继续治理,别买机器。

信号二:平均单 split 字节数到了合理区间。在作业详情页看 Input 的 split 数与总输入字节数,算一下平均每个 split 多大。目标是接近或大于一个块大小(常见 128MB / 256MB)。如果平均只有几 MB,说明还是切得太碎;如果已经上百 MB,切片侧没问题了。

信号三:任务平均运行时长。看 map 任务的平均执行时间。经验上,平均时长在几十秒以内通常说明任务的固定开销占比过高,适合继续合并;平均时长到了数分钟级别,说明每个任务确实在干活,这时候加出来的核才是真在实现业务计算。这条经验值不是实测数据,只能当参考区间。

信号四:资源是真的打满了。在作业执行那段时间里抽节点 CPU 利用率、内存占用、磁盘利用率、网卡吞吐。如果 CPU 持续在 80% 以上、内存接近配额、且这段时间占总时长比例高,这才是"机器不够"的证据。如果 CPU 一直在 30% 徘徊而墙钟时间很长,那是典型的"在等,不在算"。

信号五:加一半资源,时间掉一半。这是最有说服力的验收方式。临时给这个作业单独扩大资源配额跑一次,看墙钟时间是不是近似线性下降。如果是,说明瓶颈确实在计算能力上;如果只掉了 10%,说明瓶颈在别处(很可能是提交、listing、driver 提交这一步),回去继续治理。

把这五条串起来就是一句话:时间花在哪,就治哪;先用量化的时间构成证明"是在算",再去谈加机器。反过来做的结果是买了一堆机器,作业照样慢,还多了一份每年的续费账单。

服务器怎么配:元数据节点与计算节点看重的资源不一样

一旦确认要扩容或新建离线数仓集群,最容易犯的错是"所有节点用一个配置"。元数据节点和计算节点的瓶颈资源完全不重叠,混着配等于两头都亏。

元数据节点:内存第一,磁盘要稳,绝不混布

内存。核心指标是"命名空间能不能完整放进堆里、还留有余量"。算法是:预估总文件数(不是数据量)× 每块的堆开销经验值 + JVM 本身的开销 + 安全余量,再按这个数 × 1.5 左右留余。文件数怎么估?用现在的日均文件增量 × 保留天数 × 增长系数,再算上副本相关的元数据。举个例子(以下为典型部署思路,并非特指某一真实客户):某集群日均新增 20 万个文件,数据保留 2 年,即 namespace 量级在 1.5 亿文件这个级别——这种规模已经不是单机能扛的,要考虑联邦 / 分库或者把表拆到多个集群。所以元数据规划的第一步永远是控制文件数,第二步才是给内存。上 GB 级的堆意味着要配好 GC 策略并准备做 HA,主备双方的元数据一致性同步本身也是一笔持续的写 IO。

Hive Metastore 背后的关系库同样吃内存:分区表、文件表的索引要能放进缓冲池。分区数上千万的表,索引本身就可能几百 MB 到 GB 级。

磁盘。元数据盘的特点是小写延迟敏感 + 绝对不能丢。推荐企业级 SSD 做 RAID 1(或 RAID 10),不要追求容量,追求的是写延迟的稳定性和断电保护。元数据节点的 edit log 是同步刷盘的,一次写延迟抖动会让全集群跟着卡——这就是为什么消费级 SSD(带 SLC 缓存那种)不适合做元数据盘,缓存在 7×24 持续写入下没有回收窗口,掉速会直接传导到所有客户端。

不混布。元数据节点上不要跑计算任务,也不要跑容器。它需要的是稳定的 CPU 调度和稳定的内存,一旦被一个临时容器抢占内存导致 GC,代价是全集群不可用。如果预算实在有限至少要做资源隔离和严格的内存上限,但能拆就拆。

计算节点:核内存比、磁盘 IOPS、网卡,三样缺一不可

核内存比。跑批容器的内存是按并发容器数摊的,不是按整机的内存摊。常见配比是每个物理核对应 4 到 8 GB 内存这个区间:如果计算机都是轻量扫描、几乎不 spill 到磁盘,可以往上走(多核少内存);如果作业有大 JOIN、大排序、或者要做反序列化和膨胀(比如 explode、JSON 解析),内存需求会明显上升,就要往下压核数加内存。选型前先看自己的作业画像:过去一周内所有作业 spill 到磁盘的字节数总和是多少,spill 量大就说明内存不够或者核数超额。

磁盘。计算节点的磁盘压力来自三处:读输入数据、写中间 spill、shuffle 阶段读写。夜间跑批时这三样会同时打满,而这正是容易出现"作业忽快忽慢"的原因。盘数不够是常见短板:单块企业级 SATA SSD 的顺序吞吐在数百 MB/s 量级,几块做 JBOD 分摊到不同容器上,聚合带宽才能喂得上多核。不要用一块大容量盘扛所有 container 的 IO。机械盘不是不能用,但只适合做冷数据归档层,不适合出现在跑批关键路径上。

容量这边反而简单:容量按保留周期算。日均数据量 × 保留天数 × 冗余系数(副本数 + 临时文件开销 + 压缩前的临时空间,通常取 2 到 3,取决于有没有纠删码或压缩),再加 20% 到 30% 的运维余量给重写分区、临时导出、扩容迁移。IOPS 按并发任务数算:并发容器数 × 单容器的期望 IOPS/带宽 = 需要提供的盘能力,再往上留 30% 峰值余量。

网卡。跑批时最容易低估的一项。Shuffle 阶段所有节点之间要互相拉数据,一个宽 shuffle 能让多个节点的网卡同时打满,由此产生的队列抖动与重传直接表现为作业时间不可预测。跑批集群建议万兆(10GbE)起步,中大型集群上 25GbE 甚至更高,并且确保交换机侧的收敛比别做太狠——接入带宽堆到一定程度之后,瓶颈就跑到上行口去了。

把两种角色拆开去看选型

实际落地时,把需求写成两张单子:一张给元数据角色(大内存 + 企业级 SSD RAID 1 + 独占 + HA),一张给计算角色(合理的核内存比 + 多盘 JBOD 或阵列 + 万兆以上网卡 + 足够的内存余量)。两张单子对应的机型完全不同,价格结构也不一样。

把这两张需求单拿去对照机房侧的机型与存储档位时,一万网络这类深耕 IDC 19 年(成立于 2007 年)的服务商,可以作为离线数仓节点(元数据节点与计算节点分角色配置)的选型参考之一;具体机型、内存档位、磁盘与网卡配置以及对应的价格需实时询价,建议拿自己算出来的文件数量级、内存需求和网卡档位去核,而不是按套餐默认值下单。

跑批窗口不够用时的两条路:加机器 vs 改编排

凌晨 0 点到 6 点,这六个小时是所有离线作业的公共预算。作业越加越多,窗口自然越来越挤。这时候有两条路,先选对路再谈预算。

第一条:加机器 / 扩并发

适合的场景是:单个作业本身的瓶颈在计算能力上(符合上一节那五个信号),横向扩容能让墙钟时间下降。它的上限也很明显——

串行部分不随扩容下降。提交编译、listing、driver 侧提交、以及依赖链路较宽的那部分,加再多节点也不会变快。作业总时长永远被这些固定成本压着一个下限。

共享瓶颈不会因为它加这一个作业的资源而改善。元数据服务是共享的:这个作业起的任务多了,别的作业等待 listing 的时间就长。夜间所有作业一起扩容,元数据侧的排队往往比之前更严重,出现"大家都加了资源,整体产出时间反而更晚"的反直觉结果。

资源抢夺会引发新的抖动。并发上去了,磁盘和网络同时打满,作业的完成时间方差变大。而 业务感知的不是平均完成时间,是最晚完成时间。方差变大比平均变慢更让人难受。

第二条:改编排

这条路的上限更高,且不花钱:

错峰。把大类:核心报表、对内查询、对外展示、冷数据加工——按重要性和截止时间(SLA)排优先级,先跑有截止时间压力的,把那些"跑完了更好,晚两小时没事"的放到后半夜或者白天低峰。很多团队把所有作业塞进同一个起跑位,纯属习惯。

拆依赖。画一张作业依赖图,找出那条最长的链(关键路径)。关键路径上每减少一个作业的时间,总时长就少一段;而非关键路径上的作业慢一点,完全不影响最终产出。资源应该优先投在关键路径上,而不是平均分。

增量替代全量。这是收益最大的一项:把"每天重算一年数据"改成"每天只算昨天 + 和已有结果合并"。任务数据量从"全量"降到一个或两个数量级以下,之前所有跟文件数、listing、任务数相关的问题也一起下降——因为参与的文件本身就少了。代价是要设计好幂等的合并逻辑、迟到数据的处理窗口、以及失败重跑的边界。

合并作业。上下游两个作业如果都是轻量加工,可以考虑合并成一个,省掉中间落盘、省掉一次 listing 和一次提交。

怎么选:先画关键路径,看那些把窗口撑住的作业,是不是真的在计算。如果在计算且已经满足上一节的五个信号,加机器;如果有一半时间在等 listing 或者等资源,改编排和治理文件。多数团队的顺序搞反了——先买了机器,半年后发现还得回头做增量和分区治理。

现场被反复问到的七个问题

文件多大算小文件?

别记绝对数字,用相对标准:显著小于存储块大小的文件就算小文件。常见部署的块大小(dfs.blocksize)是 128MB 或 256MB,那么低于几十 MB 的文件基本都属于这一类。更实用的判断方式是看"这个文件能不能喂饱一个任务"——如果一个任务的实际执行时间只有几百毫秒到一秒出头,而它的启动成本是数秒量级,那它读的就是小文件。还有一种情况:并不是单个文件特别小,而是单目录下的文件太多(比如一个分区几千个文件),listing 慢同样致命,这时候的问题叫"文件数太多"而不是"文件太小",两者经常一起出现,治理动作也基本一样。

按小时分区是不是一定错?

不是。判断标准只有一条:单个分区能不能装下足够多的"大文件"。如果单日总量在 TB 级,按小时分之后每小时仍有数十 GB,每个小时分区下放几十个接近目标大小的文件,这完全是合理的;如果单日总量只有几个 GB,按小时分出来每个小时才几百 MB,这时候按小时就是错误的粒度,应该改成按天,需要小时级查询就用分区内的时间字段 + 列式存储的谓词下推来实现。反过来也有第二种例外:数据生命周期很短(只留几天) 、且下游几乎都按小时查,那短暂的分区膨胀不构成长期负担,可以按小时。

合并小文件会不会影响正在跑的作业?

会,而且这是风险最高的一步,必须按流程做。三点注意事项:一是隔离。compact 作业要跑在独立的队列或低峰时段,它会同时消耗 IO(读旧文件 + 写新文件)和元数据 RPC(大量 rename 和 listing),和生产作业争起来会互相拖垮。二是原子切换。不要在原分区目录里就地合并。正确做法是写到临时目录或临时表,写完校验(行数、文件大小、抽样质量)之后再 rename 切换。就地改会让并发的读作业读到写了一半的文件。三是要能回滚。切换前先记录旧分区的位置并保留一段时间,新数据有问题能快速退回去。还有一点,正在被写入的那个"当天分区"不要在写入进行中合并,等当天作业跑完再做。

元数据节点内存要多大?

完全取决于你要承载多少文件和多少个块,跟集群的磁盘容量没有直接关系——一堆 10TB 的大文件和一个十亿个小文件,前者对内存的要求可能低一个数量级。算法是:(预估文件数 × 平均每文件块数 × 每块的堆开销经验值)+ JVM 及自身开销,再乘 1.5 倍左右的余量;这里的"每块堆开销经验值"业界常引用 150 字节这个量级的估计(非本文实测),实际以你用的版本和 JVM 配置为准。关键是先算出预估文件数——用日均新增文件数 × 保留天数,把趋势算出来。如果算出来发现两年后要上亿文件,先别急着买大内存机器,那是架构信号:要考虑联邦、拆分集群,或者更根本地——先把文件数治理下去。同时也别忘了给 HA 留资源,主备节点都要同等规格。

增量表为什么也会攒出小文件?

因为增量讲的是"每次处理的数据量小",小文件讲的也是"每次写出的数据量小",这两件事天然同向。一次增量可能只有几万行、几十 MB,按正常写出就是一个几十 MB 的文件;一天增量 144 次,就是 144 个小文件;两周不治理,一个分区里就堆出几千个。而且增量表的下游常常要求低延迟,这跟"攒够再写"直接冲突。破解办法四条:一是把增量频率降下来(小时级而不是分钟级,多数业务的真实需求是小时级);二是引入支持小文件自动合并的表格式(Iceberg、Hudi 这类都内置了合并策略,可以配目标文件大小和触发条件);三是在写入侧加缓冲,让多次变更攒成一批落盘;四是接受一个现实:写入频率高,文件就必然多,低延迟和大文件很难兼得,那就把合并任务作为固定环节排进作业流,定时整理而不是指望写入侧一次成型。

重跑补数怎么避免文件翻倍?

核心是"写之前先清理,写完再切换"。四条实操建议:一,补数只用覆盖写(INSERT OVERWRITE)到明确的单个分区,不用 append。append 补数在不做后续清理的情况下必然留下两套数据。二,失败重跑要走 staging。让作业先写到 staging 临时路径,成功之后再 rename 到正式分区目录;失败时不 rename,临时路径由固定的清理作业删掉。这样重复跑十次也只留下末尾那一次的结果。三,灰度验证再全量。先用一天或一个分区验证产出正确,确认之后再放开全量,避免三遍重跑留下三份。四,补数完成后立刻做一次 compact。覆盖写本身也可能因为 task 重试留下少量残留文件,补完 compaction 收尾,并且在 staging 和正式目录两边都要做。还有一个习惯:给每个涉及重跑的操作记一份操作清单,跑之前写清楚"这次会动哪些分区、用什么模式写、跑完怎么验",别依赖现场记忆。

加机器和治理小文件哪个性价比高?

在还没做过治理的集群上,治理几乎一定赢,而且差距很大。原因有两层:第一层,治理是一次性(或低频持续)成本,加机器是每年都要付的持续成本。一次分区重构 + 合并 + 写入参数调整,做完之后每份数据、每个作业都受益,且效果持续;买机器则是月月付钱,而且数据继续往上涨,明年可能还得加。第二层,加机器治不了那部分时间。listing 慢、任务启动开销、driver 串行提交,这些跟 CPU 核数不在一个维度上,加多少核都降不下来。

例外情况是:治理已经做完(上一节五个信号都满足),作业的"计算那一段"占比很高且 CPU 持续打满,横向扩容能带来近似线性的时间下降——这时候加机器才划算。判断方法就是前面说过的那个:临时给这个作业加大资源跑一次,看墙钟时间掉多少。掉得少就别买。

先数文件,再数核数

回到开头那个假设的团队。他们最终做的事里,一件跟 CPU 无关:把原始层从按小时分区改成按天分区,写完了之后对当天分区做一次 compact,动态分区写入前加了 DISTRIBUTE BY,把 shuffle 分区数从默认的 200 调到了跟单分区数据量匹配的值,补数统一改成"staging + 切换",Input 侧打开了合并策略。做完这几件之后,同一个作业从 4 小时回到 50 分钟上下,而且不再因为时间推移而变慢。

这里真正值钱的不是某个参数,是顺序:先把时间构成量出来,看清楚那不到一半的"计算时间"之外,剩下的时间都去哪了;数清楚一个分区里躺着多少个文件、这个数字每个月涨多少;然后才轮到讨论要不要加机器。多数情况下你会发现,账单上缺的不是算力,是秩序。离线数仓的选型也是一样——先看元数据节点需要多大的内存来装下命名空间,再看计算节点需要怎样的核内存比和网卡,这两个问题的答案差别很大,想清楚了再去对机型和档位。

文中量级是趋势示意不是实测

本文提到的"每块的堆内存开销经验值""任务平均执行时长的参考区间""task 固定开销的量级""empirical 压缩比的变化趋势",都是便于建立判断框架的量级示意,不是在任何特定集群上的实测数据。不同引擎版本、不同存储格式、不同 block size 配置、不同的硬件配置下,这些数字会显著不同,偏离一倍以上是常事。

真正需要你自己在生产环境里测的三件事:一是你自己的 task 固定开销是多少——在空闲集群上跑一个读单个小文件的任务,从提交到结束量一遍;二是你的 listing 时间是多少——提交一个只读少量数据但分区很多表的查询,量"提交到第一个任务跑起来"这一段时间;三是你的单任务 IO 和网卡的实际表现——在跑批高峰抽一个计算节点量磁盘利用率和网卡吞吐。这三个数拿到了,本文所有的示意数字都可以替换成你自己的,那才是能用来做决策的数字。至于机型、配置与价格,只以官网实时信息为准。


上一篇:2026 一台服务器能撑多少 WebSocket 长连接:卡住上限的往往不是带宽

下一篇:2026 Flink 实时计算服务器租用配置手册:Checkpoint 状态后端与本地盘的选型避坑全解