人人都会AI编程

20.5 大数据处理:PySpark、Dask 分布式计算基础

更新时间:2026-07-12

前面几节我们围绕 NumPy 和 Pandas 做了大量单机数据分析,但当数据量增长到几十 GB、几百 GB,甚至 TB 级别时,单台机器的内存和 CPU 就不够用了。这时我们需要将计算拆分到多台机器或多核 CPU 上并行处理。Python 生态里有两个主流的大数据处理框架:PySparkDask

它们都提供类似 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 上并行,也可以扩展到多台机器的分布式集群。

核心特性

  • 熟悉的 APIdask.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 分析从“单机”跃迁到“集群”的能力,真正覆盖现代数据工程的主流需求。