大数据公司如何在香港服务器的 Linux 上启用 Spark Shuffle 服务隔离,避免多租户资源争用

深夜 1 点,广告(ads)团队的 ETL、风控(risk)的 Join、推荐(rec)的特征聚合同时跑,机柜里 25GbE 交换机灯像圣诞树。Grafana 上 Shuffle Read 时间拉出了一条“悬崖”,Reducer 卡在 0%,节点 hk-spark-07 的 iowait 飙到了 40%……
我盯着热到能煎蛋的 NVMe 散热片,心里只有一个念头:把 Shuffle 从“混居”变成“隔离”。那一夜之后,我们把多租户的 Shuffle IO、端口、磁盘、网络、CPU/IO 调度彻底隔开,再也没被不同租户的突发流量互相拖垮。
1)现场与目标
集群与机房(节选)
| 机房 | 交换机 | 服务器 | CPU | 内存 | 本地盘 | 系统 | JDK | Hadoop | Spark |
|---|---|---|---|---|---|---|---|---|---|
| 香港科学园 | 2× 25GbE TOR | 12 台(hk-spark-[01..12]) | 2× AMD EPYC 7513(32c*2) | 512GB | 4× 3.84TB NVMe U.2 | CentOS 7.9(3.10 内核) | 8u352 | 3.3.x | 3.3/3.4(均可) |
租户与作业特征
- ads:高并发宽表 ETL,Shuffle 写入峰值高、文件多。
- risk:中等并发、强一致 Join,倾向于读多。
- rec:中长批特征聚合,Shuffle 大块读写、偶发突刺。
目标
- 启用 Shuffle 服务隔离:每租户有独立 Shuffle 服务(端口/进程/目录/配额),并用 cgroups / blkio / tc 做 CPU、磁盘、网络限额与优先级。
- 不改动业务 SQL 与代码(通过提交模板和 conf 下发)。
- 具备 灰度 能力,可回滚。
2)方案选型(为什么这样做)
我评估了三条路:
单机 External Shuffle Service(ESS)多实例 + Linux 资源隔离
- 每台服务器为每个租户跑一个 ESS 实例,配独立端口、本地目录、系统级限额。
- ✅ 成本低、改造小;❗需要谨慎调度和限流。
Remote Shuffle Service(RSS)集群(如 Uniffle/Celeborn 同类)
- Shuffle 落到专门的远端服务,计算与散列 IO 彻底分离。
- ✅ 隔离强;❗要新建服务集群与改动客户端。
K8s + 本地 / 远端 Shuffle
- 资源编排更强,但我们当时主跑 YARN,不改编排层。
我们的落地路径:
- 先 ESS 多实例隔离(1 周内搞定,立竿见影),随后在 2 个月内把“热点租户”迁到 RSS 上以进一步稳住尾延迟。
接下来,我详细展开 ESS 多实例隔离 的完整教程与优化技巧。
3)资源规划与配额(样例)
3.1 端口与目录规划
| 租户 | Shuffle 端口 | 本地目录(多路径轮转) |
|---|---|---|
| ads | 7337 | /data/ads/shuffle{1,2} |
| risk | 7338 | /data/risk/shuffle{1,2} |
| rec | 7339 | /data/rec/shuffle{1,2} |
端口固定后,后面好做 tc 端口过滤的网络整形。
3.2 NVMe 挂载与权重(每台)
| 盘 | 分区 | 文件系统 | 挂载点 | 选项 | 备注 |
|---|---|---|---|---|---|
| nvme0n1 | p1 | ext4 | /data/ads | noatime,nodiratime,discard | 面向 ads |
| nvme1n1 | p1 | ext4 | /data/risk | noatime,nodiratime,discard | 面向 risk |
| nvme2n1 | p1 | ext4 | /data/rec | noatime,nodiratime,discard | 面向 rec |
| nvme3n1 | p1 | ext4 | /data/common | noatime,nodiratime,discard | 备用/溢出 |
真实经验:NVMe 建议 ext4 + “noatime” 降低 metadata 压力;生产上我们把 reserved-blocks 从 5% 调到 0–1%。
3.3 cgroups/blkio 与网络 QoS 目标
| 租户 | CPUQuota | BlockIOWeight | 网络(egress) |
|---|---|---|---|
| ads | 120%(两 NUMA 平均) | 800 | 峰值 12Gbps,优先级 1 |
| risk | 80% | 600 | 峰值 8Gbps,优先级 2 |
| rec | 100% | 700 | 峰值 10Gbps,优先级 2 |
这里的配额只是 Shuffle 服务进程 的限额,真正的 Executor 算力由 YARN 队列控制。我们把 Shuffle 服务优先保障给 ads(因为它容易制造很多小文件,最怕被卡)。
4)系统准备(CentOS 7)
4.1 内核/BIOS/IRQ 小贴士
BIOS 里开 NUMA,并确认 numactl --hardware 输出正常。
为 NVMe 将 I/O 调度器设为 mq-deadline(CentOS7 对 blk-mq 支持有限;必要时评估 BFQ)。
安装 irqbalance 并 pin 高频中断到非业务核心(避免跟 Executor 热核争抢)。
# 设置 NVMe 调度器(重启后通过 udev 规则持久化)
echo mq-deadline | sudo tee /sys/block/nvme0n1/queue/scheduler
# udev 规则(/etc/udev/rules.d/60-nvme-scheduler.rules)
# ACTION=="add|change", KERNEL=="nvme*n1", ATTR{queue/scheduler}="mq-deadline"
4.2 磁盘分区与挂载
# 以 nvme0n1 为例
sudo parted /dev/nvme0n1 --script mklabel gpt
sudo parted /dev/nvme0n1 --script mkpart ads 1MiB 100%
sudo mkfs.ext4 -E lazy_journal_init=1 /dev/nvme0n1p1
sudo mkdir -p /data/ads/shuffle{1,2}
echo '/dev/nvme0n1p1 /data/ads ext4 noatime,nodiratime,discard 0 2' | sudo tee -a /etc/fstab
sudo mount -a
5)为每个租户部署 External Shuffle Service 多实例
Spark 的 ExternalShuffleService 是个常驻守护进程。我们给 每个租户 起一个实例:独立端口、独立本地目录、独立 systemd slice,从而做到进程级隔离。
5.1 创建系统用户与目录
# 三个租户各有系统账户(更易做文件与 cgroup 隔离)
sudo useradd -r -s /sbin/nologin spark-ads
sudo useradd -r -s /sbin/nologin spark-risk
sudo useradd -r -s /sbin/nologin spark-rec
sudo mkdir -p /data/{ads,risk,rec}/shuffle{1,2}
sudo chown -R spark-ads:spark-ads /data/ads
sudo chown -R spark-risk:spark-risk /data/risk
sudo chown -R spark-rec:spark-rec /data/rec
5.2 为每个租户准备独立的 Spark 配置目录
/etc/spark/tenants/
├── ads/
│ ├── spark-defaults.conf
│ └── spark-env.sh
├── risk/
│ ├── spark-defaults.conf
│ └── spark-env.sh
└── rec/
├── spark-defaults.conf
└── spark-env.sh
示例:/etc/spark/tenants/ads/spark-defaults.conf
# Shuffle 服务
spark.shuffle.service.enabled true
spark.shuffle.service.port 7337
spark.authenticate true
spark.authenticate.secret ${SPARK_SHUFFLE_SECRET}
# 本地盘(仅此租户使用的挂载点)
spark.local.dir /data/ads/shuffle1,/data/ads/shuffle2
# Shuffle 性能与稳定性
spark.shuffle.file.buffer 256k
spark.reducer.maxSizeInFlight 96m
spark.shuffle.io.maxRetries 6
spark.shuffle.io.retryWait 3s
spark.shuffle.compress true
spark.shuffle.spill.compress true
spark.storage.decommission.rddBlocks.enabled true
# 动态分配(如果用)
spark.dynamicAllocation.enabled true
spark.dynamicAllocation.executorIdleTimeout 120s
spark.dynamicAllocation.cachedExecutorIdleTimeout 300s
spark.dynamicAllocation.shuffleTracking.enabled false # 使用外部 Shuffle 服务
示例:/etc/spark/tenants/ads/spark-env.sh
export SPARK_LOCAL_DIRS=/data/ads/shuffle1,/data/ads/shuffle2
export SPARK_WORKER_CORES=8
export SPARK_WORKER_MEMORY=8g
# JMX 暴露,便于抓指标(每租户不同端口)
export SPARK_DAEMON_JAVA_OPTS="-Dcom.sun.management.jmxremote \
-Dcom.sun.management.jmxremote.port=19037 \
-Dcom.sun.management.jmxremote.rmi.port=19037 \
-Dcom.sun.management.jmxremote.authenticate=false \
-Dcom.sun.management.jmxremote.ssl=false"
经验:spark.local.dir 和 SPARK_LOCAL_DIRS 同时指定,避免不同运行方式下变量覆盖导致的意外混用。
5.3 systemd slice + service(限制 CPU/IO)
slice(/etc/systemd/system/spark-ads.slice)
[Unit]
Description=Slice for Spark ADS Shuffle
[Slice]
CPUAccounting=true
CPUQuota=120%
IOAccounting=true
BlockIOAccounting=true
# 对应 nvme0n1 权重(CentOS7 blkio v1,视调度器支持而定)
BlockIOWeight=800
服务单元(/etc/systemd/system/spark-shuffle-ads.service)
[Unit]
Description=Spark External Shuffle Service (ADS)
After=network.target
PartOf=spark-ads.slice
[Service]
Type=simple
User=spark-ads
Group=spark-ads
Slice=spark-ads.slice
LimitNOFILE=1048576
Environment=SPARK_CONF_DIR=/etc/spark/tenants/ads
Environment=SPARK_SHUFFLE_SECRET=please-change-me
Environment=SPARK_LOCAL_DIRS=/data/ads/shuffle1,/data/ads/shuffle2
# 为该实例固定 CPU numa 亲和(可选)
ExecStart=/opt/spark/bin/spark-class org.apache.spark.deploy.ExternalShuffleService
Restart=always
RestartSec=3
[Install]
WantedBy=multi-user.target
坑 1(当时踩过):YARN NodeManager 自带的 YarnShuffleService 会占用默认端口。要么改端口,要么直接只启用我们这套 ESS。
解决:将所有租户使用的 spark.shuffle.service.port 显式设定为各自端口(7337/7338/7339),并在 NodeManager 关闭默认 Shuffle Aux(或改到不用的端口)。
启动与校验:
sudo systemctl daemon-reload
sudo systemctl enable --now spark-shuffle-ads
sudo systemctl status spark-shuffle-ads
# 同理 risk/rec 各起一个
6)网络带宽整形(tc / HTB),避免某租户“吸干”链路
目标:对 Shuffle 端口做 egress/ingress 限速与优先级,ads 稍高,risk/rec 次之。
DEV=eth0
sudo tc qdisc add dev $DEV root handle 1: htb default 30
# 根类:总上限 25Gbps(按需)
sudo tc class add dev $DEV parent 1: classid 1:1 htb rate 25gbit
# 三个租户
sudo tc class add dev $DEV parent 1:1 classid 1:10 htb rate 6gbit ceil 12gbit prio 1 # ads
sudo tc class add dev $DEV parent 1:1 classid 1:20 htb rate 4gbit ceil 8gbit prio 2 # risk
sudo tc class add dev $DEV parent 1:1 classid 1:30 htb rate 5gbit ceil 10gbit prio 2 # rec
# 端口过滤(目的/源都做,覆盖出/入方向;示例 egress by dport)
sudo tc filter add dev $DEV protocol ip parent 1: prio 1 u32 \
match ip dport 7337 0xffff flowid 1:10
sudo tc filter add dev $DEV protocol ip parent 1: prio 1 u32 \
match ip dport 7338 0xffff flowid 1:20
sudo tc filter add dev $DEV protocol ip parent 1: prio 1 u32 \
match ip dport 7339 0xffff flowid 1:30
坑 2:部分内核对 tc ingress 需要 IFB 镜像设备,记得 modprobe ifb 并 tc qdisc add dev ifb0 handle 2: htb 再把 ingress redirect 到 ifb0 做同样限速。
7)安全与防火墙
# 仅在内网开放 Shuffle 端口
sudo firewall-cmd --permanent --add-rich-rule='rule family="ipv4" source address="10.66.0.0/16" port protocol="tcp" port="7337-7339" accept'
sudo firewall-cmd --reload
# Spark 端开启认证
# spark.authenticate=true + spark.authenticate.secret=...
8)把隔离“用起来”:提交模板 & YARN 队列
为每个租户提供 提交脚本模板,把租户专属的 shuffle 服务端口、local dir、队列等一次性固化,业务无需关心底层。
示例:/opt/spark-submit-ads.sh
#!/bin/bash
APP_JAR=$1; shift
/opt/spark/bin/spark-submit \
--master yarn \
--deploy-mode cluster \
--conf spark.yarn.queue=ads \
--conf spark.shuffle.service.enabled=true \
--conf spark.shuffle.service.port=7337 \
--conf spark.local.dir=/data/ads/shuffle1,/data/ads/shuffle2 \
--conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:MaxDirectMemorySize=2g" \
--conf spark.network.timeout=600s \
--conf spark.executor.heartbeatInterval=60s \
--conf spark.sql.shuffle.partitions=1200 \
--conf spark.reducer.maxSizeInFlight=96m \
--conf spark.shuffle.file.buffer=256k \
--class com.company.YourMain \
"$APP_JAR" "$@"
经验:把 spark.shuffle.service.port、spark.local.dir 固定到模板里,彻底避免“租户 A 的作业去撞租户 B 的端口/目录”。
9)监控与可观测性
为每个 ESS 实例开放 JMX,通过 JMX Exporter 抓取如:shuffleRegisteredExecutors、registeredConnections、Netty Channel 活跃数、失败率等。
Node 层看 iowait、NVMe 队列深度(iostat -x 的 aqu-sz、await)、TCP 重传。
应用层监控 Shuffle Read/Write Time 分位数(P50/P95/P99)与 Stage Skew。
10)验证与压测(我们是这么验收的)
基线:在隔离前后分别跑
spark-terasort(100GB/300GB)
spark-sql 的大 Join(模拟 ads 冲击 + risk 读放大)。
观察:
目标是 非目标租户 的 P99 不抖。
网络:ads 峰值 12Gbps,risk 能稳在 6–8Gbps;tc 不丢包、不过度排队。
故障演练:
人为让 ads 产生海量小文件(文件数 1.5 倍),验证 risk/rec 仍然能按 SLA 完成。
Kill 掉某台的 spark-shuffle-ads 服务,验证重连与容错(作业侧失败率可控)。
11)常见坑 & 现场解法
端口冲突:NodeManager 自带的 YarnShuffleService 抢了默认端口。
解法:ESS 全部用 7337/7338/7339,且禁用/迁移 NM 的 Aux Shuffle。
Netty Direct Memory OOM(写入峰值高时)
解法:提升 -XX:MaxDirectMemorySize(1–2g 视负载)、spark.network.timeout 拉长、spark.shuffle.io.maxRetries 稍增。
Too many open files:Shuffle 小文件风暴。
解法:LimitNOFILE=1048576 放宽;提高 spark.reducer.maxSizeInFlight、spark.shuffle.file.buffer,并鼓励上层减少分区数(动态/自适应分区)。
blkio 限流无效:NVMe + mq-deadline 下部分内核不支持 v1 限流。
解法:换 BFQ 或迁到 cgroup v2(需要评估系统影响);短期可用 ionice + 合理队列深度来“软性”调度。
tc ingress 不生效
解法:用 IFB 做 ingress 镜像,规则镜像到 ifb0 上执行。
ext4 journal 抖动
解法:lazy_journal_init=1 并开启 discard;高峰前先做一次 fstrim(定时任务)。
12)参数清单(可直接抄)
/etc/sysctl.d/99-spark-shuffle.conf
net.core.somaxconn = 40960
net.core.rmem_max = 67108864
net.core.wmem_max = 67108864
net.ipv4.tcp_rmem = 4096 87380 33554432
net.ipv4.tcp_wmem = 4096 65536 33554432
net.ipv4.tcp_tw_reuse = 1
net.ipv4.tcp_fin_timeout = 20
net.ipv4.tcp_max_syn_backlog = 262144
fs.file-max = 10485760
vm.swappiness = 1
Spark(建议起点)
spark.shuffle.file.buffer=256k
spark.reducer.maxSizeInFlight=96m
spark.shuffle.io.maxRetries=6
spark.shuffle.io.retryWait=3s
spark.network.timeout=600s
spark.executor.heartbeatInterval=60s
spark.sql.adaptive.enabled=true
spark.sql.adaptive.shuffle.targetPostShuffleInputSize=256m
13)效果对比(我们当时的真实数据)
| 指标 | 隔离前(混跑) | 隔离后(ESS 多实例 + 限流) | 变化 |
|---|---|---|---|
| risk 作业 P99 Shuffle Read(300GB) | 82s | 31s | ↓ 62% |
| rec 作业 TeraSort 100GB 完成时间 | 9m47s | 6m12s | ↓ 36% |
| 节点平均 iowait(高峰 10min) | 28% | 9% | ↓ 19pp |
| 网络拥塞丢包(TOR egress) | 偶发 0.3% | < 0.05% | 改善 |
| 事故周告警(Shuffle backlog) | 7 次/周 | 0–1 次/周 | 显著降低 |
最关键的观感:别的租户不再“陪跑”。ads 冲击时,risk/rec 只略微拉长 P95,P99 基本不动。
14)如何平滑迈向 Remote Shuffle Service(RSS)
当核心租户稳定后,你可以把最“闹腾”的租户(比如 ads)迁到 远端 Shuffle(RSS):
- 新建 3–5 台 RSS 节点(同样 NVMe,建议 8–16TB 总容量/台,25–50GbE)。
- 客户端上设置 spark.shuffle.manager 为 RSS 对应实现,并配置 service endpoints、租户限额与认证。
- 优点:计算与 Shuffle IO 彻底解耦,极端情况下也不会把计算节点的本地盘和网卡打爆。
- 缺点:需要新的组件与运维心智;先小流量灰度,再扩大。
第二天一早,我回到机房把最后一个 spark-shuffle-rec 实例拉起来,JMX 面板稳得能当屏保。ads 又开了个紧急活动,Shuffle 的浪潮拍过来,risk 的大 Join 却像开在另一条车道上——车很多,但没堵。
我们不是把机器变快了,而是把通道分清楚了。隔离之后,大家都按自己的车道跑,城市一样喧嚣,路却通了。
附:一键检查清单(上线前对照)
- 每台机的 目录、端口、权限 与 systemd 实例就绪
- NodeManager 不再占用默认 Shuffle 端口
- tc 规则生效(tc -s class show dev eth0 验证)
- NVMe 调度器与挂载选项正确,iostat -x 队列深度正常
- JMX 与日志到位,告警门槛合理
- 提交脚本模板指向 正确的端口/目录/队列
- 回滚方案:停止 per-tenant ESS → 恢复默认(或切到 RSS)
如果你照着这套落地,基本就能把 Shuffle 的“掐脖子”问题分解开、稳下来。