Dask 并行计算替代 Pandas
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_parquet 或 to_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()