前面几节我们围绕 NumPy 和 Pandas 做了大量单机数据分析,但当数据量增长到几十 GB、几百 GB,甚至 TB 级别时,单台机器的内存和 CPU 就不够用了。这时我们需要将计算拆分到多台机器或多核 CPU 上并行处理。Python 生态里有两个主流的大数据处理框架:PySpark 和 Dask。
它们都提供类似 DataFrame 的操作接口,让你用熟悉的语法处理超越单机内存的数据集。
PySpark:Apache Spark 的 Python 接口
PySpark 是 Apache Spark 的官方 Python API。Spark 是一个通用分布式计算引擎,支持在成百上千台机器上并行处理数据。通过 PySpark,你可以在 Python 里调用 Spark 的完整能力。
核心概念
- RDD(弹性分布式数据集):Spark 最基础的抽象,代表一个可分区的、可容错的分布式数据集合。现在大部分场景更推荐直接使用更高层的 DataFrame API。
- DataFrame:类似于 Pandas 的二维表结构,但数据分布在集群的多台节点上,支持类 SQL 操作和列式存储,性能更好。
- Lazy Evaluation(惰性计算):转换操作(如
select,filter,groupBy)并不会立即执行,只是构建执行计划;只有遇到action(如count,collect,show)时才会真正触发集群计算。这让 Spark 可以整体优化执行流程。
快速上手示例
from pyspark.sql import SparkSession
# 1. 创建 SparkSession(入口)
spark = SparkSession.builder \
.appName("demo") \
.getOrCreate()
# 2. 读取一个大 CSV 文件(来自本地、HDFS、S3 等)
df = spark.read.csv("data/large_file.csv", header=True, inferSchema=True)
# 3. 类似 SQL 的操作
df.filter("age > 30") \
.groupBy("city") \
.avg("salary") \
.show()
适用场景
- 真正的分布式集群处理,数据存储在 HDFS、S3、Hive 等。
- 需要与其他大数据组件(如 Kafka、HBase、Spark SQL)集成。
- 有数据工程师团队维护 Spark 集群。
需要注意的点
- 环境搭建比单机库复杂,需要 Java 环境和集群配置。
- 学习曲线相对较陡,需要理解分区、shuffle、序列化等概念。
- 本地模式可以用
spark-submit或 PyCharm 直接运行,适合小数据量调试。
Dask:更 Pythonic 的并行计算库
Dask 是一个纯 Python 库,设计的初衷就是让熟悉的 NumPy/Pandas 代码在遇到大规模数据时能无缝并行化,而无需学习新的框架或生态。它既可以在单机的多核 CPU 上并行,也可以扩展到多台机器的分布式集群。
核心特性
- 熟悉的 API:
dask.dataframe几乎完整克隆了 Pandas 接口,dask.array模仿 NumPy 接口。你只需少修改代码,就能处理超出内存的数据。 - Lazy Execution(惰性执行):和 Spark 类似,操作不会立即执行,而是构建任务图,调用
.compute()时才并行计算。 - 自动分块:DataFrame 被自动切分成多个 Pandas DataFrame 分片,每次计算只加载必要的分片到内存。
快速上手示例
import dask.dataframe as dd
# 1. 读取大 CSV,自动分成多个分区
df = dd.read_csv("data/large_file-*.csv")
# 2. 与 Pandas 几乎一样的操作
result = df[df.age > 30].groupby("city").salary.mean()
# 3. 触发计算(此时才会并行执行)
print(result.compute())
单机多核与集群
- 单机模式(默认):利用多核 CPU 并行处理,数据可以溢出到磁盘,突破内存限制。
- 分布式模式:通过
dask.distributed可以在多台机器上部署 Scheduler 和 Worker,组成一个轻量级计算集群。
适用场景
- 你手头已经有一套成熟的 Pandas 分析代码,现在数据量变大了一两个数量级。
- 需要在单台高性能服务器上充分利用多核资源,不用搭建复杂集群。
- 偏爱 Python 原生生态,希望少引入外部依赖。
- 数据科学实验、特征工程、ETL 过程的快速扩展。
需要注意的点
- 并非所有 Pandas 操作都完美支持,某些高阶方法可能需要通过
map_partitions写自定义函数。 - 对内存占用的控制需要合理设置分区大小和数量。
- 分布式模式下的网络通信会增加一定开销。
选型建议
| 维度 | PySpark | Dask |
|------------------|--------------------------------------------|--------------------------------------------|
| 生态系统 | 大数据组件集成广泛,有 Spark SQL/MLlib | 纯 Python 生态,与 NumPy/Pandas 深度绑定 |
| 学习成本 | 较高,需要 JVM + Spark 知识体系 | 较低,Pandas 用户几乎无缝切换 |
| 部署复杂度 | 需要集群管理(Yarn/K8s 等) | 单机直接 pip install 即可,集群也较简单 |
| 适用数据规模 | 数百 TB 至 PB 级 | 数十 GB 至 TB 级(单机多核或小集群) |
| 典型用户 | 数据工程师、大数据平台 | 数据分析师、数据科学家、Python 全栈 |
实战建议:
- 如果数据已经在大数据平台里(HDFS、Hive),且团队有 Spark 基础设施,用 PySpark 是自然的选择。
- 如果只是现有 Pandas 脚本跑不动了,数据量还在单机硬盘可存储的范围内,优先试试 Dask,因为迁移成本最低,开发和调试速度快。
- 两者并不是非此即彼的关系,很多团队同时使用:用 Dask 做探索和不频繁的实验,用 PySpark 跑定期生产任务。
掌握了 PySpark 或 Dask 的基础,你就具备了让 Python 分析从“单机”跃迁到“集群”的能力,真正覆盖现代数据工程的主流需求。