本章会用到的数据
公开数据:百万行教学交易数据用于分块计算;恒瑞医药分钟行情用于真实数据练习。
这些数据能做什么:比较 Pandas 与 Dask 的读取方式,并学习在数据较大时分块计算。
分析时注意:先只读取需要的列,观察运行时间和内存;时间预测仍要按先后顺序划分样本。
示例说明:41 MB 的教学文件用于理解方法,不代表真正的超内存计算。
【课堂核心】3 学时学习安排(1/2)
| 动机与先修 |
20 |
| 概念与公式 |
35 |
| 输入说明与检查 |
20 |
【课堂核心】3 学时学习安排(2/2)
| 独立任务与反馈 |
45 |
| 完整答案与完整示例 |
40 |
| 本章小结 |
20 |
【可选拓展】:集群实际使用、调度器细节与大规模机器学习框架比较。
欢迎来到大数据时代
- 我们正处在一个数据爆炸的时代。
- 从金融交易到社交媒体,海量数据正在重塑商业决策。
- 但我们迄今为止所学的工具,在面对“真正的大数据”时,会遇到一堵看不见的墙。
- 今天,我们的任务就是——翻越这堵墙。
本章学习议程
我们将循序渐进,探索大数据的世界:
- 问题的根源:为什么传统方法会失效?
- 核心思想:解决大数据问题的通用策略。
- 关键工具(Dask):深入了解Python生态的并行计算利器。
- 动手实践:用 Dask 对公开真实分钟数据执行多分区惰性计算。
- 高级话题:区分教学缩放接口、真实分区计算与需要另行资源验证的集群机器学习。
学习目标 (1): 理解挑战
在这堂课结束时,我希望你们能够理解:
为什么处理海量数据是一个巨大的挑战,以及它与传统数据分析的根本区别。
学习目标 (2): 了解工具
并且,能够了解:
Dask和Spark这两个业界主流的大数据处理框架的核心思想。
学习目标 (3): 掌握实践边界
最终,能够实践并掌握:
- 使用Dask执行基本的数据操作,并体会其与Pandas的异同。
- 解释海量数据机器学习的分区、增量与分布式训练策略,并在单机多分区示例中核对接口和执行边界。
核心问题:当数据大到内存装不下时,该怎么办?
我们之前所有的数据分析工作,都有一个隐含的前提…
那个我们从未质疑过的前提…
所有数据都能一次性加载到计算机的内存(RAM)中。
pandas 就是这样工作的。当你调用 pd.read_csv() 时,它会尝试将整个文件读入内存。
这在小数据集上表现优异,但在大数据时代,这个前提本身就是问题的根源。
内存之墙:理想与现实
我们希望用我们有限的内存,去分析无限增长的数据。
个人电脑的内存是极其有限的
让我们来看一下现实:
- 学生笔记本电脑:通常是 8GB 到 16GB 内存。
- 高端工作站:也许能达到 64GB 或 128GB,但很少超过 200GB。
这就像…
一个形象的比喻
…试图用一个水杯去装整个游泳池的水。
假设情景:十亿用户级交易数据
下面只做容量级别的课堂估算,不对应任何具名平台的披露数据。四项输入均为假设,用途是判断单机内存约束,而不是陈述真实业务规模。
- 假设月活跃用户: \(10^9\)
- 假设每位用户每月交易: 50 次
- 假设每次交易的数值字段: 10 个
- 假设每个字段存储大小: 8 字节(float64)
这会产生多大的数据量?
触目惊心的数据量估算
让我们来做个简单的计算:
\[
\large{
\text{数据量} = \underbrace{10^9}_{\text{用户数}} \times \underbrace{50}_{\text{交易数}} \times \underbrace{10}_{\text{变量数}} \times \underbrace{8}_{\text{字节/变量}}
}
\]
问题的严重性:4TB
计算结果是:
\[
\large{
\text{结果} \approx 4 \times 10^{12} \text{ 字节} \approx \textbf{4 Terabytes}
}
\]
4TB 的裸数值容量通常远超学生常用设备的内存,且还未计入索引、字符串和运行时开销。
可选策略包括按列/分区读取、数据库下推、流式聚合和分布式计算;是否需要 Dask 取决于数据布局、操作类型与硬件预算。
解决方案:分而治之 (Divide and Conquer)
当数据无法一次性装入可用内存时,一个常见办法是:
将一个大任务(处理整个数据集)分解成许多可以在小块数据上执行的小任务,然后将结果汇总起来。
这正是大数据处理框架的精髓所在。
“分而治之”的可视化理解
Python生态中的两大主流工具
在Python世界里,有两个强大的程序库专门用来实现“分而治之”的思想:
- Dask: 一个轻量级、原生支持Python的并行计算库。
- Apache Spark: 一个功能更强大、生态更成熟的分布式计算引擎。
我们今天将主要以Dask为例进行讲解,因为它对熟悉Pandas的你来说最容易上手。
工具介绍: Dask
Dask: 轻量级、Python原生、与现有库(Pandas, NumPy)无缝集成,学习曲线平缓。
工具介绍: Apache Spark
Apache Spark: 行业标准、功能强大、生态成熟、支持多语言(通过PySpark在Python中使用)。
Dask: 熟悉的味道,更强的能力
Dask是一个开源项目,旨在将Python数据科学生态系统(Pandas, Scikit-Learn, NumPy)的能力扩展到大数据和并行计算领域。
它的设计哲学是:尽可能地模仿现有库的API,让你用最小的学习成本处理更大的数据集。
Dask的优点:与Pandas的惊人相似性
Dask最吸引人的一点在于,你不需要学习一套全新的语法。
- Dask DataFrame: Dask的核心数据结构之一,它在内部由多个小的Pandas DataFrame组成。
- API兼容: Dask DataFrame 覆盖许多常用 Pandas 操作,但并非完整替代;不受支持的操作、全局排序与大规模 shuffle 需要改写或换用数据库/分布式方案。
这使得从Pandas迁移到Dask变得异常平滑。
Dask的五大核心特点
- 懒惰执行 (Lazy Execution)
- 任务调度 (Task Scheduling)
- 智能内存管理 (Smart Memory Management)
- 并行与分布式计算 (Parallel & Distributed)
- 无缝生态集成 (Ecosystem Integration)
我们将逐一深入探讨这些概念。
核心特点 (1): 懒惰执行 (Lazy Execution)
这是Dask与Pandas最本质的区别。
- Pandas: 你写的每一行代码都会立即执行。
- Dask: 多数 collection 的读取与变换先构建任务图,不会在定义表达式时立刻物化全部结果。
相反,Dask会构建一个任务图(Task Graph),记录下你想要执行的所有计算步骤。
什么是任务图?
任务图是 Dask 的“作战计划”:
需要执行的所有任务(例如:从文件中读取一个数据块,过滤行,计算平均值)。
任务之间的依赖关系(例如:必须先读取数据块,然后才能对其进行过滤)。
触发执行:.compute()、.persist()、.head() 等查看或缓存操作。
模型触发:需要统计量归约的 estimator .fit() 也会执行任务图。
判断原则:逐项识别 API 的执行与物化边界,不能简化为“只有 .compute() 才计算”。
思想实验:懒惰执行就像是写一份菜谱
想象一下做一道复杂的菜:
立即执行 (Pandas)
- 读第一步“切洋葱”
- ✅ 立刻切洋葱
- 读第二步“热锅”
- ✅ 立刻热锅
- …一步一动…
懒惰执行 (Dask)
- 读完整份菜谱
- 🧠 在脑中规划
- “先切好所有蔬菜,再混合酱料,最后一起下锅。”
- 直到决定“开火”时,才真正开始动手。
Dask的这种方式让它有机会在执行前对整个计算流程进行优化。
核心特点 (2): 任务调度 (Task Scheduling)
懒惰执行构建了任务图,而任务调度器(Task Scheduler)则是负责执行这个图的“大脑”。
Dask会将大的计算任务分解成许多小的、独立的任务,然后智能地安排它们的执行顺序,并分配到不同的CPU核心或机器上并行处理。
任务调度器的工作原理
核心特点 (3): 智能内存管理
Dask可按分区调度任务,但内存行为取决于图依赖、分区大小、缓存与调度器配置。
计算时相关分区会进入内存;不再被依赖的中间结果可能被释放,分布式调度器也可按阈值 spill 到磁盘,但持久化对象、过大分区或 shuffle 仍可能造成内存压力。
因此必须监控峰值内存并合理划分分区,不能假定每个数据块计算后都会立即释放。
内存管理的“流水线”模式
核心特点 (4): 并行与分布式计算
Dask的强大之处在于它的可伸缩性:
并行计算 (Parallel)
- 单机多核是可选执行方式
- 只有任务可并行、分区足够大且调度、序列化与 I/O 开销相对较小时,合适的 Dask scheduler 才可能带来加速;Python GIL、任务过碎或数据搬运都可能抵消收益。
分布式计算 (Distributed)
- 多机集群扩展容量与并发
- Dask 可以协调多台机器处理分区任务,但收益必须用同一工作负载的 wall time、峰值内存、吞吐与调度开销实测;数据规模大不保证更快。
并行 vs. 分布式
核心特点 (5): 与现有生态系统无缝集成
Dask不仅仅是模仿了Pandas,它与整个Python科学计算栈都深度集成。
这意味着你可以用很小的代码改动,就将现有分析流程从处理小数据升级为处理大数据。
Dask 与 Python科学计算全家桶
Dask 与 Pandas 的操作对比
下表直观地展示了Dask与Pandas在基本操作上的相似与不同之处。
| 数据结构 |
DataFrame, Series |
Dask DataFrame, Dask Series |
| 导入库 |
import pandas as pd |
import dask.dataframe as dd |
| 读取数据 |
pd.read_csv(...) |
dd.read_csv(...) |
| 查看数据 |
df.head(), df.tail() |
df.head(), df.tail() |
| 数据过滤 |
df[df['col'] > 0] |
df[df['col'] > 0] |
| 分组聚合 |
df.groupby('col').sum() |
df.groupby('col').sum() |
| 数据合并 |
pd.concat(), df.merge() |
dd.concat(), df.merge() |
| 计算执行 |
立即执行 |
多数变换先建任务图;.compute()、.persist()、.head() 或 estimator .fit() 等操作可触发执行 |
理解.compute(): Dask操作的“执行”按钮
.compute()方法是Dask懒惰执行模型的关键。
- 当你写下
df_filtered = df[df['column'] > 0] 时,如果df是一个Dask DataFrame,df_filtered并不是一个包含结果的数据帧。
df_filtered此时只是一个“任务图”,一个“计算计划”,它描述了如何得到过滤后的结果。
result = df_filtered.compute() 触发本例的读取、过滤与汇总,并返回 Pandas DataFrame。
- 其他流程也可能由
.persist()、查看操作或 estimator .fit() 触发。
.compute() 的作用:从计划到现实
实践环节目标
- 安装 Dask 库。
- 检查仓库随附的教学缩放交易数据版本,并明确它不代表真实市场或超内存规模。
- 使用Dask读取并分析该数据集。
- 亲身体验懒惰执行和
.compute()。
步骤1: 安装 Dask 库
打开你的终端或Anaconda Prompt,根据你的需求选择安装命令:
步骤2: 检查确定的教学缩放数据集
本章使用约 41MB 的公开教学 CSV,观察分区读取、惰性执行和分块计算。它通常小于机器内存,因此本章只练习方法,不声称完成了 TB 级数据测试。
我们的目标是对同一份一百万行数据版本完成“输入检查—分区读取—年度归约—期望输出”的可自查小结。
开始前先检查数据文件
代码
# 导入 `pandas`,用于查看字段、行数与时间范围。
import pandas as pd
# 为后续 Dask-ML 与公开分钟验证导入 `numpy` 并绑定 `np`,用于数组和数值计算。
import numpy as np
# 导入 `os`,用于检查文件是否存在与查看文件大小。
import os
from pathlib import Path
from urllib.request import urlretrieve # 复用本章隐藏设置单元安装的浏览器标识下载器
# 固定教学 CSV 的期望行数,作为截断或追加文件的结果不理想时如何解释。
n_rows = 1_000_000 # 100万行
# 文件不存在时,从与其他章节相同的公开数据域下载。
# 按 Linux 共享数据、Windows 共享数据、项目缓存的顺序选择大型交易教学文件。
file_path = next((candidate_path for candidate_path in [Path('/home/ubuntu/r2_data_mount/data/course/large_transactions.csv'), Path('C:/qiufei/data/course/large_transactions.csv'), Path('data/course/large_transactions.csv')] if candidate_path.exists()), Path('data/course/large_transactions.csv'))
if not file_path.exists():
file_path.parent.mkdir(parents=True, exist_ok=True)
urlretrieve('https://assets.qiufei.site/data/course/large_transactions.csv', file_path)
查看数据的字段与范围
先查看字段、行数与时间范围,确认下载的数据适合后续练习。
代码
# 读取字段名而不载入数据行,核对教学 CSV 的五列说明与顺序。
observed_columns = pd.read_csv(file_path, nrows=0).columns.tolist()
# 要求字段说明完全一致,防止旧版或手工修改文件进入 Dask 练习。
assert observed_columns == ['transaction_id', 'timestamp', 'stock_id', 'price', 'volume'], '教学 CSV 字段或顺序不符合课程说明'
# 初始化逐块累计的行数与时间上下界,避免为检查一次性物化全部字段。
validated_row_count, validated_time_min, validated_time_max = 0, None, None
# 分块读取交易编号与时间字段,检查完整样本范围而不依赖文件字节数。
for validation_chunk in pd.read_csv(file_path, usecols=['transaction_id', 'timestamp'], chunksize=100_000):
# 把当前块时间解析为时间类型,供全文件上下界比较。
validation_times = pd.to_datetime(validation_chunk['timestamp'])
# 累计当前块行数,最终应严格等于一百万行。
validated_row_count += len(validation_chunk)
# 更新已读分块中的最早时间,核对数据版本说明的确定时间窗口。
validated_time_min = validation_times.min() if validated_time_min is None else min(validated_time_min, validation_times.min())
# 更新已读分块中的最晚时间,核对数据版本说明的确定时间窗口。
validated_time_max = validation_times.max() if validated_time_max is None else max(validated_time_max, validation_times.max())
# 要求输入严格包含一百万行,防止截断文件或追加文件通过检查。
assert validated_row_count == n_rows, '教学 CSV 行数不符合课程说明'
# 要求所有时间均落在确定的左闭右开秒级窗口内。
assert validated_time_min >= pd.Timestamp('2017-07-14 02:40:00') and validated_time_max < pd.Timestamp('2020-09-13 12:26:40'), '教学 CSV 时间范围不符合课程说明'
# 显示数据概况,随后进入 Dask 读取与练习。
print({'字段数': len(observed_columns), '行数': validated_row_count, '时间范围': (validated_time_min, validated_time_max)})
{'字段数': 5, '行数': 1000000, '时间范围': (Timestamp('2017-07-14 02:42:24'), Timestamp('2020-09-13 12:25:54'))}
代码
# 读取教学 CSV 的实际字节数并换算为 MB,核对 8MB 分区参数会产生多个块。
file_size_mb = os.path.getsize(file_path) / (1024**2)
# 报告这份 CSV 占用的兆字节数,核对 8MB 分区设置能产生多个分区。
print(f'文件大小约为: {file_size_mb:.2f} MB')
为什么先看数据概况?
- 字段、行数与时间范围能帮助我们发现下载不完整或文件用错等问题。
- 分块检查只让一个十万行检查块同时驻留,演示降低检查阶段峰值对象的模式。
- 本次约 41MB 文件没有证明一次性读取会 OOM;是否 OOM 必须结合实测峰值和机器预算。
- 后续年度
groupby 仍需汇总所有源分区,最终四行输出不等于只读取少数分区。
步骤3: 使用 Dask 读取数据
现在,关键时刻来了。我们不用pandas,而是用dask.dataframe来读取这个大文件。
代码
# 为“步骤3: 使用 Dask 读取数据”导入 `dask.dataframe` 并绑定 `dd`,用于建立当前任务的分区表格与惰性任务图。
import dask.dataframe as dd
# 使用Dask读取CSV
# 8MB分块可保证本教学文件形成多个分区
ddf = dd.read_csv('large_transactions.csv', blocksize='8MB')
# 此时尚未物化全量分区;`read_csv` 可能已抽取少量样本推断字段和类型。
print('Dask DataFrame已创建!')
# 打印 Dask DataFrame 的逻辑分区数,核对 8MB 块设置确实形成多分区教学输入。
print(f'分区数: {ddf.npartitions}')
# 要求教学文件确实跨越多个 Dask 分区,否则不能据此演示分块执行。
assert ddf.npartitions > 1, '教学文件必须形成多个分区才能演示分块执行'
Dask DataFrame已创建!
分区数: 5
解读 dd.read_csv
import dask.dataframe as dd: 这是Dask DataFrame模块的惯用别名,就像import pandas as pd一样。
blocksize='8MB':该教学参数确保约 25—45MB 的文件形成多个分区。实际使用时的参数应依据文件大小、可用内存和调度开销实测选择。
审视Dask DataFrame对象
让我们看看刚刚创建的ddf对象到底是什么。
代码
# 输出惰性 Dask DataFrame 元数据,查看分区、字段与尚未物化的任务对象。
print(ddf)
Dask DataFrame Structure:
transaction_id timestamp stock_id price volume
npartitions=5
int64 string int64 float64 int64
... ... ... ... ...
... ... ... ... ... ...
... ... ... ... ...
... ... ... ... ...
Dask Name: to_string_dtype, 2 expressions
Expr=ArrowStringConversion(frame=FromMapProjectable(b606a15))
解读输出结果
- 没有数据! 和Pandas不同,它不会打印出数据内容。
- 列名和类型: 它显示了
Dask DataFrame Structure,包含了列名和dtype。Dask通过智能地读取文件的一小部分来推断这些元信息。
- 分区 (npartitions): 这是核心。它告诉你这个Dask DataFrame在逻辑上被分成了多少块(分区)。这个数量取决于文件总大小和
blocksize。
逻辑依赖示意:从分区读取到聚合
Dask 先构造惰性任务计划,再由 .compute() 触发本例计算。下面的图只用 3 个代表性分区说明过程;当前 CSV 实际形成 5 个分区。
代码
ddf['amount'] = ddf['price'] * ddf['volume'] # 为每个分区建立“价格乘数量”的惰性列变换
total_amount = ddf['amount'].sum() # 建立跨分区归约任务,此时尚未执行全量计算
print({'当前实际分区数': ddf.npartitions, '归约对象类型': type(total_amount).__name__}) # 检查五个实际分区与等待执行的 Dask 标量,不把简化 SVG 冒充完整任务图
{'当前实际分区数': 5, '归约对象类型': 'Scalar'}
如何解读逻辑依赖示意
本图用三个代表性分区解释逻辑依赖;当前对象实际有 5 个分区,完整任务键、优化层和调度细节未画出。
- 左列圆角矩形:读取源分区;
… N 表示其余分区采用同一逻辑。
- 中列圆角矩形:在各分区内计算
price × volume,仍只形成惰性列变换。
- 右列圆角矩形:汇总全部分区的局部结果,形成跨分区归约。
- 从左向右的箭头:表示读取 → 局部变换 → 全局汇聚的依赖方向。
构建 amount 与 sum 只形成计划;本例调用 .compute() 才触发该归约。其他查看操作、.persist() 或 estimator .fit() 也可能在各自流程中触发执行。
步骤4: 执行计算并获取结果
现在,让我们通过调用.compute()来触发这些计算。我们将计算price列的基本描述性统计。
代码
# 计算描述性统计
# 这个过程会花费一些时间,因为Dask现在正在读取所有数据块并进行计算
print("开始计算描述性统计...")
# 由 `.compute()` 触发全部价格分区的描述统计归约,并返回物化的 Pandas Series。
summary_stats = ddf['price'].describe().compute()
# 报告聚合结果已经物化为 Pandas 对象。
print("计算完成!")
# 输出物化后的价格描述统计,检查跨分区归约的最终 Pandas 结果。
print(summary_stats)
开始计算描述性统计...
计算完成!
count 1000000.000000
mean 255.179348
std 141.401938
min 10.000000
25% 133.150000
50% 255.700000
75% 377.970000
max 500.000000
Name: price, dtype: float64
.compute() 幕后发生的事
- Dask Scheduler接收到
describe 任务图。
- 它将任务(为每个分区计算统计量)分配给可用的CPU核心。
- 每个核心读取一个数据块,计算
count, mean, std, min, max 等。
- 所有分区的中间结果被汇总起来,计算出最终的全局统计量。
- 峰值内存取决于分区大小、聚合中间量、任务图、调度器与 spill 配置;必须实测,不能由“分区执行”推断严格的小内存上界。
Dask的Pandas-like操作:分组聚合
Dask的groupby操作也和Pandas几乎一样。让我们计算每只股票的平均交易量。
代码
# 计算每只股票的平均交易量
# 同样,这只是定义了计算计划
mean_volume_by_stock = ddf.groupby('stock_id')['volume'].mean()
# 使用查看操作把输出限制为前十个分组结果,以节省屏幕空间。
# `head(10)` 返回 Pandas 对象,但全局 groupby 均值仍可能需要所有源分区提供分子与分母。
print("开始分组聚合计算...")
# 取得前十个已聚合分组;`head` 限制返回行与输出分区,不承诺只读取前几个源分区。
result = mean_volume_by_stock.head(10)
# 报告受限输出已由 `head` 触发并物化;全局归约的源分区工作量不能由返回十行判断。
print("计算完成!")
# 输出前十个股票的平均成交量,检查分组键与聚合值。
print(result)
开始分组聚合计算...
计算完成!
stock_id
1 4951.722832
2 5243.609553
3 5261.428571
4 5137.991812
5 5016.143145
6 5053.174588
7 4856.658436
8 5126.087500
9 5171.106833
10 5086.707676
Name: volume, dtype: float64
简单小结:Dask vs. Pandas
- 相似性: API非常接近,如果你会Pandas,你几乎就会用Dask DataFrame。
- 核心差异: Dask 的多数 collection 变换先构建任务图,再由
.compute()、.persist()、部分查看操作或 estimator .fit() 等触发执行;分区与调度使其能够处理大于单机内存的数据。
- 适用场景: Dask 常适合 Python 原生、可分区的单机或集群任务;是否采用它必须同时核对数据格式、操作与 shuffle 图、峰值内存、容错、集群条件、团队生态和基准测试,文件大小本身不是决定规则。
简介:Apache Spark 与 PySpark
当任务需要成熟的多语言数据平台、强容错与大规模集群治理时,可以评估 Apache Spark;数据规模只是输入之一,不能仅凭 TB 数量替代对格式、操作图、倾斜、资源与团队生态的判断。
- Apache Spark: 是一个为大规模数据处理而设计的统一分析引擎。它最初用Scala编写,拥有卓越的性能和可伸缩性。
- PySpark: 是Spark为Python提供的官方API,让我们可以在Python中利用Spark的强大功能。
Spark的核心特点:内存集群计算
Spark最著名的特点是其内存计算(In-memory Computing)能力。
与早期的大数据系统(如Hadoop MapReduce)不同,Spark可以将中间计算结果保存在集群中各台机器的内存里,大大加快了迭代算法(如机器学习)和交互式查询的速度。
Spark的集群计算模型
Dask 与 Spark: 何时选择哪个?
这是一个常见的问题。下表先比较两种工具;实际选择还要结合数据格式、运行时间、内存和已有学习环境,不能只看文件大小。
| 生态系统 |
纯Python,与Pandas/Numpy/Scikit-Learn紧密集成 |
多语言 (Scala, Java, Python, R),更庞大、更复杂的生态 |
| 设置复杂度 |
非常简单 (pip install dask) |
相对复杂,通常需要配置Java环境和集群管理器 |
| 适用场景 |
单机并行、中等规模集群,Python原生项目 |
大型企业级集群,超大规模数据,多语言团队 |
| 学习曲线 |
平缓,对Pandas用户友好 |
较陡峭,需要理解Spark的架构和概念 |
总结:如何选择
统一决策证据:
- 先记录数据格式与分区键、操作和 shuffle 图、倾斜风险、实测峰值内存与墙钟时间。
- 再比较失败恢复与容错要求、可用集群、既有 Python/Spark 生态和团队运维能力。
- Dask 与 PySpark 都必须通过代表性切片和目标规模压测;数据量不单独决定框架。
大数据机器学习的挑战
- 失败条件:estimator 要求全量物化,且输入或关键中间量超过内存预算。
- 典型后果:
.fit(X, y) 发生 OOM 或严重换页。
- 可扩展接口:
partial_fit、核外 solver,或可接受 Dask 分区输入的 estimator。
- 仍需核对:具体 solver、训练规则与中间量规模。
何时全量 .fit(X, y) 可能失败?
对策1 (简单但有缺陷): 随机抽样
最直接的方法是将大数据集缩小到内存可以容纳的大小。
- 做法: 从海量数据中随机抽取一部分样本(例如1%)。
- 优点: 简单易行,可以使用所有标准的
scikit-learn工具。
随机抽样的代价
缺点:
- 信息丢失: 丢弃了99%的数据,可能会错过重要的模式,尤其是在处理稀有事件(如金融欺诈)时。
- 模型性能: 对于需要大量数据才能学习的复杂模型(如深度学习),抽样后的数据可能不足以训练一个高性能的模型。
对策2:一种常见的同步数据并行训练
下面展示的是同步数据并行的一种常见设计,而不是分布式训练的完整定义:各 worker 持有模型副本、处理不同批次,并在每步聚合梯度后同步参数。模型并行、流水线并行与异步参数服务器等设计不遵循同一流程。
它可利用多个计算单元扩展训练容量或缩短训练时间,但实际收益取决于计算/通信比、互连带宽、负载均衡、同步等待和容错开销,必须以单机基线实测。
同步数据并行示例 (1/4):分区与分发
- 准备环境: 配置一个包含多个计算节点(worker)的集群。
- 数据分区: 将大型数据集分割成许多小批次(batches)。
- 数据分配: 将这些数据批次均匀地分配给每个节点。
同步数据并行示例 (2/4):模型复制
- 模型复制: 在每个节点上初始化一个完全相同的模型副本(拥有相同的初始权重)。
同步数据并行示例 (3/4):并行计算
- 并行计算: 每个节点使用自己分配到的数据批次,独立地对模型进行训练,并计算出模型参数的更新量(梯度)。
同步数据并行示例 (4/4):聚合与更新
- 梯度聚合: 所有节点计算出的梯度被汇总起来。
- 模型更新: 使用聚合后的梯度来更新所有节点上的模型参数。
- 迭代: 重复并行计算和更新的步骤。
分布式训练的优化与工具:Dask-ML
幸运的是,我们有现成的工具来处理这些复杂性。
Dask-ML是 Dask 生态中扩展机器学习步骤的模块,部分接口接近 scikit-learn;是否真正形成分布式或核外训练,仍取决于 estimator、solver、collection、scheduler 与实际使用方式。
Dask-ML 的三大应用方式
- 扩展 Scikit-Learn
- 提供了许多与Scikit-Learn API兼容的并行算法。
- 增量学习 (核外学习)
- 可扩展的预处理
- 提供并行
StandardScaler 等工具,可由训练分区样本逐特征归约均值与标准差。
Dask-ML实践: 目标
回到之前检查的确定 large_transactions.csv 数据集,尝试训练一个线性回归模型,预测交易量 volume。
我们将用 dask-ml 演示多分区接口;该约 40MB 合成文件通常可一次性载入内存,因此不能作为核外或超内存证据。本章的真实数据验证见后文恒瑞医药分钟数据版本。
步骤1 - 准备数据
首先,我们像之前一样用Dask读取数据,并定义我们的特征(X)和目标(y)。
代码
# 为“步骤1 - 准备数据”导入 `dask.dataframe` 并绑定 `dd`,用于建立当前任务的分区表格与惰性任务图。
import dask.dataframe as dd
# 为“步骤1 - 准备数据”,导入 Dask-ML 缩放器,逐特征归约全量机制拟合样本的均值与标准差;本页不作样本外验证。
from dask_ml.preprocessing import StandardScaler
# 为“步骤1 - 准备数据”,导入 Dask-ML 线性回归器;本章显式确定 L2 惩罚、C=1.0 与 ADMM 求解器,不把它称为普通 OLS。
from dask_ml.linear_model import LinearRegression
# 读取数据
ddf = dd.read_csv('large_transactions.csv', blocksize='8MB')
# 为“步骤1 - 准备数据”,定义特征和目标
# 为了简化,我们只用 price 和 stock_id 作为特征
# 注意:直接使用 stock_id 是不好的做法,这里仅为演示
X = ddf[['price', 'stock_id']]
# 提取示例目标数组,供 Dask-ML 缩放与线性回归。
y = ddf['volume']
# 把两列分区特征转换为保留分块信息的 Dask Array,供 Dask-ML 估计器消费。
X = X.to_dask_array(lengths=True)
# 把目标 Series 转换为与特征行分区对齐的 Dask Array,而非物化为 NumPy 数组。
y = y.to_dask_array(lengths=True)
# 报告特征与目标仍保持为分块的 Dask Array,而非已物化数组。
print("X 和 y 已准备为 Dask Arrays.")
# 输出 Dask Array 的分块与形状表示,确认特征没有被物化成 NumPy 数组。
print(X)
X 和 y 已准备为 Dask Arrays.
dask.array<read-_to_string_dtype-getitem-values, shape=(1000000, 2), dtype=float64, chunksize=(202096, 2), chunktype=numpy.ndarray>
步骤2 - 数据标准化
fit:跨分区归约拟合样本的均值与方差,并物化拟合属性。
transform:使用已拟合统计量,返回的 Dask Array 仍保持惰性。
- 检查点:统计量只能来自机制拟合样本。
代码
# 创建将由机制拟合样本逐特征估计统计量的 Dask 缩放器。
scaler = StandardScaler()
# 逐特征汇总机制拟合样本的分区统计量,并物化均值与尺度属性。
scaler.fit(X)
# 使用确定的拟合样本逐特征统计量构建惰性标准化数组。
X_scaled = scaler.transform(X)
# 打印已物化均值属性与惰性转换结果的类型,核对 `fit` 和 `transform` 的执行边界。
print({'均值属性类型': type(scaler.mean_).__name__, '转换结果类型': type(X_scaled).__name__}) # 区分已物化拟合属性与惰性转换集合
# 打印惰性 Dask Array 的表示、形状与分块元数据;此句不计算列均值或验证缩放数值。
print(X_scaled)
{'均值属性类型': 'ndarray', '转换结果类型': 'Array'}
dask.array<truediv, shape=(1000000, 2), dtype=float64, chunksize=(202096, 2), chunktype=numpy.ndarray>
步骤3 - 训练模型
- 固定规格:
penalty='l2'、C=1.0、solver='admm'。
- 模型含义:在 Dask-ML 2025.1.0 中,这是 \(\lambda=1/C\) 的 L2 惩罚线性回归,不是普通最小二乘。
- 执行边界:是否跨分区仍取决于输入、solver、scheduler 与实际使用。
代码
# 显式确定 L2 惩罚、C=1.0 与 ADMM 求解器,避免类默认值随版本变化并误称为普通 OLS。
lr = LinearRegression(penalty='l2', C=1.0, solver='admm')
# 训练模型
# 这一步会触发计算!
# Dask-ML会在后台读取数据块,进行计算
print('开始训练模型...')
# `.fit()` 触发分区计算并物化线性回归系数与截距,完成本页机制样本拟合。
lr.fit(X_scaled, y)
# 在拟合完成后报告模型参数已由分区归约计算得到。
print('模型训练完毕!')
# 输出本例确定的模型族、正则强度与求解器,使结果不依赖未报告默认值。
print({'模型族': 'L2惩罚线性回归', 'penalty': lr.penalty, 'C': lr.C, 'solver': lr.solver})
# 打印模型系数
print(f'模型系数 (Coefficients): {lr.coef_}')
# 打印已拟合线性回归的截距,核对 `.fit()` 已物化模型参数而非返回惰性预测集合。
print(f'模型截距 (Intercept): {lr.intercept_}')
开始训练模型...
模型训练完毕!
{'模型族': 'L2惩罚线性回归', 'penalty': 'l2', 'C': 1.0, 'solver': 'admm'}
模型系数 (Coefficients): [ 4.93338112 -2.6660765 ]
模型截距 (Intercept): 5044.245107269162
解读 Dask-ML 的 .fit()
- 当我们调用
.fit()时,Dask-ML和Dask在后台协同工作。
- Dask负责高效地从磁盘读取数据块(partitions)。
- Dask-ML则负责在这些数据块上执行并行的机器学习算法(例如,计算每个块的梯度)。
- 最后,Dask-ML负责聚合所有块的结果,得到一个统一的最终模型。
- 整个过程对用户来说是透明的,API与Scikit-Learn高度一致。
阶段小结 (1): 思维的转变
从“一次性加载”到“分块处理”: 这是从传统数据分析迈向大数据分析最根本的思维转变。
你必须放弃“所有数据都在我眼前”的假设,开始用“流水线”和“分治”的眼光看待问题。
阶段小结 (2): 工具的力量
懒惰执行的力量: 通过先规划(构建任务图)再按需物化(如 .compute()、.persist()、.head() 或 .fit()),Dask 可优化执行并在合适分区、算法、调度和 spill 配置下处理超出单机内存的数据;本章小文件不构成该能力的实测证据。
这是一种更智能、更具伸缩性的计算范式。
阶段小结 (3): 生态的威力
Python生态的威力: Dask和PySpark让Python这门通用语言具备了处理工业级大数据的能力,而Dask-ML则将这种能力无缝延伸到了机器学习领域。
你所学的技能,可以直接应用于解决真实世界的大规模问题。
检查:何时不应调用 .compute()?
- 场景一: 你需要处理一个2GB的CSV文件,并且你的电脑有16GB内存。你会选择哪个库(Pandas还是Dask)?为什么?
- 场景二: 你需要处理一个5TB的用户行为日志数据集,并构建一个复杂的推荐系统,你的公司有一个大型计算集群。你应该考虑使用哪个库(Dask还是PySpark)?为什么?
- 场景三: 一个 Dask
groupby 仍是 lazy collection。请分别判断何时可以调用 .compute()、何时必须拒绝:结果预计 2KB 且任务图峰值内存 3GB;结果大小未知且 shuffle 峰值可能超过 16GB。
反馈:文件大小不能单独决定框架
框架选择
- 2GB CSV:依据解析后的峰值内存与墙钟基准决定 Pandas 是否可行。
- 5TB + 复杂 shuffle + 大集群:优先评估 PySpark,并核对容错与团队运维能力。
反馈:何时可以物化?
- 可以:结果约 2KB、任务图峰值约 3GB,且运行时间和内存均低于预算;同时记录墙钟与峰值。
- 拒绝:结果大小未知或 shuffle 峰值可能超预算;先做分区抽样、分阶段归约、
persist 或写出分区结果。
- 订正路径:混淆惰性与任务图,回看惰性执行和任务图;遗漏物化副作用,回看执行与物化。
独立任务:Dask 分区归约
使用课程确定的 large_transactions.csv,完成以下任务:
- 先报告文件大小、字段、行数与实际
npartitions;数据或分区不符即停止。
- 将
timestamp列转换为年份。
- 计算每年的最高和最低交易价格(
price)。
实现提示:dd.to_datetime()、.dt.year、groupby() 与 agg()。
独立任务:复算、资源与失败路径
- 用 Pandas 对同一数据版本、同一字段与同一年度口径复算,并以
atol=0.001 核对 Dask 结果。
- 分别记录 Dask 与 Pandas 的墙钟秒数和进程峰值内存;当前 41MB 数据版本只比较接口与口径,不证明超内存优势。
- 保留失败路径:哈希/字段/行数不符、分区数不大于 1、两框架年度结果不一致时,明确停止且不得提交性能结论。
习题 4:Dask 合并与框架选择
假设你还有一个stocks_info.csv文件,包含stock_id和stock_name两列。
- 合并设计:说明如何用 Dask 按
stock_id 与 large_transactions.csv 合并。
- 大小分支:分别说明
stocks_info.csv 为小表或大表时的策略。
- 结果用途:按
stock_name 分析交易情况。
- 框架复答:把选择与峰值内存、墙钟、分区/shuffle、容错及团队运维逐项对应。
练习:先完成再查看答案:年度归约作业结果
- 计算证据:CSV 哈希与分区、年度极值、Pandas 同口径基线、墙钟与峰值内存。
- 合并证据:小表/大表策略、
stock_id 唯一性、合并前后行数与未匹配比例。
- 决策证据:两项框架场景自查和失败路径。
- 重做条件:缺任一完成内容、缺分区,或只报速度不报口径。
课堂核心完整反馈(1/5):场景自查
- 场景一判断:2GB 文件在 16GB 内存中未必安全;解析、字符串列与中间副本会放大内存。
- 场景一证据:先做 Pandas 分列、分块基准;峰值不可控时再评估 Dask。
- 场景二默认:5TB 日志、复杂推荐系统与大型集群下,优先选择 PySpark。
- 改选条件:已有 Dask Distributed 运维能力,且分区、倾斜、shuffle 压测与容错演练均达标。
4 分评分:题面一致、分区/shuffle、容错/团队、资源与失败恢复计划各 1 分;只写框架名或只按文件大小判断不得分。
课堂核心完整反馈(2/5):用相同口径比较 Pandas 与 Dask
代码
import time
# Pandas:一次读入所需两列,再计算各年的最低价与最高价。
pandas_start = time.perf_counter()
pandas_frame = pd.read_csv(file_path, usecols=['timestamp', 'price'])
pandas_frame['year'] = pd.to_datetime(pandas_frame['timestamp']).dt.year
pandas_price_bounds = pandas_frame.groupby('year')['price'].agg(['min', 'max']).sort_index()
pandas_seconds = time.perf_counter() - pandas_start
# Dask:把同一文件分成多个部分,最后用 `.compute()` 得到结果。
dask_start = time.perf_counter()
dask_frame = dd.read_csv(file_path, usecols=['timestamp', 'price'], blocksize='8MB')
dask_frame['year'] = dd.to_datetime(dask_frame['timestamp']).dt.year
annual_price_bounds = dask_frame.groupby('year')['price'].agg(['min', 'max']).compute().sort_index()
dask_seconds = time.perf_counter() - dask_start
assert dask_frame.npartitions > 1
runtime_comparison = pd.DataFrame({
'工具': ['Pandas', 'Dask'],
'运行时间(秒)': [pandas_seconds, dask_seconds]
})
display(pd.DataFrame([{'文件大小(字节)': os.path.getsize(file_path), 'Dask 分区数': dask_frame.npartitions}]))
display(runtime_comparison.round(3))
课堂核心完整反馈(3/5):年度极值一致性
- Dask 参考代码:
代码
pd.testing.assert_frame_equal(annual_price_bounds, pandas_price_bounds, check_exact=False, atol=0.001) # 两框架必须同口径一致
display(annual_price_bounds) # 输出可检查的完整答案
print({'与 Pandas 一致': True, '允许的数值误差': 0.001})
| year |
|
|
| 2017 |
10.0 |
500.0 |
| 2018 |
10.0 |
500.0 |
| 2019 |
10.0 |
500.0 |
| 2020 |
10.0 |
500.0 |
{'与 Pandas 一致': True, '允许的数值误差': 0.001}
确认输入包含五列字段和一百万行后,以上动态四行表就是期望输出;价格按两位小数核对,绝对容差为 0.001。
课堂核心完整反馈(4/5):归约解释与失败路径
- 归约机制:每个年份先接收各源分区的局部极值,再做跨分区归约。
- 执行范围:
.compute() 触发完整年度分组任务图;最终四行不代表只读取少数分区。
- 失败路径:文件、字段或行数不符先停止;年度结果差异超过
0.001 时回查日期与分组步骤。
- 证据边界:当前运行时间只适用于这台电脑和教学数据,不证明中国市场、超内存或集群性能。
课堂核心完整反馈(5/5):合并验证
- 合并答案:若
stocks_info.csv 很小,可先读成单分区并广播式合并;若同样很大,应按 stock_id 合理分区并检查 shuffle 成本。合并后验证行数、主键唯一性和未匹配比例,防止多对多连接放大样本。
公开中国数据练习:Dask 分区计算
代码
from pathlib import Path # 管理公开 A 股分钟数据路径
from urllib.request import urlretrieve # 复用本章已安装的浏览器标识下载器
import dask.array as da # 建立惰性分区数组并执行聚合
import h5py # 访问非 Pandas HDF5 字段视图
# 按 Linux 共享数据、Windows 共享数据、项目缓存的顺序选择分钟行情文件。
minute_path = next((candidate_path for candidate_path in [Path('/home/ubuntu/r2_data_mount/data/stock/minbar/600276.XSHG.h5'), Path('C:/qiufei/data/stock/minbar/600276.XSHG.h5'), Path('data/course/600276.XSHG.h5')] if candidate_path.exists()), Path('data/course/600276.XSHG.h5'))
if not minute_path.exists():
minute_path.parent.mkdir(parents=True, exist_ok=True)
urlretrieve('https://assets.qiufei.site/data/stock/minbar/600276.XSHG.h5', minute_path)
minute_store = h5py.File(minute_path, 'r') # 保持文件打开以支持 Dask 延迟读取
close_field = minute_store['data'].fields('close') # 只暴露收盘价字段而不加载整个结构化数组
minute_record_count = int(close_field.shape[0]) # 在关闭文件前保存总记录数
partitioned_close = da.from_array(close_field, chunks=200000) # 按二十万行构造多个惰性分区
partition_count = partitioned_close.npartitions # 记录任务图中的实际分区数
assert partition_count > 1, '分区数必须大于1才能验证分区计算' # 防止单分区伪装为规模化演示
# 无论均值计算成功与否,都在退出本页计算时关闭 HDF5 句柄。
try:
mean_close = float(partitioned_close.mean().compute()) # 触发跨分区均值计算并返回标量
# 用 finally 保证计算异常不会遗留打开的公开数据文件。
finally:
minute_store.close() # 释放对分钟数据版本的文件句柄
print({'文件': minute_path.name, '数据集': 'data.close', '记录数': minute_record_count, '分区数': partition_count, '平均收盘价': round(mean_close, 4)}) # 输出实际分区与计算结果
公开真实数据练习:Dask-ML 下一分钟收益基线
- 输入规模:恒瑞医药分钟数据的前 300,001 条记录。
- 特征与目标:当前分钟收盘价、振幅、对数成交量;目标为下一分钟简单收益。
- 时间切分:前 240,000 个样本训练,后 60,000 个样本最终测试。
- 证据边界:只验证 Dask-ML 接口与时间切分,不声称超内存能力,也不作交易建议。
代码
import time # 记录真实多分区模型的墙钟时间
from dask_ml.preprocessing import StandardScaler as DaskStandardScaler # 在训练分区逐特征估计均值与标准差
from dask_ml.linear_model import LinearRegression as DaskLinearRegression # 拟合显式确定规格的真实分钟收益 L2 线性模型
from sklearn.metrics import mean_squared_error # 统一模型与零收益基线的最终测试指标
local_model_store = h5py.File(minute_path, 'r') # 保持 HDF5 打开以支持惰性字段读取
local_model_count = 300001 # 确定教学计算规模并保留 300000 个监督样本
local_close = da.from_array(local_model_store['data'].fields('close'), chunks=50000)[:local_model_count] # 惰性读取收盘价字段
local_high = da.from_array(local_model_store['data'].fields('high'), chunks=50000)[:local_model_count] # 惰性读取最高价字段
local_low = da.from_array(local_model_store['data'].fields('low'), chunks=50000)[:local_model_count] # 惰性读取最低价字段
local_volume = da.from_array(local_model_store['data'].fields('volume'), chunks=50000)[:local_model_count] # 惰性读取成交量字段
local_features = da.stack([local_close[:-1], (local_high[:-1] - local_low[:-1]) / local_close[:-1], da.log1p(local_volume[:-1])], axis=1) # 组合当分钟收盘价、振幅与对数成交量,预测下一分钟收益。
local_target = local_close[1:] / local_close[:-1] - 1 # 构造严格向前一期的简单收益目标
local_training_end = 240000 # 确定前 80% 训练与后 20% 最终测试边界
local_training_features = local_features[:local_training_end].rechunk((40000, 3)) # 对齐六个训练分区及完整特征轴
local_training_target = local_target[:local_training_end].rechunk(40000) # 与训练特征使用相同分区边界
local_testing_features = local_features[local_training_end:].rechunk((30000, 3)) # 对齐两个最终测试分区
local_testing_target = local_target[local_training_end:].rechunk(30000) # 与最终测试特征使用相同分区边界
代码
local_start_time = time.perf_counter() # 从实际拟合前开始计时
local_dask_scaler = DaskStandardScaler().fit(local_training_features) # 只从训练分区样本逐特征估计均值与标准差
local_training_scaled = local_dask_scaler.transform(local_training_features) # 用训练样本逐特征统计量转换训练分区
local_testing_scaled = local_dask_scaler.transform(local_testing_features) # 用确定的训练样本逐特征统计量转换最终测试分区
local_dask_model = DaskLinearRegression(penalty='l2', C=1.0, solver='admm').fit(local_training_scaled, local_training_target) # 在真实训练分区拟合与前例一致的 L2 惩罚线性模型
local_model_predictions = local_dask_model.predict(local_testing_scaled).compute() # 触发最终测试分区预测计算
local_target_values = local_testing_target.compute() # 只在最终计分时收集最终测试目标
local_model_seconds = time.perf_counter() - local_start_time # 记录拟合与最终测试预测墙钟时间
local_zero_predictions = np.zeros_like(local_target_values) # 构造同窗零收益基线
local_dask_results = pd.DataFrame({'模型': ['Dask-ML L2线性回归', '零收益基线'], '规格': ['penalty=l2; C=1.0; solver=admm', '固定预测=0'], '最终测试MSE': [mean_squared_error(local_target_values, local_model_predictions), mean_squared_error(local_target_values, local_zero_predictions)]}) # 同窗报告确定模型规格与最终测试误差,避免把惩罚回归误称为 OLS
local_model_store.close() # 完成惰性计算后关闭真实 HDF5 文件
print({'训练样本': local_training_end, '最终测试样本': len(local_target_values), '训练分区': local_training_features.npartitions, '最终测试分区': local_testing_features.npartitions, '运行秒': round(local_model_seconds, 2)}) # 输出任务规模与资源边界
display(local_dask_results.round(10)) # 如实展示模型是否胜过零收益基线
{'训练样本': 240000, '最终测试样本': 60000, '训练分区': 6, '最终测试分区': 2, '运行秒': 45.34}
结果不理想时如何解释:若 Dask-ML 未胜零收益基线,只能写“当前固定分钟样本与线性规格没有增量预测价值”。该实验的贡献是验证真实 HDF5 的多分区读取、惰性预处理、拟合与最终测试计分链,而不是证明超内存能力或市场可预测性。
本章小结
- 能把选择性读取、分区变换、惰性任务图、归约与一次最终测试评分串成可检查链。
- 是否
.compute() 取决于结果大小、内存预算与副作用;分区数需由数据规模和调度开销验证。
- 当前实验不证明超内存实际使用能力或市场可预测性。
- 下一章综合资产定价与机器学习,并继续执行点时、基线和失败报告说明。