引言
Python 是数据科学主流,但单线程 GIL 与内存瓶颈让它在「大数据」前吃力。Dask 与 Ray 是 Python 生态里最主流的两套分布式框架:Dask 把 Pandas/NumPy 的接口扩展到「集群内存」,Ray 则提供通用的任务/actor 并行运行时。理解它们,就能在笔记本上写代码、在集群上跑千核。
本文系统对比两大框架:Dask 的惰性图执行与 DataFrame 扩展、Ray 的 task/actor 模型、集群部署与资源调度、选型决策,以及大规模预处理、强化学习训练、服务化部署三类落地模式。
前置:/hpc-mpi-basics/(分布式编程范式)、/hpc-slurm-scheduling/(集群作业)、/hpc-memory-hierarchy/(内存层次)。
目录
- 1. 为什么 Python 需要分布式:GIL 与内存墙
- 2. Dask 概览:惰性图与集合抽象
- 3. Dask 实战:DataFrame 与数组的分布式处理
- 4. Ray 概览:task 与 actor 并行模型
- 5. Ray 实战:并行函数与有状态服务
- 6. 集群部署与资源调度
- 7. Dask vs Ray:选型决策
- 8. 真实落地:数据预处理、RL 与服务化
- 9. 性能调优与常见坑
- 10. 速查表与一句话记忆
- 延伸阅读
1. 为什么 Python 需要分布式:GIL 与内存墙
Python 的两个硬约束:
| 约束 | 影响 | 解法 |
|---|---|---|
| GIL | 多线程无法并行 CPU 密集 | 多进程/分布式 |
| 内存墙 | 单机内存装不下大数据 | 分片/集群内存 |
两条路:
□ 向量化 + 进程并行(单机):numpy 多进程、multiprocessing
□ 分布式(集群):Dask / Ray —— 把数据与计算分到多台机器
记忆:Python 分布式解决两件事——GIL 挡住的 CPU 并行,单机装不下的大数据。
2. Dask 概览:惰性图与集合抽象
Dask 用「惰性图(Lazy Graph)」组织计算:调用 API 先构建任务图,真正 compute() 时才执行。
df.groupby(...).agg(...) → 构建图(DAG)
df.compute() → 分片并行执行
三大集合:
| 集合 | 对应 | 用途 |
|---|---|---|
dask.dataframe | Pandas DataFrame | 分片 DataFrame |
dask.array | NumPy array | 分块多维数组 |
dask.bag | 泛型集合 | 无结构数据处理 |
核心心智:
□ 接口像 Pandas/NumPy,但数据分片
□ 惰性执行:链式 API 只是建图,compute 才跑
□ 图自动并行:一个图可拆到多 worker
记忆:Dask = 「Pandas/NumPy 的集群版」——接口不变,图自动并行。
3. Dask 实战:DataFrame 与数组的分布式处理
import dask.dataframe as dd
# 读取大 CSV(分片)
df = dd.read_csv("s3://bucket/data/*.csv", blocksize="128MB")
# 链式处理(惰性)
result = (df[df.status == "ok"]
.groupby("user_id")
.sales.sum())
# 触发执行
top = result.nlargest(10).compute()
print(top)
分布式数组:
import dask.array as da
x = da.random.normal(0, 1, size=(100000, 100000), chunks=(1000, 1000))
y = da.linalg.svd(x).compute() # 分块 SVD
调度器选择:
□ 本机:ThreadsScheduler / ProcessesScheduler
□ 集群:Distributed(推荐,看得到仪表盘)
from dask.distributed import Client
client = Client("tcp://scheduler:8786") # 连接集群
注意:
compute()才真正执行——新手最常见的错是「没 compute 就 print」拿到惰性对象。
4. Ray 概览:task 与 actor 并行模型
Ray 提供通用分布式运行时,两大原语:
Task(任务):无状态函数并行执行:
import ray
@ray.remote
def square(x: int) -> int:
return x * x
futures = [square.remote(i) for i in range(100)]
results = ray.get(futures) # 并行计算 100 个平方
Actor(有状态服务):跨调用保持状态的对象:
@ray.remote
class Counter:
def __init__(self):
self.n = 0
def inc(self):
self.n += 1
return self.n
counter = Counter.remote()
print(ray.get(counter.inc.remote())) # 1
Ray 的生态:
□ Ray Core:task/actor 基础
□ Ray Data:大数据处理(并行数据管道)
□ Ray Train:分布式训练(PyTorch/LLM)
□ Ray Serve:模型服务化部署
□ Ray RLlib:强化学习库
记忆:Ray = 「通用分布式运行时」——task 跑函数、actor 跑服务,生态覆盖数据处理到 RL 与 Serve。
5. Ray 实战:并行函数与有状态服务
并行函数(task):
@ray.remote(num_cpus=2) # 指定资源
def process_chunk(chunk: list) -> float:
# 模拟 CPU 密集计算
return sum(map(lambda x: x**2, chunk))
data = list(range(1000))
chunks = [data[i:i+100] for i in range(0, 1000, 100)]
futures = [process_chunk.remote(c) for c in chunks]
total = sum(ray.get(futures))
print(total)
有状态 actor(共享模型/状态):
@ray.remote
class ModelServer:
def __init__(self, weights):
self.weights = weights
def predict(self, x):
return x * self.weights
server = ModelServer.remote(2.0)
futures = [server.predict.remote(i) for i in range(10)]
print(ray.get(futures)) # [0.0, 2.0, 4.0, ...]
并行控制流:
□ ray.get:阻塞取结果
□ ray.wait:等待部分完成
□ object refs:结果引用,可传递
□ futures 的依赖:自动任务图调度
6. 集群部署与资源调度
Dask 集群部署:
# 方式一:dask scheduler + workers(手动/脚本)
dask-scheduler --port 8786 &
dask-worker tcp://scheduler:8786 --nprocs 4 --nthreads 2
# 方式二:配合 Slurm(作业脚本)
# sbatch 里启动分布式集群
Ray 集群部署:
ray start --head --port=6379 # 主节点
ray start --address=192.168.1.10:6379 # 工作节点加入
# 或用 Ray Cluster Launcher(Kubernetes / Slurm 模板)
资源调度要点:
□ 指定资源:num_cpus / num_gpus / memory(task 与 actor 都可)
□ 自动缩放:Dask 的 Adaptive Scaling、Ray autoscaler
□ 与 Slurm 集成:作业内启动分布式集群,作业结束统一回收
□ 监控:Dask 仪表盘、Ray Dashboard
提示:在超算上,通常一个 Slurm 作业内启动 Dask/Ray 集群——作业调度与分布式调度两层配合。
7. Dask vs Ray:选型决策
| 维度 | Dask | Ray |
|---|---|---|
| 定位 | 数据分析/DataFrame/数组 | 通用分布式运行时 |
| 心智 | 类 Pandas/NumPy | task/actor |
| 大数据处理 | 强(分片 DataFrame) | 有(Ray Data) |
| 模型训练 | 弱 | 强(Train/RLlib) |
| 服务化 | 弱 | 强(Serve) |
| 学习曲线 | 平缓 | 中等 |
选型建议:
□ 数据分析/ETL/聚合 → Dask(接口熟悉、上手快)
□ 分布式训练/RL/服务化 → Ray(生态完整)
□ 两者都要 → Ray Data 可兼容部分 Dask 用法
□ 已有 Slurm 超算 → Dask 集成成熟,Ray 也支持
记忆:「算数据」用 Dask、「训模型跑服务」用 Ray——分析向 Dask,应用向 Ray。
8. 真实落地:数据预处理、RL 与服务化
模式一:大规模数据预处理(Dask):
读 TB 级原始数据 → 分片清洗/特征工程 → 写回 parquet → 训练
features = (dd.read_parquet("s3://raw/*.pq")
.assign(is_new=lambda d: d.created_at > cutoff)
.compute())
模式二:分布式 RL 训练(Ray RLlib):
from ray.rllib.algorithms.ppo import PPOConfig
config = PPOConfig().environment("CartPole-v1").rollouts(num_rollout_workers=4)
algo = config.build_algo()
for _ in range(10):
algo.train()
模式三:模型服务化(Ray Serve):
from ray import serve
@serve.deployment
def echo(request):
return {"msg": request.query_params["text"]}
serve.run(echo.bind(), host="0.0.0.0", port=8000)
# curl "http://host:8000/?text=hi"
落地心法:「分片处理数据、并行训模型、无缝服务化」——Dask 做上游清洗、Ray 做训练与 Serving,是一条完整流水线。
9. 性能调优与常见坑
Dask 调优:
□ 分片大小:128MB 左右(太大拖慢 shuffle,太小调度开销)
□ shuffle 性能:groupby/join 是大头,评估分区策略
□ 避免小文件爆炸:coalesce 分区、写 parquet 前合并
Ray 调优:
□ actor 数量与并发:别过度并行,资源浪费
□ 序列化开销:小对象大量传递慢 → 用 object store/参考批处理
□ GIL 场景:CPU 密集用多进程(num_cpus),IO 密集用多线程
常见坑:
□ 忘 compute()(Dask)
□ 单机模式当集群用(没启动 scheduler)
□ 大对象反复复制(没走 object store)
□ 不设资源 → 过度调度 OOM
□ 未关闭集群(ray.shutdown / client.close)
记忆:分布式调优三问——数据分多大、并行度多少、对象怎么传;最常见事故是「忘 compute」与「资源不设就乱跑」。
10. 速查表与一句话记忆
| 需求 | 工具 |
|---|---|
| 大数据 DataFrame | Dask DataFrame |
| 分布式数组 | Dask Array |
| 通用并行函数 | Ray task |
| 有状态服务 | Ray actor |
| 分布式训练 | Ray Train/RLlib |
| 模型服务化 | Ray Serve |
| 集群启动 | dask-scheduler / ray start |
| 资源指定 | num_cpus / num_gpus / memory |
| 数据分析选型 | Dask |
| 应用/训练选型 | Ray |
一句话记忆:Python 分布式 = Dask 管「大数据分析」(接口像 Pandas、图自动并行、记得 compute),Ray 管「并行计算与应用」(task 跑函数、actor 跑服务、Train/Serve 全家桶)——分析选 Dask、训练与服务选 Ray。
延伸阅读
- /hpc-mpi-basics/ — 底层分布式编程范式
- /hpc-slurm-scheduling/ — 集群作业调度与 Dask/Ray 结合
- /hpc-memory-hierarchy/ — 内存层次与大数据分片
- /hpc-performance-profiling/ — 剖析分布式瓶颈
- /hpc-cluster-admin/ — 集群运维与资源管理
- [[hpc]] — 高性能计算专题
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。