上一篇 下一篇 分享链接 返回 返回顶部

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

发布人:Minchunlin 发布时间:2025-09-14 10:50 阅读量:585


深夜 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 的“掐脖子”问题分解开、稳下来。

目录结构
全文