Python 高性能计算:Dask 与 Ray 分布式计算实战

Python 生态的分布式计算两巨头:Dask 与 Ray。本文系统讲解 Dask 的惰性图计算与 DataFrame/数组扩展、Ray 的 actor/task 并行模型、集群启动与资源调度、两者选型对比,以及大数据预处理、RL 训练、服务化部署的真实落地模式。

引言

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 与内存墙

Python 的两个硬约束:

约束影响解法
GIL多线程无法并行 CPU 密集多进程/分布式
内存墙单机内存装不下大数据分片/集群内存

两条路:

□ 向量化 + 进程并行(单机):numpy 多进程、multiprocessing
□ 分布式(集群):Dask / Ray —— 把数据与计算分到多台机器

记忆:Python 分布式解决两件事——GIL 挡住的 CPU 并行,单机装不下的大数据。


2. Dask 概览:惰性图与集合抽象

Dask 用「惰性图(Lazy Graph)」组织计算:调用 API 先构建任务图,真正 compute() 时才执行。

df.groupby(...).agg(...)   →  构建图(DAG)
df.compute()               →  分片并行执行

三大集合:

集合对应用途
dask.dataframePandas DataFrame分片 DataFrame
dask.arrayNumPy 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:选型决策

维度DaskRay
定位数据分析/DataFrame/数组通用分布式运行时
心智类 Pandas/NumPytask/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. 速查表与一句话记忆

需求工具
大数据 DataFrameDask 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]] — 高性能计算专题

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「hpc」更多文章

  1. ARM 超算与专用加速器:A64FX 与 NVIDIA Grace
  2. HPC 与 AI 融合:超算跑大模型训练
  3. 绿色 HPC:能耗优化与功率封顶实战