跑批慢了,第一反应通常是"机器不够了"。但把作业时间摊开来看,真正花在读数据和算数上的那一段,往往连一半都不到,剩下的全交给了分区 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,收益比优化数据处理本身更直接。
这一段是本文的主角。引擎要拿到"这个表的哪些分区符合 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)在它负责的每个分区键下都会生成一个或多个文件。表面上看你只是写一张表,实际落地的文件数是"分区数 × 写并发数"。
举个例子:某天的数据分布在 24 个小时分区里,写入作业有 200 个 reducer(很多引擎的 shuffle 分区默认值就是 200),理论上最多会产生 24 × 200 = 4800 个文件。实际会少一些,因为不是每个 reducer 都持有每个分区键的数据,但量级在这。如果一个 reducer 每个分区只有 5MB 数据,那落下来的就是 4800 个 5MB 的文件。
而且这个数还会随数据量变化:数据涨了,引擎自动把 reducer 数调大,文件数跟着涨,单文件大小却没怎么变。这就是为什么"数据量明明只涨了三倍,文件数涨了十倍"。
每分钟落一批的增量写入,一年下来是 52 万批。哪怕每批只写一个分区下的一个文件,一年也是 52 万个文件。如果每批要写多个分区(按业务 key 分桶或多个指标表),文件数还要再乘。
这类写入天然和"大文件"冲突:单个批次的时间窗口决定了数据量上限,一分钟的数据量再大也大不到哪去。要破这个局只有两条路:一是拉开批次间隔(一分钟改成五分钟或十分钟),二是在批次之外加独立的合并任务,定期滚动做一次分区内的文件合并。
上游表重算、补历史数据、修 bug 之后回填,这些操作本身没问题,问题在于写入方式。如果用 append 追加写而不先清掉目标分区,同一份数据会在同一个分区里留下两套文件:一套是旧版本,一套是新版本。查询时如果逻辑是按主键取最新,逻辑上没错,物理上文件数翻倍了。
还有一种更常见的:回填作业本身跑了三次(第一次失败、第二次数据不对、第三次才成),前两次留下的中间文件没有通过 staging 机制清理干净,静静躺在分区目录里。这种情况下 SELECT COUNT(*) 的结果可能是对的(如果有去重逻辑),但 listing 和 split 的数量已经翻倍。
覆盖写在语义上是"替换整个表或分区",但这个"替换"依赖于实现是不是干净。几种常见情况会让旧文件留下来:用了不支持 overwrite 语义的存储格式或外部表配置;写到了 S3 这类对象存储上但清理逻辑因为请求失败而静默跳过;写入过程中 task 失败重试产生的残留文件没有回归一清理作业清掉。
还有一个大家心照不宣的情况:用 INSERT OVERWRITE 覆盖分区时,如果开了动态分区且某些分区键在这次写入里没有数据,那些分区不会被"清成空",而是保持原样。这本是合理设计,但如果下游以为被覆盖了,就会出现新旧数据混在一个分区里的诡异结果。
这一节把"文件数比数据量致命"拆成三笔账:调度账、打开账、尾巴账。
前面说过,每个 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 慢,这些它一个都管不了。它是止痛药,不是根治。
这一节解释那句让人最困惑的抱怨:"我什么都没改,数据量也没涨多少,为什么一天比一天慢。"答案通常在两个元数据服务上。
Metastore 的后端是一张关系库,分区信息存在 PARTITIONS 表里,每个分区一行。一个上千分区的表,一次查询要从这张表里拉出上千行,再做过滤、排序、分页返回给客户端。这张表的行数增长会让本来几十毫秒的查询变成几秒。
更糟的是后面那一轮:对每个命中的分区,客户端(或 metastore 内部)要去拿这个分区下的文件列表。分区数 × 每分区文件数,这才是真实的 listing 规模。分区数从 365 涨到 8760(按天改按小时),如果每个分区的文件数不变,listing 粒度直接翻了 24 倍。
还有一个独立性很强的问题:统计信息。很多集群配了自动收集列级统计信息,分区数一多,收集任务本身就要扫成千上万个分区,成了一个常年吃资源的后台负担,还会锁表。hive.stats.autogather 在大表上要非常谨慎。
以 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 级统计能帮你跳过大部分数据),而不是靠分区目录来切。
改分区粒度是一次性重构成本。做法通常是新建一个按天分区的目标表,把历史数据按天重写过去,切完之后删掉旧表。这一步要提前算好:读写两倍的存储临时开销、重写期间的文件数、以及对下游作业的影响窗口。
目标是让"每个 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 上,数据倾斜严重的场景要先解决倾斜。
已经攒下的存量要靠合并来清。几条路:
重写分区。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,代价是全集群不可用。如果预算实在有限至少要做资源隔离和严格的内存上限,但能拆就拆。
核内存比。跑批容器的内存是按并发容器数摊的,不是按整机的内存摊。常见配比是每个物理核对应 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 年)的服务商,可以作为离线数仓节点(元数据节点与计算节点分角色配置)的选型参考之一;具体机型、内存档位、磁盘与网卡配置以及对应的价格需实时询价,建议拿自己算出来的文件数量级、内存需求和网卡档位去核,而不是按套餐默认值下单。
凌晨 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 和网卡的实际表现——在跑批高峰抽一个计算节点量磁盘利用率和网卡吞吐。这三个数拿到了,本文所有的示意数字都可以替换成你自己的,那才是能用来做决策的数字。至于机型、配置与价格,只以官网实时信息为准。
Copyright © 2013-2020 idc10000.net. All Rights Reserved. 一万网络 科技有限公司 版权所有 深圳市科技有限公司 粤ICP备07026347号
本网站的域名注册业务代理北京新网数码信息技术有限公司的产品