Dask 并行计算替代 Pandas

FreeGuideOnline 最新 2026-07-09

bash pip install dask[complete]

或最小化安装

pip install dask


如果只需 DataFrame 部分,可单独安装:

```bash
pip install dask[dataframe]

核心概念:惰性计算与任务图

Dask 采用 惰性执行——构建计算图而非立即计算。调用 API 时只生成任务图,直到调用 .compute() 才真正触发计算。这允许 Dask 对整个执行过程优化并并行化。

对于用过 Pandas 的用户,记住这一条黄金法则:Dask 操作 = 构建蓝图,.compute() = 开工执行

Dask DataFrame:Pandas 的超集

Dask DataFrame 由许多按索引划分的 Pandas DataFrame 组成,称为 分区(partition)。每个分区是一个独立的 Pandas 对象,因此任何 Pandas 操作都能在分区上应用,Dask 负责协调并行。

创建 Dask DataFrame

从 Pandas DataFrame 转换:

import pandas as pd
import dask.dataframe as dd

pdf = pd.DataFrame({"x": range(100), "y": range(100, 200)})
ddf = dd.from_pandas(pdf, npartitions=3)  # 切分3个分区

从 CSV/Parquet 等文件读取:

# 读取多个 CSV 文件
ddf = dd.read_csv("data/*.csv")

# 读取 Parquet(推荐格式)
ddf = dd.read_parquet("data/")

按通配符读取时,每个文件成为一个(或多个)分区,天然并行。

与 Pandas API 的对比

操作 Pandas Dask DataFrame
读取 CSV pd.read_csv('file.csv') dd.read_csv('file*.csv')
显示前几行 df.head() ddf.head()
列选择 df[['col1','col2']] ddf[['col1','col2']]
过滤 df[df.col1 > 0] ddf[ddf.col1 > 0]
分组聚合 df.groupby('key').sum() ddf.groupby('key').sum().compute()
连接(merge) pd.merge(df1, df2, on='key') ddf.merge(ddf2, on='key')
写出结果 df.to_csv('out.csv') ddf.to_csv('out-*.csv')

注意:Dask 的大多数操作返回一个新的延迟对象,需调用 .compute() 得到 Pandas 或 NumPy 结果。

从 Pandas 迁移的实用指南

1. 替换读取和写入

pd.read_csv 改为 dd.read_csv,允许通配符。写出时指定 to_parquetto_csv('out-*.csv') 生成分区文件。

2. 坚持惰性计算链

尽量将数据清理和转换写成连续的操作链,仅在最终需要结果时调用 .compute() 或保存文件。

clean = (
    ddf[ddf.price > 0]
    .groupby("category")
    .price.mean()
)
result = clean.compute()   # 现在才执行所有任务

3. 注意无法直接“查看”全部数据

ddf.shape 返回的是延迟值,可使用 ddf.shape.compute()len(ddf) 得到实际值。用 .head() 快速预览,但不要直接打印整个 Dask DataFrame(会触发计算并可能内存溢出)。

4. 自定义函数需适应分区

对 Dask DataFrame 使用 apply 时,函数在单个分区上运行,输入输出都是 Pandas DataFrame/Series。可使用 meta 参数指定输出结构:

ddf.map_partitions(lambda pdf: pdf.assign(z=pdf.x * 2))

并行计算实战案例

案例1:超大数据集的聚合

import dask.dataframe as dd

ddf = dd.read_parquet("s3://my-bucket/transactions/")
result = ddf.groupby("user_id").amount.sum().compute()

案例2:多文件合并与转换

ddf = dd.read_csv("logs/2024-*.csv", parse_dates=["timestamp"])
ddf["hour"] = ddf.timestamp.dt.hour
hourly_stats = ddf.groupby("hour").size().compute()

案例3:分布式环境(本地集群模拟)

from dask.distributed import Client

client = Client()   # 启动本地分布式调度器,使用所有CPU核心
# 之后所有dask操作自动提交到集群
ddf = dd.read_csv("big_data/*.csv")
ddf.groupby("key").value.mean().compute()