1. 先算两笔账,你就明白为什么偏偏是“数据本地化”
1.1 数据搬运的隐性成本,比你想的贵得多
说到大数据架构的性能优化,很多人第一反应是换引擎、加资源、调并行度。但我干了这么多年基础架构,真正能稳定扛住大流量、大任务、且性价比最高的优化实践,往往只有四个字:数据本地化。简单说,就是让计算尽可能靠近数据所在的节点去执行,而不是把数据从存储端搬到计算端。这个概念听起来朴素,落地起来却牵涉任务调度、副本分布、文件格式、网络拓扑,甚至是存算分离时代的缓存设计,水很深。这篇文章适合正在维护生产集群、被离线跑批和即席查询折磨的开发运维同学,也适合刚入行大数据、想从原理上理解“为什么同一套SQL在A集群快、在B集群慢”的读者。
先算一笔物理账。一个普通数据节点的本地磁盘读带宽,SATA SSD大概在 500MB/s 到 1GB/s,NVMe SSD 能到 3GB/s 以上,万兆网卡的理论带宽是 1.25GB/s,但实际跑满也就 1GB/s 上下。也就是说,数据从本地盘读出来的速度,是不输给、甚至远高于从远端节点通过网卡拉数据的速度。可问题在于,绝大多数集群的网络并没有跑满,瓶颈反而出在别的地方。
网络传输一个大文件,不是简单地把字节从一个网口搬到另一个网口。数据要先从磁盘读出来,经过操作系统页缓存,序列化或封装成网络包,经过TCP/IP协议栈、网卡中断、交换机转发,到达目标节点后还要校验、反序列化,再写进目标节点的内存或磁盘。这中间每一跳都有开销,而且这些开销会随着数据量的增大被成倍放大。我见过很多集群,CPU 利用率不到20%,网络也跑不满,但任务就是慢,异常慢。做 profile 之后发现,时间全耗在数据移动的路径上了。
所以“计算靠近数据”这个策略的出发点特别简单:既然移动计算的开销通常远小于移动数据的开销,那就让计算去找数据,而不是让数据来找计算。分布式框架之所以要有一堆 task 调度、副本摆放的机制,本质上都是在为这一句话服务。
1.2 从一次凌晨跑批说起,本地化省下的时间是可以直接感知的
前两年我接手过一个离线数据仓库集群,跑批时间经常压着 SLA 的线,业务侧每天催。我们最初也以为是 Hive SQL 写得差,调了一轮 MapJoin、改了一轮分区,效果有一点,但没解决根本问题。后来我看 Spark UI 上每个 stage 的 Locality Level 分布,发现大量 task 都是 RACK_LOCAL 甚至 ANY,也就是计算节点和它要读的 HDFS 数据块根本不在同一个节点上。
那次跑批的数据规模大概是 6 万多个文件,总量在 3TB 左右。当时集群是 10GbE 网络,跨节点读的理论上限是 1.25GB/s,但实际每个 task 因为要跟远端 DataNode 建连接、做校验,真实吞吐能到 700MB/s 就算不错了。也就是说,3TB 数据如果全靠网络搬运,最理想也要跑 4200 秒以上,也就是 70 分钟。但如果让每个 task 直接读本地磁盘,3TB 从本地 NVMe 上读出来,理论上只要 1000 秒左右,差了四倍多。
当我们把副本分布、机架感知、调度等待这些参数调好,任务本地执行率从 40% 左右升到接近 90% 的时候,整个跑批时间缩短了将近一半。不是 SQL 变聪明了,也不是并行度变高了,纯粹是数据不用再来回搬家了。从那以后,我对大数据架构优化有了一个很深的体会:调度层的优化,往往比执行层的微调来得更猛,也更隐蔽。
2. 拆开数据本地化:从进程内亲和到跨机架调度
2.1 三个本地层级,先搞清楚调度器眼里的“距离”
讨论数据本地化,绕不开调度器眼中的“距离”概念。无论是 Hadoop YARN 还是 Spark Standalone,调度器都会给每个 task 打一个 locality 标签,越是靠近数据,标签优先级越高。常见框架里一般分这么几个等级:
| 本地级别 | 含义 | 相对代价 |
|---|---|---|
| PROCESS_LOCAL | task 与输入数据在同一个 JVM 进程内 | 最低,几乎无网络和序列化开销 |
| NODE_LOCAL | task 与输入数据在同一台物理机,但不在同一进程 | 低,可能需要跨进程读共享内存或本地磁盘 |
| RACK_LOCAL | task 与输入数据在同一机架内,但不同节点 | 中,需要经过接入层交换机 |
| ANY | task 与输入数据跨机架甚至跨数据中心 | 高,网络传输和延迟损耗最明显 |
很多人以为本地化就是“数据在某节点,就把任务调度到某节点”,其实还差一层。哪怕在同一个节点上,如果 task 跑在一个新起的 JVM 里,而数据在另一个 JVM 管理的磁盘上,依然要二次读盘。PROCESS_LOCAL 往往意味着执行器复用了之前已经加载过数据的进程,连反序列化和对象复用都省了,这在 Spark 这种内存计算模型里收益非常明显。
理解了这几个层级,你再去看调度器的各种等待策略就会顺畅得多。因为调度器不是“一定要”把任务调度到本地,而是在“尽可能”和“必须执行”之间做权衡。
2.2 延迟调度:调度器是怎么把任务塞到数据旁边的
调度器要做的事,通俗点说就是“先看看数据在哪,再决定任务往哪发”。但分布式集群的资源是动态的,数据块有多个副本,节点的负载也在不断变化,所以调度器不能一上来就硬等本地节点空出来,那样可能等到天荒地老。这里的关键机制是延迟调度(Delay Scheduling)。
以 Spark 为例,当 TaskSetManager 分配一个 task 时,会根据输入的 HDFS 块位置生成候选节点列表。如果当前有空闲 executor 正好在候选节点上,直接调度过去,完美。如果没有,调度器不会立刻把任务派给一个远端空闲节点,而是先把这个 task “挂起”一小段时间,比如默认的spark.locality.wait是 3 秒。在这 3 秒里,如果任何候选节点上的 executor 空出来了,就调度过去;如果一直等不到,就降级到下一级本地性,比如从 NODE_LOCAL 降为 RACK_LOCAL,再等一小段,最后实在不行才以 ANY 级别调度出去。
这段等待时间就非常讲究。设太短,本地化率会下降;设太长,集群空转的 slot 会变多,并行度受影响。我见过有团队把spark.locality.wait调到 10 秒,结果整个集群利用率肉眼可见地下降,因为大量 task 都在“等本地位置”,实际上空出来的资源反而被浪费了。延迟调度本质上是用“一点点等待”换“少挪很多数据”,这个交换通常划算,但前提是等待时间必须跟集群的负载、任务大小匹配。
2.3 一个必须做好的权衡:本地性等待 vs 集群吞吐
这个权衡没有标准答案,取决于你的任务画像和集群规模。如果一个大任务跑 2 小时,多等 3 秒换本地执行,愿意。但如果一个任务本身只要 30 秒,你让它先等 3 秒再找本地位置,就太亏了。我自己的经验是,把等待时间拆开来看更科学。
一个比较粗暴的参考建议:如果以离线大任务为主,节点数量较多、资源相对宽裕,可以用默认值甚至略微调大一点;如果集群是多个业务共用的、小任务密集,把spark.locality.wait调到 300ms 到 500ms 会更平衡。还有一种做法,是针对不同 task 类型单独设置spark.locality.wait.node和spark.locality.wait.rack,让 NODE_LOCAL 的等待短一点,PROCESS_LOCAL 的等待可以相对长一点,因为进程内复用带来的收益更大。
另外,HDFS 的副本放置策略也直接影响调度的选点。默认的 BlockPlacementPolicy 会把第一个副本放在客户端所在节点,第二个副本放在同机架另一个节点,第三个放在不同机架。如果机架感知配置错误,调度器把所有节点都当成同一个机架,那本地化级别计算就会失真,调度器大概率会做出“看似还行、实则很亏”的决策。这是很多集群本地化率上不去的隐藏原因,后面我在踩坑部分会专门讲。
3. 存算分离之后,本地化变成“缓存亲近”
3.1 当数据放进对象存储,传统本地化策略失效了
现在不少新集群已经不做存算一体了,而是把数据放在 OSS、S3 这类对象存储上,计算节点按需拉起。这种存算分离的大数据架构弹性好、成本可控,但代价是数据离计算节点“远”了,不是物理距离,而是访问路径变长。对象存储不提供 HDFS 那样机架感知的副本摆放信息,调度器很难判断“数据在哪”,因为数据根本没有固定在某台机器上。
我在做过一个基于对象存储的实时数仓项目,初期查询时发现每次跑结果集,都要从对象存储重新拉一遍数据,网络开销非常大。表面上看起来是 SQL 慢,实际上每一份数据都经过了“对象存储->计算节点内存”这条漫长链路。这时候传统意义上的数据本地化已经失效了,你能依靠的,是把“最近读过的数据”留在本地,并想办法让下一次查询尽量复用。
这就把问题从“任务调度靠近数据”变成了“数据缓存靠近计算节点”。思路是一样的,只是“本地化”的对象从 HDFS 数据块变成了本地缓存文件或缓存页。所谓计算靠近数据的优化实践,在存算分离架构里几乎等价于“缓存亲和调度”和“缓存命中率”的优化。
3.2 我实践过的一套“缓存亲和调度”方案
当时我们用的查询引擎支持自定义调度策略,核心做法分三步。第一步,给每个 worker 节点挂载一块本地 NVMe 盘作为缓存目录,这块盘不存业务数据,只放临时缓存。第二步,读取对象存储上的数据时,按“表名+分片ID+文件块ID”做一致性哈希,把每个分片的缓存固定落到 N 个 worker 上,而不是随机落。第三步,查询调度时,扫描任务会先读取分片元数据,知道哪个 worker 上有对应缓存,然后优先把 scan task 调度到这些 worker 上。
这套方案的关键在于,缓存分片和计算任务必须“绑定”在同一组 worker 上。否则即使每个节点都有缓存,下一次查询依然可能被调度到没有缓存分片的节点,那缓存就白做了。我们当时用调度队列的分区标签来实现亲和:缓存落盘时在 worker 上打标签,调度时只把对应任务的 scan 请求发到带标签的 worker 上。
实测效果很稳。热点表的缓存命中率从最初的不到 20% 提升到 70% 左右,P99 查询时间缩短了接近一半。还有一点容易被忽略:当缓存命中的时候,数据直接从本地 NVMe 读,不走对象存储 SDK,也不做网络签名和校验,CPU 的使用率也下来了。整个集群的 QPS 能力因此提了一个台阶,而不仅仅是单个查询变快。
3.3 别把“缓存亲和”做成“数据倾斜”
这套方案也有坑。最开始我们用一致性哈希把分片均匀分布到所有 worker 上,看起来没问题,但查询热点并不是均匀的。某个大卖家或某个热门时间段的数据分片可能被反复查询,而它恰恰哈希到了同一个 worker。那个 worker 的磁盘 IO 被打爆,其他 worker 闲得发慌,现象就是任务偶发抖动。
后面我们加了两个机制才解决问题:一是给一致性哈希引入虚拟节点,让每个分片可以映射到多个候选 worker,调度时轮流选;二是做权重感知,每个 worker 上报当前磁盘 IO 和缓存容量,调度器在候选 worker 中挑负载最低的。说白了,缓存亲和不能做成死绑定,要留一点弹性。否则你为了本地化牺牲了资源均衡,结果往往适得其反。
另外缓存淘汰策略也很关键。LRU 是最常见的策略,但如果你有大量周期性全表扫描任务,LRU 会被全表扫描的数据污染,热点数据反而被挤出去。我们后来改成“两级缓存”:第一级按访问频率保留热点分片,第二级才用 LRU 兜底。这个调整看起来很小,但实际对命中率的影响非常大。
4. 真正的“移动计算”,存储层比调度层更关键
4.1 文件格式与索引裁剪:把计算下推进数据内部
聊完调度和缓存,再把视角拉到存储层。数据本地化并不只是“任务在哪个节点跑”的问题,还包括“计算到底要碰多少数据”。如果存储格式很烂,即使任务完全本地执行,也要花大量时间扫无用数据。反过来,如果存储层能把大部分数据提前过滤掉,网络和磁盘的压力都会大幅下降,这相当于把“计算”下推到了离数据更近的地方。
列式存储是这一思想最典型的落地方式。拿 Parquet 或 ORC 来说,文件会被拆成多个 row group 或 stripe,每个块都带有 min/max 统计信息和布隆过滤器。执行引擎在做扫描时,会把这些统计信息当作粗粒度索引来用,直接跳过不可能满足查询条件的块。用户写一条where dt='2024-01-01',如果分区裁剪已经砍掉大部分文件,扫描层再用 min/max 跳过 row group,实际读取的数据量可能只有原始数据的 1%,甚至更低。
我见过太多团队把精力全花在调 Spark 参数上,却不去做列式存储改造。同样的 SQL,在 TextFile 上跑 20 分钟,换成 ORC 之后 2 分钟就能跑完,提升根本不是靠参数,而是靠数据本地化思想在存储层的体现:尽量让“计算”发生在离数据最近的扫描阶段。
4.2 副本感知与短路读:省掉一次可避免的“本地拷贝”
再往深一层走,数据本地化还有一个容易忽略的点:就算数据已经在本地节点了,读取路径是否足够短。HDFS 的传统读取路径是客户端通过 RPC 跟 DataNode 通信,由 DataNode 把数据从磁盘读出来再转发给客户端。这个过程即使发生在同一台机器上,也绕了一大圈,数据从磁盘到系统缓存,再被 DataNode 进程处理,再通过 socket 传回客户端进程。
为了省掉这一步,HDFS 提供了 ShortCircuit Read(短路读)机制。开启后,客户端进程可以直接用文件描述符访问本地磁盘上的块文件,不需要额外经过 DataNode 的 RPC 转发。听起来是内部小优化,但对本地化执行率高的任务来说,往往能再带来 10% 到 20% 的吞吐提升。
相关的核心参数包括dfs.client.read.shortcircuit和dfs.domain.socket.path。有一点特别注意:这两个配置必须在客户端和 DataNode 端保持一致,否则短路读会以失败告终,而失败后系统会默默回退到普通 RPC 路径,表面上不报错,但性能直接受影响。这类“静默回退”最坑人,很多人以为配好了,其实根本没生效。
4.3 一组非常有效的组合拳:列式裁剪 + 分区裁剪 + 调度亲和
说一个我印象很深的实际案例。业务侧有一条取数需求,要查一张 300GB 的事实表,关联两张维表,过滤条件就一个city_id='xxx'。最初实现是读全表,跑 20 分钟,资源占用还高。我们做优化时没有改一行业务逻辑,只做了三件事:把表改成 ORC + 分区表,按天分区;在 city_id 上建布隆过滤并开启谓词下推;再把调度改成 NODE_LOCAL 优先。
改完以后,300GB 的表单次查询实际扫描数据只剩 3GB 左右,任务跑完不到 40 秒。这 20 倍以上的提升,不是 CPU 变快了,也不是集群变大了,而是计算真正“贴”着数据走:分区裁剪先砍掉 80% 文件,列式索引再砍掉 90% 的 row group,最后调度亲和保证剩下的 3GB 尽量从本地磁盘读。这个组合让我意识到,数据本地化不是某一个开关,而是调度、存储、执行引擎三层协同的结果,缺一个都发挥不出最大价值。
5. 生产集群上,我踩过的本地化相关的坑
5.1 本地执行率只有 25%,问题出在机架感知脚本
有段时间我们一个 Spark 集群的跑批任务明显变慢,打开 Spark UI 看 Locality Level,大部分 task 都是 RACK_LOCAL 和 ANY,NODE_LOCAL 少得可怜。我当时第一反应是资源不够,加了节点也没改善,后来才去查 HDFS 的机架感知配置。
结果发现,集群里的拓扑脚本(topology script)被改坏了。脚本本身是网络团队早期用来上报交换机位置的,后来机房做了一次网络架构变更,交换机的命名规则变了,脚本返回的机架 ID 全部变成了default,导致 HDFS 认为所有节点都在同一个机架里。调度器在计算数据本地性时,自然没办法区分 NODE_LOCAL 和 RACK_LOCAL,相当于所有的节点距离都是 0,只能退到 ANY 级别去执行。
修复方法很简单:修正拓扑脚本,让每个节点能正确上报自己的机架信息,然后执行hdfs dfsadmin -printTopology验证每个节点都映射到了正确的机架上。这个配置看起来基础,但它影响的是所有上层框架的调度决策,优先级非常高。如果你发现本地率长期上不去,第一个检查项就应该是这个。
5.2 为了本地化把等待时间调到 10 秒,结果更慢了
另一个反面教材是我自己犯过的错。当时为了把大任务的本地执行率从 70% 拉到 95%,我直接把spark.locality.wait调到了 10 秒,想着多点等待就能多点本地执行。结果那个大任务是跑快了,但整个集群同时段的其他任务全部变慢,整体 SLA 反而恶化了。
原因不复杂。这个大任务占用了大量 slot,等待本地位置的时候,这些 slot 一直空着,调度器不能把它们分配给其他任务。本来集群还有余力并行跑好几个小任务,结果所有资源都被大任务“悬空占住”,并行度直接掉了一半。后来的处理是,把全局等待时间调回 500ms,单独给超大任务设置更大的spark.locality.wait.node,只在确实值得等待的场景下才牺牲一点并行度。这个教训让我明白:本地化优化必须以集群整体吞吐为前提,不能为了某个任务好看,把整个集群拖下水。
5.3 存算分离下命中率低的真相:缓存不公用
前面提到我们在对象存储上加缓存亲和调度,中间也踩过命中率低的坑。初期缓存命中率只有 20% 左右,看起来像是淘汰策略不对,但排查后才发现,命中率低不是因为数据被淘汰,而是因为同一个数据分片在不同时刻被调度到了不同的 worker 上进行缓存。每个 worker 都有一份缓存,但谁都不完整,查询的时候这个 worker miss 了,去远端读一次,过一会儿另一个 worker 又 miss,又去远端读一次。相当于每份数据在集群里被重复拉取了很多遍。
这就是缓存没有亲和性的典型症状。后来我们把分片和 worker 的映射做成了固定的哈希绑定,查询调度跟着缓存分片走,命中率才明显升上来。这也是我说的“本地化对象从数据块变成缓存块”的核心差异:HDFS 的副本位置是系统按策略摆放的,你有办法感知;而对象存储上的缓存位置是你自己设计的,如果没有规整的映射关系,缓存就是一堆散落的碎片,毫无价值。
6. 一份可以直接抄作业的本地化配置与检查清单
6.1 几个关键参数,先对着查一遍
结合这些经验,整理一份我每次排查数据本地化问题时都会查的清单。不保证覆盖所有情况,但能解决绝大多数常见问题。
| 参数/配置 | 作用 | 推荐检查方式 |
|---|---|---|
| HDFS topology script | 正确上报节点机架信息 | 用hdfs dfsadmin -printTopology检查机架映射 |
dfs.client.read.shortcircuit | 开启本地短路读 | 确认 true 且 domain socket 路径两边一致 |
spark.locality.wait | 控制延迟调度等待时间 | 离线大任务 3s 以上,混合负载建议 500ms 级别 |
spark.locality.wait.node | 指定节点级等待时间 | 对大任务可单独调大,避免全局拉长 |
| 对象存储缓存映射策略 | 缓存分片与 worker 绑定 | 检查是否有固定哈希/一致性哈希映射 |
另外,如果你的集群用了 YARN 的节点标签功能,记得确认标签调度没有把执行器限制在错误节点上。节点标签和本地化是两个独立机制,但组合使用时很容易互相干扰,我在实际运维中遇到过好几次因为标签设置导致任务只能跑在非本地节点上的问题。
6.2 什么情况下才值得“死磕本地化”
不是所有场景都值得投入资源去优化数据本地化。我的判断标准一般看三点。第一,任务是不是数据密集型,也就是读取的数据量远大于计算量。如果是,本地化收益很大;如果是计算密集型的,效果就会打折。第二,底层存储是不是 HDFS 或本地盘缓存。如果数据全在对象存储上,那就按缓存亲和那套思路做,跟传统本地化已经不是一回事了。第三,任务有没有重复性。一次性临时查询,优化完下次可能就不跑了,不值得;每天固定跑批的任务,值得花大力气调。
另一个实用判断依据是看 Spark UI 上 Locality Level 的分布。如果 NODE_LOCAL 和 PROCESS_LOCAL 占比已经超过 80%,再努力提升空间也不大,不如把精力放到 SQL 或存储格式上。如果本地率长期低于 50%,那调度和副本这块非常值得优先排查。
6.3 可落地的监控指标与验收方法
优化做完了,如何确认有效?我自己习惯盯三个指标。一是本地执行率,也就是 Locality Level 中 PROCESS_LOCAL 和 NODE_LOCAL 的比例,目标一般设在 80% 以上。二是网络吞吐,如果集群网络一直被打满,说明数据移动量大,本地化改进后会明显下降。三是任务运行时间的稳定性,特别是 P95 和 P99。数据本地化做不好,任务时长会忽高忽低,因为远端读取的性能受网络拥塞影响很大。
验证方式也简单:挑一个固定的跑批任务,在改动前后各跑三次取平均时间。注意节点资源、并发度要基本一致,否则对比结果没有说服力。改完以后不要立刻下结论,观察至少一周,因为集群负载是周期性的,只有完整跨过一轮业务高峰才能看到真实效果。
再分享一个小技巧。我一般会把 Spark UI 上的 locality 分布信息定时采集到监控系统里,而不是等到任务变慢了才去翻 UI。这样每次调优后,都可以自动对比前后趋势,定位问题的时间能从小时级压缩到分钟级。毕竟在分布式系统里,性能问题是会“传染”的,早发现远比事后补救重要。