Windows 本地 PySpark 入门(环境搭建 + 基础操作)
📖 目录
- 适用场景
- 核心原理
- PySpark 是什么
- 为什么 Windows 需要 winutils
- 懒执行:transformation 与 action
- 具体做法
- 1. 安装 JDK 和 Python
- 2. 安装 PySpark
- 3. 下载 winutils
- 4. 配置环境变量
- 5. 最小验证 test.py
- 6. 入门代码骨架:建会话 + 建 DataFrame
- 7. 五个基础操作
- 8. 写入结果
- 决策辅助
- 版本匹配表
- 报错速查表
- 操作选择表
- 实战案例
- 逐段运行与输出对照
- 映射表:代码行 ↔ 概念 ↔ 输出
- 安全锁清单
- 用 Chrome 下载被拦截怎么办
- 为什么代码里还要再设一遍环境变量
- 为什么写了 filter 却没有输出
- 数据写不出来或版本对不上怎么办
- 要不要下载完整 Hadoop 发行版
- 进阶方向
适用场景
这篇文档写给想在 Windows 电脑上把 PySpark 跑起来、并学会用 DataFrame 做入门级数据处理的同学。
关键判断条件:
- 系统是 Windows
- 本地单机学习/开发,不涉及集群
- Python 环境可用,能
pip install
反例(什么时候不需要下面这些步骤):
- 用的是 Docker / WSL / Linux 环境 → Hadoop 有官方 Linux 版,不需要 winutils
- 已经有远程 Spark 集群 → 直接连集群即可,不用本地补丁
- 只在 notebook 里看数据、不落盘 → 大多场景连 winutils 都可以省
判断口诀:"Windows 本地 + 要跑通 + 要落盘" → 走本文全流程;否则按需跳过 winutils 环节。
核心原理
PySpark 是什么
30 秒本质:PySpark 是 Spark 的 Python 接口。你的 Python 代码通过 Py4J 桥接,驱动一个运行在 JVM 里的 Spark 干活;数据操作的对象是 DataFrame——一张分布式的表。
flowchart LR
PY["你的 Python 代码"] -->|Py4J| DRV["JVM 里的 Spark Driver"]
DRV --> EX["Executor"]
EX -->|Hadoop API| HW["Hadoop 本地库"]
HW -->|winutils.exe| WIN["Windows 文件系统"]
图:PySpark 代码到最终落盘的调用链——Python 只是壳,真正干活的是 JVM 里的 Spark,最底层通过 Hadoop 访问文件系统。
为什么 Windows 需要 winutils
Spark 落盘、读文件都要走 Hadoop 的文件系统 API,但 Hadoop 官方只发布 Linux 版。所以 Windows 上要么报 Could not locate winutils.exe,要么在深层文件操作时报 UnsatisfiedLinkError: access0。
- 朴素做法(不推荐):只
pip install pyspark就开跑 → 报错。缺的不是 pyspark 包,而是 Windows 上没有的 Hadoop 本地库。 - 优化做法:补一个社区编译的 Windows 适配补丁(winutils.exe + hadoop.dll),再让 Spark 找到它——这就是环境搭建的全部意义。
懒执行:transformation 与 action
DataFrame 的操作分两类,这是入门最容易懵的地方:
| 类型 | 特点 | 例子 |
|---|---|---|
| transformation(转换) | 懒执行,只是"记下要干什么",不真正算 | select / filter / withColumn / groupBy |
| action(动作) | 触发计算,真正跑出结果 | show() / count() / write |
flowchart LR
A["df.select(...)"] --> B["df.filter(...)"]
B --> C["df.withColumn(...)"]
C --> D["df.show()"]
D --> E["触发计算并输出结果"]
图:一连串 transformation 只是搭流水线,最后一个 action 才真正执行。
一句话记忆:transformation 记账,action 结账。 写了
df.filter(...)没看到输出,不是错了,是没"结账"——补一个.show()即可。
具体做法
本节分两块:前 5 步把环境跑通,6~8 步是 PySpark 入门的代码骨架(完整可跑版见仓库 env/apps/learn_spark.py)。
flowchart TD
A["安装 JDK + Python"] --> B["pip 安装 PySpark"]
B --> C["下载 winutils 补丁"]
C --> D["配置环境变量"]
D --> E["跑 test.py 验证"]
图:环境搭建五步流程。
1. 安装 JDK 和 Python
- JDK:安装 OpenJDK 17(或 11),双击安装,记住安装路径(如
C:\Program Files\Java\jdk-17),后面环境变量要用。 - Python:安装 Python 3.8+ 或 Miniconda,安装时勾选 Add Python to PATH。
2. 安装 PySpark
pip install pyspark==3.5.0
PySpark 自带 Spark,不需要单独下载 Spark 发行版。装完用
pip show pyspark可查版本和内置 Hadoop 版本。
3. 下载 winutils
目标:把 winutils.exe 和 hadoop.dll 放进一个统一目录。下文以 D:\tools\hadoop-3.3.6\bin 为例,目录名要和你的 PySpark 内置 Hadoop 版本对应(如 3.3.6)。
⚠️ 用 Edge 浏览器下载:Chrome 会把这两个文件识别为"危险文件"直接拦截,Edge 不会。
方法一:git clone 整个仓库(推荐)
git clone https://github.com/cdarlint/winutils.git
克隆完成后,进入 winutils/hadoop-3.3.6/bin/,把里面的文件全部复制到目标目录:
mkdir -p D:/tools/hadoop-3.3.6/bin
cp winutils/hadoop-3.3.6/bin/* D:/tools/hadoop-3.3.6/bin/
这样 winutils.exe 和 hadoop.dll 一次拿全,不用分别下载。
方法二:DownGit 打包下载
打开 DownGit 链接,点 Download,解压后把 bin 目录内容复制到 D:\tools\hadoop-3.3.6\bin:
https://minhaskamal.github.io/DownGit/#/home?url=https://github.com/cdarlint/winutils/tree/master/hadoop-3.3.6/bin
备选仓库:steveloughran/winutils(目录结构相同,可同样操作)。
4. 配置环境变量
右键"此电脑" → 属性 → 高级系统设置 → 环境变量 → 系统变量:
| 变量 | 示例值 | 作用 |
|---|---|---|
JAVA_HOME |
C:\Program Files\Java\jdk-17 |
Spark 找 Java |
HADOOP_HOME |
D:\tools\hadoop-3.3.6 |
Spark 找 winutils / hadoop.dll |
Path 追加 |
%JAVA_HOME%\bin、%HADOOP_HOME%\bin |
让 java、winutils 命令可用 |
新开一个 PowerShell 窗口验证:
java -version # 能输出版本号
winutils # 输出 Usage 帮助即成功
5. 最小验证 test.py
新建 env/apps/test.py,环境变量必须在 import pyspark 之前设置:
import os
os.environ["HADOOP_HOME"] = "D:/tools/hadoop-3.3.6"
os.environ["PATH"] = os.environ["PATH"] + ";D:/tools/hadoop-3.3.6/bin"
# Windows 上 Spark 默认找 python3,实际命令是 python
os.environ["PYSPARK_PYTHON"] = "python"
os.environ["PYSPARK_DRIVER_PYTHON"] = "python"
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("LocalTest") \
.getOrCreate()
df = spark.range(10)
print("总行数:", df.count())
spark.stop()
运行:
python env/apps/test.py
输出 总行数: 10 即环境跑通。
6. 入门代码骨架:建会话 + 建 DataFrame
不用每次从零写,记住这个骨架即可(完整版见 learn_spark.py):
from pyspark.sql import SparkSession
# 1. 建会话(Spark 程序的"入口")
spark = SparkSession.builder \
.appName("MyLearning") \
.config("spark.driver.extraLibraryPath", "D:/tools/hadoop-3.3.6/bin") \
.config("spark.executor.extraLibraryPath", "D:/tools/hadoop-3.3.6/bin") \
.getOrCreate()
# 2. 造一张表:列名 + 行数据
columns = ["Name", "Age", "Job"]
data = [("Alice", 25, "Engineer"),
("Bob", 30, "Data Scientist"),
("Cathy", 28, "Product Manager")]
df = spark.createDataFrame(data, columns)
df.show()
要点:
.config("spark.driver.extraLibraryPath", ...)让 JVM 找到hadoop.dll,Win 环境必备;createDataFrame(data, columns)把 Python 列表变成 DataFrame(分布式表);- 每个
spark会话对应一个程序,用完spark.stop()。
7. 五个基础操作
| 操作 | 作用 | 示例 | 注意 |
|---|---|---|---|
select |
选列 | df.select("Name", "Age") |
返回新 DataFrame |
filter |
按条件筛行 | df.filter(df["Age"] > 25) |
返回新 DataFrame |
withColumn |
加/改一列 | df.withColumn("Age+10", df["Age"] + 10) |
返回新 DataFrame |
groupBy |
分组聚合 | df.groupBy("Job").avg("Age") |
聚合后列名自动带 avg(...) |
| SQL | 当表来查 | df.createOrReplaceTempView("people") + spark.sql("...") |
先注册视图 |
三个"返回新 DataFrame"的操作都不会改原表——想看结果,记得补
.show()。
8. 写入结果
df.write.mode("overwrite").parquet("output/people.parquet")
df.write.mode("overwrite").option("header", "true").csv("output/people_csv")
mode("overwrite"):目录已存在时覆盖,防止重复跑报错;option("header", "true"):CSV 带表头;write是 action,会真正触发计算并把数据落盘。
决策辅助
版本匹配表
| Spark | Hadoop | JDK | PySpark |
|---|---|---|---|
| 3.5.0 | 3.3.6 | 11 / 17 | 3.5.0 |
| 3.4.0 | 3.3.4 | 11 / 17 | 3.4.0 |
| 3.3.0 | 3.3.0 | 8 / 11 | 3.3.0 |
口诀:Spark 与 PySpark 同号,winutils 目录名跟 Hadoop 走。
pip show pyspark可确认内置 Hadoop 版本。
报错速查表
环境类报错
| 报错 | 原因 | 解决 |
|---|---|---|
java not found |
JAVA_HOME 未设置或路径错 | 安装 JDK,检查 JAVA_HOME |
Cannot run program "python3" |
Windows 上叫 python 不叫 python3 | 设 PYSPARK_PYTHON=python |
UnsatisfiedLinkError: access0 |
hadoop.dll 没被 JVM 加载 | 确认文件在 bin 目录,检查 extraLibraryPath |
Could not locate winutils.exe |
HADOOP_HOME 不对或 bin 不在 Path | 检查 HADOOP_HOME 和 Path |
Python worker failed to connect back |
Python 子进程连不上 JVM | PYSPARK_PYTHON 用 python 绝对路径 |
代码类报错
| 报错 | 原因 | 解决 |
|---|---|---|
写了 filter 却没输出 |
transformation 是懒执行 | 后面补 .show() |
| 重复写目录报错 | 目录已存在 | 加 mode("overwrite") |
操作选择表
想做什么 → 用哪个 API:
| 需求 | 用哪个 | 一句话 |
|---|---|---|
| 只要某几列 | select |
选列 |
| 筛掉不满足条件的行 | filter |
筛行 |
| 基于旧列算新列 | withColumn |
加列 |
| 按类目求均值/计数 | groupBy + avg / count |
分组 |
| 想写熟悉的 SQL | createOrReplaceTempView + spark.sql |
变表来查 |
| 把结果存下来 | write.mode(...).parquet / .csv |
落盘 |
口诀:选列 select、筛行 filter、加列 withColumn、分组 groupBy、想看结果 show。
实战案例
用仓库里的 env/apps/learn_spark.py 完整跑一遍。环境配置部分和 test.py 一样,从建会话开始逐段看:
逐段运行与输出对照
① 建会话 + 造表
spark = SparkSession.builder \
.appName("MyLearning") \
.config("spark.driver.extraLibraryPath", "D:/tools/hadoop-3.3.6/bin") \
.config("spark.executor.extraLibraryPath", "D:/tools/hadoop-3.3.6/bin") \
.getOrCreate()
data = [("Alice", 25, "Engineer"),
("Bob", 30, "Data Scientist"),
("Cathy", 28, "Product Manager")]
columns = ["Name", "Age", "Job"]
df = spark.createDataFrame(data, columns)
② 看全表 + 结构(show 触发计算)
df.show()
+-----+---+---------------+
| Name|Age| Job|
+-----+---+---------------+
|Alice| 25| Engineer|
| Bob| 30| Data Scientist|
|Cathy| 28|Product Manager|
+-----+---+---------------+
df.printSchema()
root
|-- Name: string (nullable = true)
|-- Age: long (nullable = true)
|-- Job: string (nullable = true)
③ 选列、筛行、加列
df.select("Name", "Age").show()
+-----+---+
| Name|Age|
+-----+---+
|Alice| 25|
| Bob| 30|
|Cathy| 28|
+-----+---+
df.filter(df["Age"] > 25).show()
+-----+---+---------------+
| Name|Age| Job|
+-----+---+---------------+
| Bob| 30| Data Scientist|
|Cathy| 28|Product Manager|
+-----+---+---------------+
df.withColumn("Age_After_10_Years", df["Age"] + 10).show()
+-----+---+---------------+--------------------+
| Name|Age| Job|Age_After_10_Years|
+-----+---+---------------+--------------------+
|Alice| 25| Engineer| 35|
| Bob| 30| Data Scientist| 40|
|Cathy| 28|Product Manager| 38|
+-----+---+---------------+--------------------+
④ 分组聚合
df.groupBy("Job").avg("Age").show()
+---------------+--------+
| Job|avg(Age)|
+---------------+--------+
| Engineer| 25.0|
| Data Scientist| 30.0|
|Product Manager| 28.0|
+---------------+--------+
⑤ SQL 风格查询
df.createOrReplaceTempView("people")
spark.sql("SELECT Job, COUNT(*) as cnt FROM people GROUP BY Job").show()
+---------------+---+
| Job|cnt|
+---------------+---+
| Engineer| 1|
| Data Scientist| 1|
|Product Manager| 1|
+---------------+---+
⑥ 落盘 + 收尾
df.write.mode("overwrite").parquet("output/people.parquet")
df.write.mode("overwrite").option("header", "true").csv("output/people_csv")
spark.stop()
输出目录 output/ 下会出现 people.parquet/ 和 people_csv/ 两个文件夹。
映射表:代码行 ↔ 概念 ↔ 输出
| 概念/API | 在脚本里的位置 | 干了什么 | 对应输出 |
|---|---|---|---|
| 建会话 | 第 16 行 | 启动 Spark | 控制台 Spark 日志 |
createDataFrame |
第 27 行 | 内存列表 → DataFrame | 后续所有操作的数据来源 |
show(action) |
第 30 行 | 触发计算并打印全表 | 3 行表格 |
printSchema |
第 31 行 | 打印列结构与类型 | 列名 + 类型 |
select |
第 33 行 | 选列 | 2 列表格 |
filter |
第 34 行 | 筛行 | 剩下 Age>25 的 2 行 |
withColumn |
第 36 行 | 加一列 | 3 列表格(含新列) |
groupBy |
第 39 行 | 分组求均值 | 3 组 + avg(Age) |
createOrReplaceTempView |
第 41 行 | 注册成临时表 | 无 |
spark.sql |
第 42 行 | SQL 查询 | 每组计数 |
write(action) |
第 45~46 行 | 落盘 | output 目录生成文件 |
安全锁清单
用 Chrome 下载被拦截怎么办
换 Edge 下载;或直接用方法一 git clone——克隆不经过浏览器,不会被拦截。
为什么代码里还要再设一遍环境变量
系统变量只在 Shell 里生效,PyCharm / IDE 启动的进程不一定继承。脚本开头用 os.environ 强制设置,可以避免 UnsatisfiedLinkError。两边都配,双保险。
为什么写了 filter 却没有输出
不是没生效,是 transformation 懒执行——没接 .show() / count() 等 action 就不会真的算。看完加一句 .show()。
数据写不出来或版本对不上怎么办
先看报错类型:
- 权限 /
access0→ 检查hadoop.dll是否在%HADOOP_HOME%\bin,extraLibraryPath路径是否一致; - 目录已存在 → 加
mode("overwrite"); - 版本对不上 → 对照"版本匹配表",winutils 目录名要和 PySpark 内置 Hadoop 版本一致。
要不要下载完整 Hadoop 发行版
不要。Apache 官方 tar.gz(约 696MB)里根本没有 winutils.exe——那是社区补丁。clone 或 DownGit 只取一个 bin 目录就够了。
进阶方向
本文边界:只覆盖 Windows 本地单机 + 入门操作。以下不在范围内,需要时另开专题:
- 部署方式变种:Docker 跑 Spark、WSL2、远程集群(不再需要 winutils)
- 数据量变大:分区(partition)、缓存(cache)、广播变量(broadcast)
- 更多 IO:读取 CSV / JSON / JDBC、与 Pandas 互转
- 进阶 SQL:窗口函数(
row_number/rank)、多表join
下一步建议:把 learn_spark.py 跑通后,照着"操作选择表"把每个 API 都替换一遍练手感,再去看 Spark 官方 DataFrame 文档。
No comments yet. Be the first!