本节摘要:数据装不下内存时换大灶台——Dask 把 pandas 写法平移到分块世界,PySpark 把数据料理搬上集群。本节讲两间厨房的开启姿势:惰性执行与 compute 时机、DataFrame 算子的链式写法,以及"该不该换"的量级判断。承接 7.3 的调料柜,为中央厨房收官。
为什么同样的 groupby 脚本,百万行时一分钟,十亿行时直接内存爆掉?因为 pandas 的世界观是"全表进内存"——数据一旦越过内存,任何技巧都只是延缓。换大灶台才是正解:Dask 的思路是"把大表切成很多块,每块都是一张标准 pandas 表,分而治之";PySpark 的思路更进一步,"切块并搬到多台机器上同时做"。好消息是两间的灶位布置都刻意模仿 pandas——groupby 还是 groupby,merge 还是 merge,要换的主要是心智:从"立即执行"换成"先记账、后开工"。本节往前接 6.3 的量级阶梯,往后第 8 章的流水线实战会用到 Dask 的最小形态。
Dask:dask.dataframe.read_csv(路径, blocksize=...) 按块读入成惰性大表,blocksize 控制每块字节数;表被切成多个分区(partitions),head(5) 只碰第一块所以飞快;真正的计算要显式 .compute() 才落地成 pandas 表——记账与开工的界线就在这一刀;persist() 把反复使用的中途结果钉进内存,避免重复记账。PySpark:spark.read.csv(路径, header=True) 建惰性 DataFrame;select、filter、groupBy、agg、withColumn 全是算子,链式记的账到 show 或 write 才真正执行;join 的 broadcast 提示、repartition 的分桶调整是进阶旋钮,入门先记住"惰性"与"算子链"。
# Dask:pandas 用户的几乎无缝迁移 import dask.dataframe as dd ddf = dd.read_csv("big_orders_*.csv", blocksize="256MB") # 多个分块文件 print(ddf.columns.tolist()) # 只看结构,不触发计算 monthly = (ddf.assign(月份=ddf["日期"].str.slice(0, 7)) .groupby("月份")["金额"].sum()) print(monthly.compute()) # 此刻才真正开工
# PySpark:集群上的同一道菜 from pyspark.sql import SparkSession import pyspark.sql.functions as F spark = SparkSession.builder.appName("monthly").getOrCreate() sdf = spark.read.csv("big_orders.csv", header=True, inferSchema=True) result = (sdf.withColumn("月份", F.substring("日期", 1, 7)) .groupBy("月份").agg(F.sum("金额").alias("总额"))) result.show() # 触发执行
两段代码做同一件事:按月汇总金额。注意 Dask 里借了字符串切片、PySpark 里换了内置函数库(F.sum 而非 Python 的 sum)——大厨房有自己的厨具规范,越贴近原生算子,越不吃亏。
场景:日积月累的订单库到了内存的临界点,按量级阶梯做一次判断,然后执行最小迁移。
# 第一步:量级判断(6.3 的阶梯落地) # 估算内存:行数 × 每行字节 import pandas as pd sample = pd.read_csv("big_orders.csv", nrows=100000) # 先抽样看结构 per_row = sample.memory_usage(deep=True).sum() / len(sample) total_gb = per_row * 500_000_000 / 1024**3 print(f"估算全量内存:{total_gb:.1f} GB") # 结论:远超机器内存 -> 换 Dask;需要多机并行 -> 再上 PySpark # 第二步:迁移最小改动面 # 旧:df = pd.read_csv(...) 新:ddf = dd.read_csv(...) # 旧:df.groupby(...).sum() 新:ddf.groupby(...).sum().compute() # 中间步骤原样保留,只在"要结果"处补 compute
迁移的心法是"改动面最小化":读入与取结果两端换写法,中间的清洗整形聚合尽量原样——这正是 Dask 刻意与 pandas 对齐接口的用意。
**翻车一:小数据硬上大厨房。**千万行以内的任务,Dask 的分块调度开销可能比收益大——先确认真的"装不下",再换灶。**翻车二:compute 时机乱放。**循环里对每个分区各 compute 一次,账本反复重算,比单机还慢——记账到最后一刻,中途要复用就 persist。**翻车三:PySpark 里写 Python 函数。**对列 apply 自定义 Python 函数会把数据在 JVM 与 Python 间来回搬运,性能塌方——优先用内置函数库(F.sum、F.when、F.substring),Python 函数是最后手段。**翻车四:join 后没控制分区数。**两张大表 join 出的分区又多又斜,个别任务拖死全场——入门阶段至少知道"慢了先查分区",进阶再学 repartition 与广播提示。
量级"略超内存"的尴尬区间,还有更轻的中间档:6.3 的 chunksize 分块批处理,或把大 CSV 预先切成按月分块的多个小文件(Dask 的通配符读法天然欢迎这种布局);数据库里能完成的聚合,让 SQL 先算、pandas 只取结果;列存格式 parquet 比 CSV 省得多,Dask 与 PySpark 都吃——很多"内存不够"的案例,换存储格式就解了一半。真正的流式数据,批处理框架不接客,那是流处理引擎的地界。
中央厨房四间全部开门。下一章后厨动线:效率实测、内存算账、并行开灶、排错手册,最后整条流水线实战——见 8.1。